KafkaNotes

第 09 章:Spring Boot 集成 Kafka

zjc 于 2026-01-09 发布

这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 Spring Kafka 是 Java 生态最常用的 Kafka 封装。本章从发送、消费、手动确认、错误处理讲到死信队列与批量监听,给出一套可以直接落地的工程模板。

9.1 依赖与配置

<dependency>
    <groupId>org.springframework.kafka</groupId>
    <artifactId>spring-kafka</artifactId>
</dependency>

Spring Boot 3.x 的依赖管理会自动匹配合适的 spring-kafka 版本。

application.yml

spring:
  kafka:
    bootstrap-servers: localhost:9092
    producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
      acks: all
      properties:
        enable.idempotence: true
        linger.ms: 10
        compression.type: zstd
    consumer:
      group-id: order-service
      auto-offset-reset: earliest
      enable-auto-commit: false
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
      properties:
        spring.json.trusted.packages: "com.example.dto.*"
        isolation.level: read_committed
    listener:
      ack-mode: manual
      concurrency: 3

要点:

9.2 发送消息:KafkaTemplate

@Service
@RequiredArgsConstructor
public class OrderEventPublisher {
    private final KafkaTemplate<String, OrderCreated> kafkaTemplate;

    public void publish(OrderCreated event) {
        kafkaTemplate.send("order-events", event.orderId(), event)
                .whenComplete((result, ex) -> {
                    if (ex != null) {
                        log.error("send failed, orderId={}", event.orderId(), ex);
                        // 落补偿表或触发重试
                    } else {
                        var meta = result.getRecordMetadata();
                        log.debug("sent partition={} offset={}",
                                meta.partition(), meta.offset());
                    }
                });
    }
}

同步发送(需要强确认时):

kafkaTemplate.send("order-events", key, event).get(3, TimeUnit.SECONDS);

9.3 消费消息:@KafkaListener

@Component
@RequiredArgsConstructor
@Slf4j
public class OrderEventListener {
    private final OrderService orderService;

    @KafkaListener(
        topics = "order-events",
        groupId = "order-service",
        concurrency = "3"
    )
    public void onMessage(ConsumerRecord<String, OrderCreated> record,
                          Acknowledgment ack) {
        try {
            orderService.handle(record.value());
            ack.acknowledge();   // 处理成功才提交
        } catch (Exception e) {
            log.error("handle failed, offset={}", record.offset(), e);
            throw e;             // 交给错误处理器决定重试或进 DLT
        }
    }
}

9.4 错误处理与死信队列

直接抛异常会触发容器重试,默认行为可能导致无限循环。生产推荐“有限重试 + 死信”:

@Configuration
public class KafkaErrorConfig {

    @Bean
    public DefaultErrorHandler errorHandler(KafkaTemplate<Object, Object> template) {
        // 重试 3 次,间隔 1s,指数退避倍数 2
        var recoverer = new DeadLetterPublishingRecoverer(
                template,
                (record, ex) -> {
                    // 死信目标:原 topic + ".DLT",保持原分区
                    return new TopicPartition(record.topic() + ".DLT", record.partition());
                });

        var backOff = new ExponentialBackOffWithMaxRetries(3);
        backOff.setInitialInterval(1000L);
        backOff.setMultiplier(2.0);

        var handler = new DefaultErrorHandler(recoverer, backOff);
        // 这些异常不重试,直接进死信
        handler.addNotRetryableExceptions(
                DeserializationException.class,
                IllegalArgumentException.class);
        return handler;
    }
}

配套监听死信:

@KafkaListener(topics = "order-events.DLT", groupId = "order-dlt-watcher")
public void onDead(ConsumerRecord<String, OrderCreated> record) {
    log.error("dead letter key={} offset={} value={}",
            record.key(), record.offset(), record.value());
    // 人工处理、告警、回补
}

死信设计建议:

9.5 批量消费

下游是数据库或搜索引擎时,批量写入能带来数量级的性能提升:

@KafkaListener(topics = "order-events", batch = "true")
public void onBatch(List<ConsumerRecord<String, OrderCreated>> records,
                    Acknowledgment ack) {
    List<OrderCreated> events = records.stream()
            .map(ConsumerRecord::value).toList();
    orderService.batchHandle(events);
    ack.acknowledge();
}

spring.kafka.consumer:
  max-poll-records: 200
  fetch-min-size: 1KB
  fetch-max-wait: 500ms

9.6 消费者多线程

concurrency=3 相当于起 3 个 KafkaMessageListenerContainer,各自拥有独立 consumer,各自分到一部分分区。这是首选的并发模型,因为它保留了标准的位移与再平衡语义。

只有当单条消息处理极重(比如调用多个慢接口)且分区数已到上限时,才考虑在监听器内部再套线程池。这时要自己解决:

  1. 位移提交粒度(一批任务全完成后再 ack);
  2. 优雅停机(关闭线程池、等待任务完成);
  3. 异常时的消息路由(重试/死信)。

9.7 事务生产者

需要“多条消息原子写入”或 consume-transform-produce 场景时:

kafkaTemplate.executeInTransaction(ops -> {
    ops.send("order-events", key, orderCreated);
    ops.send("audit-events", key, auditEvent);
    return true;
});

配置:

spring:
  kafka:
    producer:
      transaction-id-prefix: order-tx-

消费端记得 isolation.level: read_committed,否则会读到未提交数据(第 14 章展开)。

9.8 一个完整的消费端模板

@Component
@Slf4j
public class ReliableOrderListener {

    private final OrderService orderService;
    private final MeterRegistry metrics;

    public ReliableOrderListener(OrderService orderService, MeterRegistry metrics) {
        this.orderService = orderService;
        this.metrics = metrics;
    }

    @KafkaListener(topics = "order-events", groupId = "order-service")
    public void listen(ConsumerRecord<String, OrderCreated> record, Acknowledgment ack) {
        long start = System.currentTimeMillis();
        try {
            orderService.handleIdempotent(record.value());
            ack.acknowledge();
            metrics.counter("kafka.consume.success").increment();
        } catch (Exception e) {
            metrics.counter("kafka.consume.failure").increment();
            throw e;
        } finally {
            metrics.timer("kafka.consume.latency")
                  .record(System.currentTimeMillis() - start, TimeUnit.MILLISECONDS);
        }
    }
}

这套结构包含:手动 ack、幂等处理、指标埋点、异常上抛交给统一错误处理。可以直接作为业务消费者模板复制。

本章小结

思考题

  1. concurrency=6 但 topic 只有 3 个分区,会发生什么?
  2. 监听器里 ack.acknowledge() 之后进程立即崩溃,这条消息会怎样?不 ack 又会怎样?
  3. 如何设计死信消息的元信息,才能让“重新投递”变得可操作?