.NET Core で Apache Kafka を利用した信頼性のあるメッセージ処理実装

オフセット管理を活用した重複処理防止の実装

Apache Kafka を .NET Core で利用する際、コンシューマーがメッセージを受信した後、Kafka が停止/再起動した際に重複処理が発生する問題があります。この問題を解決するため、Redis を活用したオフセット管理を実装します。

問題の背景

コンシューマーがメッセージを受信後、StoreOffset() を呼び出してオフセットを確認しますが、Kafka サービスがダウンした際に再起動時に同じメッセージを再度処理する可能性があります。このため、Redis に現在のオフセット値を永続化し、再起動時のオフセットリセットを制御します。

実装コード

以下は修正済みのメッセージブローカー実装です。

namespace MessageBroker
{
    public class BrokerSettings
    {
        public string BootstrapServers { get; set; }
    }

    public static class TopicNames
    {
        public const string OrderTopic = "order-events";
    }

    public class MessageBroker
    {
        private readonly BrokerSettings _brokerSettings;
        private readonly IDistributedCache _cache;
        private readonly ILogger<MessageBroker> _logger;

        public MessageBroker(
            IOptionsMonitor<BrokerSettings> settings,
            IDistributedCache cache,
            ILogger<MessageBroker> logger)
        {
            _brokerSettings = settings.CurrentValue;
            _cache = cache;
            _logger = logger;
        }

        public ProducerConfig CreateProducer()
        {
            return new ProducerConfig
            {
                BootstrapServers = _brokerSettings.BootstrapServers,
                MessageTimeoutMs = 5000,
                EnableIdempotence = true
            };
        }

        public ConsumerConfig CreateConsumer()
        {
            return new ConsumerConfig
            {
                BootstrapServers = _brokerSettings.BootstrapServers,
                AutoOffsetReset = AutoOffsetReset.Earliest,
                GroupId = "order-processing"
            };
        }

        public void Publish(string message)
        {
            using var producer = new ProducerBuilder<string, string>(CreateProducer())
                .SetDefaultPartitioner(RoundRobinPartitioner)
                .Build();

            var deliveryResult = producer.Produce(
                TopicNames.OrderTopic,
                new Message<string, string> { Key = "order", Value = message });

            _logger.LogInformation($"メッセージ送信成功: {deliveryResult.Offset}");
        }

        public void Consume()
        {
            var consumer = new ConsumerBuilder<string, string>(CreateConsumer())
                .Build();

            consumer.Subscribe(TopicNames.OrderTopic);
            var storedOffset = GetStoredOffset();
            consumer.Assign(new TopicPartitionOffset(
                TopicNames.OrderTopic, 
                new Partition(0), 
                storedOffset));

            while (true)
            {
                var message = consumer.Consume();
                _logger.LogInformation($"受信: {message.Value} (オフセット: {message.Offset})");

                // Redis に現在のオフセットを保存
                _cache.SetString(TopicNames.OrderTopic, message.Offset.ToString());
            }
        }

        private long GetStoredOffset()
        {
            var offsetStr = _cache.GetString(TopicNames.OrderTopic);
            return string.IsNullOrEmpty(offsetStr) 
                ? 0 
                : long.Parse(offsetStr);
        }

        private static int _partitionCounter;
        private Partition RoundRobinPartitioner(
            string topic, 
            int partitionCount, 
            ReadOnlySpan<byte> key, 
            bool isKeyNull)
        {
            var partition = _partitionCounter % partitionCount;
            _partitionCounter++;
            return new Partition(partition);
        }
    }
}

API サービスの設定例

Web API プロジェクトでの設定例です。

var builder = WebApplication.CreateBuilder(args);

builder.Services.Configure<BrokerSettings>(
    builder.Configuration.GetSection("KafkaSettings"));

builder.Services.AddDistributedRedisCache(options =>
{
    options.Configuration = builder.Configuration["RedisConnection"];
});

builder.Services.AddTransient<MessageBroker>();
builder.Services.AddControllers();

バックグラウンドサービスの実装

メッセージ処理をバックグラウンドで実行するサービス例です。

public class OrderProcessor : BackgroundService
{
    private readonly MessageBroker _broker;

    public OrderProcessor(MessageBroker broker) 
        => _broker = broker;

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        _broker.Consume();
    }
}

タグ: .NET Core Apache Kafka redis Distributed Systems Message Queue

7月31日 01:00 投稿