Consume the same message again if processing of the message fails

Viewed 2185

I am using Confluent.Kafka .NET client version 1.3.0. I am following the docs:

var consumerConfig = new ConsumerConfig
{
    BootstrapServers = "server1, server2",
    AutoOffsetReset = AutoOffsetReset.Earliest,
    EnableAutoCommit = true,
    EnableAutoOffsetStore = false,
    GroupId = this.groupId,
    SecurityProtocol = SecurityProtocol.SaslPlaintext,
    SaslMechanism = SaslMechanism.Plain,
    SaslUsername = this.kafkaUsername,
    SaslPassword = this.kafkaPassword,
};

using (var consumer = new ConsumerBuilder<Ignore, string>(consumerConfig).Build())
{
    var cancellationToken = new CancellationTokenSource();
    Console.CancelKeyPress += (_, e) =>
    {
        e.Cancel = true;
        cancellationToken.Cancel();
    };

    consumer.Subscribe("my-topic");
    while (true)
    {
        try
        {
            var consumerResult = consumer.Consume();
            // process message
            consumer.StoreOffset(consumerResult);
        }
        catch (ConsumeException e)
        {
            // log
        }
        catch (KafkaException e)
        {
            // log
        }
        catch (OperationCanceledException e)
        {
            // log
        }
    }
}

The problem is that even if I comment out the line consumer.StoreOffset(consumerResult);, I keep getting the next unconsumed message the next time I Consume, i.e. the offset keeps increasing which doesn't seem to be what the documentation claims it does, i.e. at least one delivery.

Even if I set EnableAutoCommit = false and remove 'EnableAutoOffsetStore = false' from the config, and replace consumer.StoreOffset(consumerResult) with consumer.Commit(), I still see the same behavior, i.e. even if I comment out the Commit, I still keep getting the next unconsumed messages.

I feel like I am missing something fundamental here, but can't figure what. Any help is appreciated!

3 Answers

You may want to have a re-try logic for processing each of your messages for a fixed number of times like say 5. If it doesn't succeed during these 5 retries, you may want to add this message to another topic for handling all failed messages which take precedence over your actual topic. Or you may want to add the failed message to the same topic so that it will be picked up later once all those other messages are consumed.

If the processing of any message is successful within those 5 retries, you can skip to the next message in the queue.

I had the same situation, and here is my solution:

Set a configuration of max retries for each operation.

  • For consuming, just retry.
  • For Saving, re-assign the current offset, and then retry.

Here is the code:

var saveRetries = 0;
var consumeRetries = 0;
ConsumeResult<string, string> consumeResult;

while (true)
{
    try
    {
        consumeResult = consumer.Consume();
        consumeRetries = 0;
    }
    catch (ConsumeException e)
    {
        //Log and retry to consume, up to {MaxConsumeRetries} times
        if (consumeRetries++ >= MaxConsumeRetries)
        {
            throw new OperationCanceledException($"Too many consume retries ({MaxConsumeRetries}). Please check configuration and run the service agian.");
        }
        continue;
    }
    catch (OperationCanceledException oe)
    {
        //Log
        consumer.Close();
        break;
    }

    try
    {
        SaveResult(consumeResult);
        saveRetries = 0;
    }
    catch (ArgumentException ae)
    {
        //Log and retry to save, up to {MaxSaveRetries} times
        if (saveRetries++ < MaxSaveRetries)
        {
            //Assign the same offset, and try again.
            consumer.Assign(consumeResult.TopicPartitionOffset);
            continue;
        }
    }

    try
    {
        consumer.StoreOffset(consumeResult);
    }
    catch (KafkaException ke)
    {
        //Log and let it continue
    }
}
Related