.NET Core で Kafka トランザクション付きバッチ送信を実装する

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());
    }
}

タグ: .NET Core Kafka トランザクション

7月24日 17:21 投稿