Skip to content
第 63 / 250 章后端⏱ 10 分钟阅读

第 63 章:Kafka 与流处理

学习目标

  • 理解 Kafka 的核心架构与设计哲学
  • 掌握 Spring Kafka 集成与生产者消费者模式
  • 学会 Kafka 在日志、事件溯源、CDC 中的应用

一、为什么需要 Kafka?

Kafka 的本质:分布式持久化日志

维度RabbitMQKafka
定位消息队列(任务分发)分布式日志(事件流)
吞吐量万级百万级
消费模式队列(消费即删除)日志(可重放,保留 N 天)
延迟微秒级毫秒级
场景业务异步、可靠投递日志、事件溯源、CDC、流计算

二、Kafka 核心概念

概念说明
Producer生产者
Consumer消费者
Topic主题(消息分类)
Partition分区(Topic 的并行单位)
Offset消息在分区里的位置(类似数组下标)
BrokerKafka 服务器节点
Consumer Group消费者组(一个分区只被组内一个消费者消费)

三、Spring Boot 集成

xml
<dependency>
    <groupId>org.springframework.kafka</groupId>
    <artifactId>spring-kafka</artifactId>
</dependency>
yaml
spring:
  kafka:
    bootstrap-servers: localhost:9092,localhost:9093,localhost:9094
    producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
      acks: all                              # ① 最强一致性
      retries: 3
      properties:
        enable.idempotence: true             # ② 幂等生产者(防重复)
        max.in.flight.requests.per.connection: 5
        linger.ms: 20                        # 批量发送延迟
        compression.type: snappy             # 压缩
    consumer:
      group-id: taskflow-consumer
      auto-offset-reset: earliest
      enable-auto-commit: false              # ③ 手动提交 offset
      max-poll-records: 100
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
      properties:
        spring.json.trusted.packages: "com.taskflow.**"
    listener:
      ack-mode: manual_immediate
      concurrency: 3                         # 并发消费者数

四、生产者

java
@Component
@RequiredArgsConstructor
@Slf4j
public class UserEventProducer {

    private final KafkaTemplate<String, Object> kafkaTemplate;

    // ① 发送事件(key 保证同一用户的消息有序)
    public void publishUserEvent(UserEvent event) {
        String key = String.valueOf(event.getUserId());    // ① 同用户的事件进同一分区

        CompletableFuture<SendResult<String, Object>> future =
                kafkaTemplate.send("user.events", key, event);

        // ② 回调处理发送结果
        future.whenComplete((result, ex) -> {
            if (ex != null) {
                log.error("发送失败", ex);
                // 失败重试或入本地消息表
            } else {
                log.info("发送成功 partition={} offset={}",
                        result.getRecordMetadata().partition(),
                        result.getRecordMetadata().offset());
            }
        });
    }

    // ② 同步发送(业务关键场景)
    public void publishSync(UserEvent event) {
        try {
            SendResult<String, Object> result =
                    kafkaTemplate.send("user.events", String.valueOf(event.getUserId()), event)
                            .get(5, TimeUnit.SECONDS);
            log.info("同步发送成功 offset={}",
                    result.getRecordMetadata().offset());
        } catch (Exception e) {
            log.error("同步发送失败", e);
            throw new BusinessException("事件发布失败");
        }
    }
}

KafkaTemplate 的事务支持

java
@Configuration
public class KafkaConfig {

    @Bean
    public KafkaTemplate<String, Object> kafkaTemplate(
            ProducerFactory<String, Object> factory) {
        KafkaTemplate<String, Object> template = new KafkaTemplate<>(factory);
        // ① 开启事务(生产消息 + DB 操作可以事务化)
        template.setTransactionIdPrefix("taskflow-tx-");
        return template;
    }
}

// 用法
@Transactional("kafkaTransactionManager")
public void placeOrder(OrderDTO dto) {
    orderMapper.insert(order);
    kafkaTemplate.send("order.events", order);   // 与 DB 同事务
}

五、消费者

java
@Component
@Slf4j
public class UserEventConsumer {

    // ① 监听主题
    @KafkaListener(
        topics = "user.events",
        groupId = "user-service",
        concurrency = "3"
    )
    public void onMessage(
            ConsumerRecord<String, UserEvent> record,
            @Header(KafkaHeaders.RECEIVED_TOPIC) String topic,
            @Header(KafkaHeaders.RECEIVED_PARTITION) int partition,
            @Header(KafkaHeaders.OFFSET) long offset,
            Acknowledgment ack) {

        try {
            log.info("收到消息 topic={} partition={} offset={} key={} value={}",
                    topic, partition, offset, record.key(), record.value());

            UserEvent event = record.value();
            processEvent(event);

            // ② 手动提交 offset
            ack.acknowledge();

        } catch (Exception e) {
            log.error("处理失败", e);
            // ③ 不提交 offset,下次重新消费
            // 或转入死信主题
            kafkaTemplate.send("user.events.dlq", record.value());
            ack.acknowledge();   // 即使失败也 ack,避免死循环
        }
    }

    private void processEvent(UserEvent event) {
        switch (event.getType()) {
            case REGISTER -> sendWelcomeEmail(event);
            case LOGIN    -> updateLastLogin(event);
            case PURCHASE -> awardPoints(event);
        }
    }
}

