Skip to content
第 199 / 250 章中间件⏱ 12 分钟阅读

第 199 章:Kafka 进阶特性

学习目标

  • 掌握 Kafka 事务
  • 学会消息幂等设计
  • 深入分区与副本机制
  • 了解 Kafka Streams

一、Kafka 幂等性

1.1 什么是幂等 Producer

默认 Producer 可能因重试导致重复消息,幂等 Producer 保证单分区单会话内不重复

yaml
spring:
  kafka:
    producer:
      properties:
        enable.idempotence: true   # 开启幂等性
        acks: all
        retries: 5
        max.in.flight.requests.per.connection: 5

1.2 原理

  • ProducerId:每个 Producer 实例唯一
  • SequenceNumber:每条消息单调递增
  • Broker 端去重:已存在的 (PID, seq) 直接丢弃

二、Kafka 事务

2.1 事务 Producer

java
@Autowired
private KafkaTemplate<String, String> kafkaTemplate;

public void sendInTransaction(List<String> messages) {
    kafkaTemplate.executeInTransaction(ops -> {
        for (String msg : messages) {
            ops.send("topic1", msg);
        }
        // 原子提交:要么全成功,要么全失败
        return true;
    });
}

2.2 多 Topic 事务

java
kafkaTemplate.executeInTransaction(ops -> {
    ops.send("order-topic", order);
    ops.send("audit-topic", auditLog);
    ops.send("notification-topic", notification);
    // 三条消息要么全成功,要么全回滚(对消费者而言)
    return true;
});

2.3 读已提交(read_committed)

yaml
spring:
  kafka:
    consumer:
      properties:
        isolation.level: read_committed   # 只读已提交

事务中的消息对消费者不可见,直到事务提交。

三、消费幂等性

3.1 业务层面幂等

java
@KafkaListener(topics = "order-events")
public void handleOrder(OrderEvent event, Acknowledgment ack) {
    // 1. 唯一键去重(Redis / DB)
    String dedupKey = "kafka:order:" + event.getOrderId();
    if (!redisTemplate.opsForValue().setIfAbsent(dedupKey, "1", 24, TimeUnit.HOURS)) {
        log.info("重复消息,跳过: {}", event.getOrderId());
        ack.acknowledge();
        return;
    }

    // 2. 业务处理
    inventoryService.deduct(event);

    ack.acknowledge();
}

3.2 数据库唯一约束

java
@Transactional
public void handle(OrderEvent event) {
    // 唯一约束兜底
    try {
        orderRepository.save(event.toEntity());
    } catch (DuplicateKeyException e) {
        log.warn("重复消息");
    }
}

四、消息积压处理

4.1 排查

bash
kafka-consumer-groups --bootstrap-server localhost:9092 \
  --describe --group inventory-group

# 看 LAG(滞后量)
# CURRENT-OFFSET   LOG-END-OFFSET   LAG
# 100              10000            9900    ← 积压

4.2 解决方案

扩容消费者

yaml
spring:
  kafka:
    listener:
      concurrency: 10    # 10 个消费者线程

批量消费

yaml
spring:
  kafka:
    consumer:
      max-poll-records: 500   # 一次拉 500 条
    listener:
      type: batch
java
@KafkaListener(topics = "events", containerFactory = "batchFactory")
public void handleBatch(List<Event> events, Acknowledgment ack) {
    // 批量处理
    service.batchProcess(events);
    ack.acknowledge();
}

加分区

bash
kafka-topics --bootstrap-server localhost:9092 \
  --alter --topic events --partitions 10

五、DLQ(Dead Letter Queue)

处理失败的消息放到 DLQ,主流程不受影响。

java
@Configuration
public class KafkaConfig {

    @Bean
    public DefaultErrorHandler errorHandler(KafkaTemplate<String, String> template) {
        // 失败 3 次后发到 DLQ
        DeadLetterPublishingRecoverer recoverer =
            new DeadLetterPublishingRecoverer(template,
                (record, ex) -> new TopicPartition("events.DLT", record.partition()));

        return new DefaultErrorHandler(recoverer, new FixedBackOff(1000L, 3));
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory(
            ConsumerFactory<String, String> consumerFactory,
            DefaultErrorHandler errorHandler) {
        ConcurrentKafkaListenerContainerFactory<String, String> factory =
            new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory);
        factory.setCommonErrorHandler(errorHandler);
        return factory;
    }
}
java
@KafkaListener(topics = "events.DLT")
public void handleDLT(String message) {
    log.error("DLQ 消息: {}", message);
    // 人工处理或重新入队
}

