Kafka プロデューサーを使用して複数メッセージを一括送信し、トランザクションにより全メッセージの成功または失敗を保証します。トランザクション ID の設定は必須です。
public void SendMessagesInTransaction(List<string> messages)
{
var config = BuildProducerConfig();
config.EnableIdempotence = true;
config.MessageTimeoutMs = 5000;
config.TransactionalId = Guid.NewGuid().ToString(); // トランザクションIDは必須
using var producer = new ProducerBuilder<string, string>(config)
.SetDefaultPartitioner(RoundRobinPartitioner)
.Build();
producer.InitTransactions(TimeSpan.FromSeconds(60));
try
{
producer.BeginTransaction();
foreach (var payload in messages)
{
var deliveryResult = producer.ProduceAsync(
KafkaTopic.Topic,
new Message<string, string> { Key = "order", Value = payload }
).Result;
_logger.LogInformation("メッセージ '{Value}' を {TopicPartitionOffset} へ送信完了",
deliveryResult.Value, deliveryResult.TopicPartitionOffset);
}
producer.CommitTransaction();
}
catch (ProduceException<string, string> ex)
{
_logger.LogError(ex, "トピック 'order' への送信に失敗: {Reason}", ex.Error.Reason);
producer.AbortTransaction();
}
}
Web API エンドポイントから呼び出す例:
[HttpPost("batch")]
public IActionResult BatchSend([FromBody] List<string> payloads)
{
_kafkaService.SendMessagesInTransaction(payloads);
return Ok();
}
コンシューマ側は前回と同様、Redis でオフセットを管理します:
public void ConsumeWithOffsetTracking()
{
var config = BuildConsumerConfig();
config.AutoOffsetReset = AutoOffsetReset.Earliest;
config.GroupId = "order";
using var consumer = new ConsumerBuilder<string, string>(config).Build();
var cachedOffset = _distributedCache.GetString(KafkaTopic.Topic);
var currentOffset = string.IsNullOrEmpty(cachedOffset) ? 0 : int.Parse(cachedOffset);
consumer.Assign(new TopicPartitionOffset(
KafkaTopic.Topic,
new Partition(0),
currentOffset + 1
));
while (true)
{
consumer.Subscribe(KafkaTopic.Topic);
var message = consumer.Consume();
_logger.LogInformation("受信: Offset={Offset}, Partition={Partition}, Value={Value}",
message.Offset, message.Partition, message.Value);
_distributedCache.SetString(KafkaTopic.Topic, (message.Offset + 1).ToString());
}
}