オフセット管理を活用した重複処理防止の実装
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();
}
}