...
|
...
|
@@ -41,6 +41,7 @@ func (c *MessageConsumer) ConsumeClaim(groupSession sarama.ConsumerGroupSession, |
|
|
err error
|
|
|
)
|
|
|
for message := range groupClaim.Messages() {
|
|
|
|
|
|
logs.Debug("Done Message claimed: timestamp = %v, topic = %s offset = %v value = %v \n",
|
|
|
message.Timestamp, message.Topic, message.Offset, string(message.Value))
|
|
|
for i := range c.beforeHandles {
|
...
|
...
|
@@ -118,7 +119,7 @@ func (r *Runer) Start(ctx context.Context) { |
|
|
r.consumerGroup.Close()
|
|
|
return
|
|
|
default:
|
|
|
|
|
|
|
|
|
}
|
|
|
if err := r.consumerGroup.Consume(ctx, r.msgConsumer.topics, r.msgConsumer); err != nil {
|
|
|
logs.Error("consumerGroup err:%s \n", err)
|
...
|
...
|
|