Skip to content

Commit e4598b7

Browse files
authored
Merge pull request #846 from Shopify/4-0-fix
Use sarama version to build the correct OffsetRequest call
2 parents a15e5c1 + 18dc676 commit e4598b7

2 files changed

Lines changed: 20 additions & 1 deletion

File tree

core/internal/cluster/kafka_cluster.go

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -227,6 +227,22 @@ func (module *KafkaCluster) generateOffsetRequests(client helpers.SaramaClient)
227227
}
228228
if _, ok := requests[broker.ID()]; !ok {
229229
requests[broker.ID()] = &sarama.OffsetRequest{}
230+
// Match the version of the client as sarama's getOffset function does
231+
// https://github.com/IBM/sarama/blob/main/client.go#L863-L876
232+
if client.Config().Version.IsAtLeast(sarama.V2_1_0_0) {
233+
// Version 4 adds the current leader epoch, which is used for fencing.
234+
requests[broker.ID()].Version = 4
235+
} else if client.Config().Version.IsAtLeast(sarama.V2_0_0_0) {
236+
// Version 3 is the same as version 2.
237+
requests[broker.ID()].Version = 3
238+
} else if client.Config().Version.IsAtLeast(sarama.V0_11_0_0) {
239+
// Version 2 adds the isolation level, which is used for transactional reads.
240+
requests[broker.ID()].Version = 2
241+
} else if client.Config().Version.IsAtLeast(sarama.V0_10_1_0) {
242+
// Version 1 removes MaxNumOffsets. From this version forward, only a single
243+
// offset can be returned.
244+
requests[broker.ID()].Version = 1
245+
}
230246
}
231247
brokers[broker.ID()] = broker
232248
requests[broker.ID()].AddBlock(topic, partitionID, sarama.OffsetNewest, 1)

core/internal/cluster/kafka_cluster_test.go

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -172,6 +172,7 @@ func TestKafkaCluster_generateOffsetRequests(t *testing.T) {
172172
// Set up the mock to return the leader broker for a test topic and partition
173173
client := &helpers.MockSaramaClient{}
174174
client.On("Leader", "testtopic", int32(0)).Return(broker, nil)
175+
client.On("Config").Return(sarama.NewConfig())
175176

176177
requests, brokers := module.generateOffsetRequests(client)
177178

@@ -199,6 +200,7 @@ func TestKafkaCluster_generateOffsetRequests_NoLeader(t *testing.T) {
199200
var nilBroker *helpers.BurrowSaramaBroker
200201
client.On("Leader", "testtopic", int32(0)).Return(nilBroker, errors.New("no leader error"))
201202
client.On("Leader", "testtopic", int32(1)).Return(broker, nil)
203+
client.On("Config").Return(sarama.NewConfig())
202204

203205
requests, brokers := module.generateOffsetRequests(client)
204206

@@ -233,6 +235,7 @@ func TestKafkaCluster_getOffsets(t *testing.T) {
233235
var nilBroker *helpers.BurrowSaramaBroker
234236
client.On("Leader", "testtopic", int32(0)).Return(broker, nil)
235237
client.On("Leader", "testtopic", int32(1)).Return(nilBroker, errors.New("no leader error"))
238+
client.On("Config").Return(sarama.NewConfig())
236239

237240
go module.getOffsets(client)
238241
request := <-module.App.StorageChannel
@@ -274,7 +277,7 @@ func TestKafkaCluster_getOffsets_BrokerFailed(t *testing.T) {
274277
// Set up the mock to return the leader broker for a test topic and partition
275278
client := &helpers.MockSaramaClient{}
276279
client.On("Leader", "testtopic", int32(0)).Return(broker, nil)
277-
280+
client.On("Config").Return(sarama.NewConfig())
278281
module.getOffsets(client)
279282

280283
broker.AssertExpectations(t)

0 commit comments

Comments
 (0)