批量消费

yaml
spring:
  kafka:
    listener:
      type: batch                       # ① 批量消费
      ack-mode: manual
java
@KafkaListener(topics = "user.events", groupId = "user-batch-consumer")
public void onBatch(List<UserEvent> events, Acknowledgment ack) {
    log.info("批量收到 {} 条消息", events.size());

    // ② 批量处理(数据库批量插入)
    userEventMapper.batchInsert(events);

    // ③ 一次性 ack
    ack.acknowledge();
}

六、消费者组与 Rebalance

关键规则

  • 一个分区只能被组内一个消费者消费
  • 消费者数 > 分区数 → 多余的消费者闲置
  • 增减消费者会触发 Rebalance,期间整个组暂停消费

七、Offset 管理

java
// ① 自动提交(不推荐)
enable-auto-commit: true
// 风险:处理失败但 offset 已提交 → 消息丢失

// ② 手动提交(推荐)
enable-auto-commit: false
ack.acknowledge();    // 处理成功才提交

// ③ 提交特定 offset
ack.acknowledge();

// ④ 从指定 offset 开始消费
@KafkaListener(topicPartitions = @TopicPartition(
    topic = "user.events",
    partitions = {"0", "1"},
    partitionOffsets = @PartitionOffset(partition = "0", initialOffset = "100")
))
public void onMessage(...) { ... }

八、Kafka 应用场景

1. 日志收集

yaml
# Filebeat → Kafka → Logstash → ES
filebeat.inputs:
  - type: log
    paths:
      - /var/log/app/*.log

output.kafka:
  hosts: ["localhost:9092"]
  topic: app-logs

2. 事件溯源(Event Sourcing)

java
// 不存当前状态,存所有事件,从头重放就能恢复状态
public void placeOrder(OrderDTO dto) {
    OrderCreatedEvent event = new OrderCreatedEvent(...);
    kafkaTemplate.send("order.events", event);  // 存事件
}

// 查询订单状态:消费所有相关事件,重建状态
public Order getOrderState(Long orderId) {
    Order order = new Order();
    // 读所有 OrderCreatedEvent, OrderPaidEvent, OrderShippedEvent...
    // 按顺序应用,重建订单状态
    return order;
}

3. CDC(Change Data Capture)

典型场景:Canal/Debezium 监听 MySQL binlog,把变更事件发到 Kafka,下游各系统各自消费。

4. 用户行为追踪

java
// 用户点击 → Kafka → 实时分析
@KafkaListener(topics = "user.click")
public void trackClick(ClickEvent event) {
    // 实时统计热点内容
    // 推荐系统特征计算
    // 实时大屏
}

九、Kafka 高级特性

消息顺序性

java
// 保证同一用户的事件有序:partition key 用 userId
kafkaTemplate.send("user.events", String.valueOf(event.getUserId()), event);
// 同 userId → 同一 partition → 严格有序

Exactly-Once 语义

yaml
producer:
  acks: all                                    # ① Leader + 所有 ISR 都确认
  properties:
    enable.idempotence: true                   # ② 幂等(防重复)
    transactional.id: "taskflow-tx-1"         # ③ 事务 ID(跨分区原子)

consumer:
  isolation-level: read_committed              # ④ 只读已提交

消息压缩

yaml
producer:
  properties:
    compression.type: snappy    # snappy/lz4/gzip,吞吐量提升 2-3 倍

十、Kafka 运维要点

监控指标

bash
# Kafka 自带命令
kafka-topics.sh --list --bootstrap-server localhost:9092
kafka-consumer-groups.sh --list --bootstrap-server localhost:9092
kafka-consumer-groups.sh --describe --group taskflow --bootstrap-server localhost:9092
关键指标健康值
消息积压(Lag)< 1000
生产 TPS取决于业务
消费 TPS接近生产 TPS
分区数= 业务并行度

常见问题

bash
# ① 消费者启动后不消费
# 检查 auto-offset-reset 配置

# ② Lag 一直涨
# 增加消费者实例(不超过分区数)
# 优化消费逻辑

# ③ Rebalance 频繁
# 调大 session.timeout.ms
# 优化处理逻辑,避免超时

十一、本章小结

要点关键
本质分布式持久化日志(区别于消息队列)
吞吐量百万级
核心概念Topic / Partition / Offset / Consumer Group
分区并行单位,同 key 进同一分区保证有序
消费者组一个分区只被组内一个消费者消费
可靠投递acks=all + idempotence + transaction
幂等消费enable.auto.commit=false + 手动 ack
应用场景日志、事件溯源、CDC、用户行为、实时分析

动手练习

练习 1:基础题

实现用户注册事件:服务 A 写入数据库后发 Kafka 消息,服务 B 消费消息发送欢迎邮件。

练习 2:进阶题

实现订单事件的 Exactly-Once:生产者事务 + 数据库操作,消费者手动提交 offset + 业务幂等。

练习 3:思考题

你的系统每天产生 1 亿条用户行为日志。如何设计 Kafka 集群和消费者,保证数据不丢、可重放、可被多个下游系统独立消费?


下一章第 64 章:Spring Security 入门

本站基于 VitePress 构建 · 由 Codebook 团队维护