Skip to content

Commit

Permalink
Patch sarama_kafka rebalance fix
Browse files Browse the repository at this point in the history
  • Loading branch information
nbajaj90 committed Nov 3, 2022
1 parent 6784a56 commit 01c4f13
Showing 1 changed file with 17 additions and 11 deletions.
28 changes: 17 additions & 11 deletions protocol/kafka_sarama/v2/receiver.go
Original file line number Diff line number Diff line change
Expand Up @@ -53,19 +53,25 @@ func (r *Receiver) Close(context.Context) error {
}

func (r *Receiver) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
for message := range claim.Messages() {
msg := message
m := NewMessageFromConsumerMessage(msg)

r.incoming <- msgErr{
msg: binding.WithFinish(m, func(err error) {
if protocol.IsACK(err) {
session.MarkMessage(msg, "")
}
}),
for {
select {
case msg, ok := <-claim.Messages():
if !ok {
return nil
}
m := NewMessageFromConsumerMessage(msg)

r.incoming <- msgErr{
msg: binding.WithFinish(m, func(err error) {
if protocol.IsACK(err) {
session.MarkMessage(msg, "")
}
}),
}
case <-session.Context().Done():
return nil
}
}
return nil
}

func (r *Receiver) Receive(ctx context.Context) (binding.Message, error) {
Expand Down

0 comments on commit 01c4f13

Please sign in to comment.