六、分区策略

6.1 默认 Hash

java
// KafkaTemplate.send(topic, key, value)
// key 为 null → 轮询
// key 不为 null → hash(key) % partitions

6.2 自定义分区器

java
public class CustomPartitioner implements Partitioner {
    @Override
    public int partition(String topic, Object key, byte[] keyBytes,
                        Object value, byte[] valueBytes, Cluster cluster) {
        List<PartitionInfo> partitions = cluster.partitionsForTopic(topic);
        int numPartitions = partitions.size();

        if (key == null) {
            return ThreadLocalRandom.current().nextInt(numPartitions);
        }

        // 自定义规则:VIP 用户到独立分区
        String keyStr = key.toString();
        if (keyStr.startsWith("vip_")) {
            return 0;   // VIP 分区
        }

        return Math.abs(keyStr.hashCode()) % numPartitions;
    }
}

七、Kafka Streams

7.1 引入

xml
<dependency>
  <groupId>org.apache.kafka</groupId>
  <artifactId>kafka-streams</artifactId>
</dependency>

7.2 单词计数

java
@Configuration
public class StreamsConfig {

    @Bean
    public KStream<String, Long> wordCount(StreamsBuilder builder) {
        KStream<String, String> source = builder.stream("text-input");

        KTable<String, Long> counts = source
            .flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\s+")))
            .groupBy((key, value) -> value)
            .count();

        counts.toStream().to("word-count-output", Produced.with(Serdes.String(), Serdes.Long()));

        return source;
    }
}

八、Spring Cloud Stream(屏蔽底层)

xml
<dependency>
  <groupId>org.springframework.cloud</groupId>
  <artifactId>spring-cloud-stream-binder-kafka</artifactId>
</dependency>
yaml
spring:
  cloud:
    stream:
      bindings:
        output:
          destination: orders
        input:
          destination: orders
          group: inventory-group
      kafka:
        binder:
          brokers: localhost:9092
java
// 生产
@Service
public class OrderService {
    @Autowired
    private StreamBridge streamBridge;

    public void send(OrderEvent event) {
        streamBridge.send("output", event);
    }
}

// 消费
@Component
public class InventoryConsumer {
    @Bean
    public Consumer<OrderEvent> input() {
        return event -> inventoryService.deduct(event);
    }
}

九、Kafka 监控

9.1 JMX 指标

指标含义
kafka.server:type=BrokerTopicMetrics,name=MessagesInPerSec入速率
kafka.server:type=ReplicaManager,name=UnderReplicatedPartitions副本不足分区
kafka.server:type=BrokerTopicMetrics,name=BytesInPerSec入字节
kafka.network:type=RequestMetrics,name=TotalTimeMs请求延迟

9.2 关键告警

promql
# Lag 告警
kafka_consumergroup_lag{group="my-group"} > 10000

# 副本不足
sum(kafka_topic_partition_in_sync_replica_count) by (topic) == 0

# 磁盘使用率
100 - (kafka_log_log_size / kafka_log_log_capacity * 100) < 10

十、性能调优

10.1 Producer 调优

yaml
spring:
  kafka:
    producer:
      batch-size: 65536           # 64KB 批量
      linger-ms: 10               # 等 10ms 凑批
      compression-type: snappy    # 压缩
      buffer-memory: 67108864     # 64MB 缓冲区

10.2 Broker 调优

properties
# 内存
KAFKA_HEAP_OPTS=-Xms4g -Xmx4g

# 磁盘
log.dirs=/data/kafka
log.retention.hours=72

# 网络
num.network.threads=8
num.io.threads=16

# 分区数预估
# 期望吞吐量 / 单分区吞吐 ≈ 分区数

10.3 Consumer 调优

yaml
spring:
  kafka:
    consumer:
      fetch.min.bytes: 1024       # 至少 1KB 才返回
      fetch.max.wait.ms: 500      # 最多等 500ms
      max.partition.fetch.bytes: 1048576  # 1MB

十一、本章小结

特性解决问题
幂等性Producer 重试重复
事务跨 topic 原子性
DLQ失败消息隔离
批量消费提高吞吐
Kafka Streams实时流处理

动手练习

  1. 实现幂等 Producer(防重复)
  2. 用事务保证多 topic 原子性
  3. 配 DLQ 处理失败消息
  4. 监控 Consumer Lag 并告警

推荐阅读


下一章:第 200 章:RabbitMQ 入门与实战

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