之前做一个数据同步的服务,需要在 .NET 5 的控制台项目里跑一个后台任务去消费 Kafka,要求挂了能自动重连,关闭的时候能优雅退出不丢消息。折腾了一圈,记录一下。


依赖

dotnet add package Confluent.Kafka
dotnet add package Microsoft.Extensions.Hosting

Confluent.Kafka 是 .NET 里用得最多的 Kafka 客户端,Microsoft.Extensions.Hosting 提供 BackgroundService,用它来托管后台任务比自己 new Thread 靠谱得多。

Program.cs

using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;

Host.CreateDefaultBuilder(args)
    .ConfigureServices(services =>
    {
        services.AddHostedService<KafkaConsumerService>();
    })
    .Build()
    .Run();

Host 托管有个好处,它帮你处理了进程信号,Docker 发 SIGTERM 过来的时候,它会通知你的 CancellationToken,不用自己写 while(true) 卡着。

KafkaConsumerService.cs

这个是核心,直接上完整代码:

using Confluent.Kafka;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;

public class KafkaConsumerService : BackgroundService
{
    private readonly ILogger<KafkaConsumerService> _logger;
    private readonly ConsumerConfig _config;
    private const string Topic = "your-topic";

    public KafkaConsumerService(ILogger<KafkaConsumerService> logger)
    {
        _logger = logger;
        _config = new ConsumerConfig
        {
            BootstrapServers = "localhost:9092",
            GroupId = "your-consumer-group",
            AutoOffsetReset = AutoOffsetReset.Earliest,
            EnableAutoCommit = false,
            SessionTimeoutMs = 10000,
            HeartbeatIntervalMs = 3000
        };
    }

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        await Task.Yield();

        while (!stoppingToken.IsCancellationRequested)
        {
            try
            {
                await ConsumeLoop(stoppingToken);
            }
            catch (OperationCanceledException)
            {
                // 正常退出,不用管
            }
            catch (Exception ex)
            {
                _logger.LogError(ex, "Kafka 消费异常,5 秒后重连");
                try
                {
                    await Task.Delay(5000, stoppingToken);
                }
                catch (OperationCanceledException)
                {
                    break;
                }
            }
        }
    }

    private async Task ConsumeLoop(CancellationToken stoppingToken)
    {
        using var consumer = new ConsumerBuilder<Ignore, string>(_config)
            .SetErrorHandler((_, e) =>
            {
                _logger.LogWarning("Kafka 错误: {Reason}", e.Reason);
            })
            .Build();

        consumer.Subscribe(Topic);
        _logger.LogInformation("开始消费 Topic: {Topic}", Topic);

        try
        {
            while (!stoppingToken.IsCancellationRequested)
            {
                var result = consumer.Consume(stoppingToken);
                if (result?.Message == null) continue;

                try
                {
                    await HandleMessage(result.Message.Value);
                    consumer.Commit(result);
                }
                catch (Exception ex)
                {
                    _logger.LogError(ex, "处理消息失败: {Value}",
                        result.Message.Value);
                }
            }
        }
        finally
        {
            consumer.Close();
        }
    }

    private Task HandleMessage(string message)
    {
        _logger.LogInformation("收到消息: {Message}", message);
        // 业务逻辑写这里
        return Task.CompletedTask;
    }
}

说几个细节

自动重连是怎么做的

ExecuteAsync 里的结构,外层是一个 while 循环套着 try-catchConsumeLoop 里面不管是 Broker 挂了、网络断了还是 Rebalance 出了问题,异常都会被外层接住,等 5 秒重新走一遍 ConsumeLoop,也就是重新创建 Consumer、重新 Subscribe。这就是所谓的"守护",不用单独开一个线程去监控,这样写反而更干净。

await Task.Yield() 是干嘛的

ExecuteAsync 一进来如果直接走到 consumer.Consume(),这个调用是阻塞的,会把宿主的启动流程卡住,别的 HostedService 都没法启动。加一行 Task.Yield() 让出执行权,后面的代码会跑到线程池线程上去,不影响宿主启动。

为什么关掉自动提交

EnableAutoCommit = false,消息处理完了手动 Commit。好处是如果 HandleMessage 里抛异常了,offset 不会被提交,下次重启还能消费到这条消息。

坏处是如果你的业务不幂等,重复消费可能有问题。如果业务是幂等的或者你不在意偶尔丢一两条,直接把 EnableAutoCommit 设成 true 也行,省事。

consumer.Close() 放在 finally 里

不管是正常退出还是异常退出,Close() 都会执行,它会通知 Kafka Broker 这个消费者要离开消费组了,Broker 就能立刻触发 Rebalance,不用傻等 session.timeout.ms 超时。

Docker 里的关闭链路

docker stop 会给容器发 SIGTERM,Host 默认监听这个信号 → stoppingToken 被取消 → consumer.Consume()OperationCanceledException → 跳出循环 → finallyconsumer.Close() → 外层 while 检查到 IsCancellationRequestedExecuteAsync 返回 → 进程退出。

不会卡住也不会丢消息。注意 Docker 默认给 10 秒的优雅关闭时间,如果你的 HandleMessage 处理比较慢,可以在 docker-compose.yml 里加个 stop_grace_period: 30s 延长一下。

生产上要注意的

配置别写死在代码里,放到 appsettings.json 或者配置中心。上面写死是为了看着清楚。

HandleMessage 如果涉及到数据库操作,记得用 IServiceScopeFactory 创建 Scope 来解析 Scoped 服务,BackgroundService 本身是 Singleton 的,直接注入 DbContext 之类的 Scoped 服务会出问题。

如果消费量大,可以考虑批量消费 + 批量提交,不用每条消息都 Commit 一次,减少和 Broker 的交互。