mirror of
https://github.com/Syngnat/GoNavi.git
synced 2026-08-11 01:03:51 +08:00
🐛 fix(kafka): 修复主题预览偏移量越界
- 使用绝对 offset 定位 Kafka 分区消息 - 避免 earliest offset 非零时重复叠加 - 增加直连预览读取的回归测试
This commit is contained in:
@@ -970,6 +970,15 @@ func (r *kafkaGoRuntime) partitionOffsets(ctx context.Context, topic string, par
|
||||
return conn.ReadOffsets()
|
||||
}
|
||||
|
||||
type kafkaOffsetSeeker interface {
|
||||
Seek(offset int64, whence int) (int64, error)
|
||||
}
|
||||
|
||||
func seekKafkaAbsoluteOffset(conn kafkaOffsetSeeker, offset int64) error {
|
||||
_, err := conn.Seek(offset, kafka.SeekAbsolute)
|
||||
return err
|
||||
}
|
||||
|
||||
func (r *kafkaGoRuntime) fetchPartitionMessages(ctx context.Context, topic string, partitionID int, startOffset int64, limit int) ([]kafkaMessageRecord, error) {
|
||||
conn, err := r.dialer.DialLeader(ctx, "tcp", r.bootstrap, topic, partitionID)
|
||||
if err != nil {
|
||||
@@ -977,7 +986,7 @@ func (r *kafkaGoRuntime) fetchPartitionMessages(ctx context.Context, topic strin
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
if _, err := conn.Seek(startOffset, io.SeekStart); err != nil {
|
||||
if err := seekKafkaAbsoluteOffset(conn, startOffset); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
deadline := time.Now().Add(r.readWait)
|
||||
|
||||
@@ -21,6 +21,30 @@ type fakeKafkaRuntime struct {
|
||||
lastPublishCommand kafkaPublishCommand
|
||||
}
|
||||
|
||||
type kafkaOffsetSeekerRecorder struct {
|
||||
firstOffset int64
|
||||
lastOffset int64
|
||||
seekOffset int64
|
||||
seekWhence int
|
||||
}
|
||||
|
||||
func (s *kafkaOffsetSeekerRecorder) Seek(offset int64, whence int) (int64, error) {
|
||||
s.seekOffset = offset
|
||||
s.seekWhence = whence
|
||||
switch whence {
|
||||
case kafka.SeekStart:
|
||||
offset += s.firstOffset
|
||||
case kafka.SeekAbsolute:
|
||||
// offset is already absolute.
|
||||
default:
|
||||
return 0, kafka.OffsetOutOfRange
|
||||
}
|
||||
if offset < s.firstOffset || offset > s.lastOffset {
|
||||
return 0, kafka.OffsetOutOfRange
|
||||
}
|
||||
return offset, nil
|
||||
}
|
||||
|
||||
func (f *fakeKafkaRuntime) Close() error { return nil }
|
||||
|
||||
func (f *fakeKafkaRuntime) Ping(ctx context.Context) error { return nil }
|
||||
@@ -44,6 +68,23 @@ func (f *fakeKafkaRuntime) Publish(ctx context.Context, command kafkaPublishComm
|
||||
return f.publishAffected, nil
|
||||
}
|
||||
|
||||
func TestSeekKafkaAbsoluteOffsetUsesAbsoluteKafkaOffset(t *testing.T) {
|
||||
seeker := &kafkaOffsetSeekerRecorder{
|
||||
firstOffset: 100,
|
||||
lastOffset: 150,
|
||||
}
|
||||
|
||||
if err := seekKafkaAbsoluteOffset(seeker, 100); err != nil {
|
||||
t.Fatalf("seek at retained first offset failed: %v", err)
|
||||
}
|
||||
if seeker.seekOffset != 100 {
|
||||
t.Fatalf("expected absolute offset 100, got %d", seeker.seekOffset)
|
||||
}
|
||||
if seeker.seekWhence != kafka.SeekAbsolute {
|
||||
t.Fatalf("expected kafka.SeekAbsolute, got %d", seeker.seekWhence)
|
||||
}
|
||||
}
|
||||
|
||||
func TestNormalizeKafkaConfigParsesURIAndParams(t *testing.T) {
|
||||
config := normalizeKafkaConfig(connection.ConnectionConfig{
|
||||
URI: "kafka://alice:secret@127.0.0.1:9092,127.0.0.2:9093/orders.events?topology=cluster&tls=true&skip_verify=true",
|
||||
|
||||
Reference in New Issue
Block a user