KafkaNotes

第 05 章:生产者深度实战

zjc 于 2026-01-05 发布

这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 本章进入代码实战:用一个完整的 Java 项目理解生产者的发送流程、分区策略、拦截器、核心参数与顺序性保证。这些知识同样适用于 Python(kafka-python / confluent-kafka)、Go(franz-go / sarama)等客户端,因为参数语义是 Kafka 协议层统一的。

5.1 Maven 依赖

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-clients</artifactId>
    <version>3.9.0</version>
</dependency>

学习时选与集群相同或更高的 3.x/4.x 客户端版本即可。

5.2 第一个生产者

import org.apache.kafka.clients.producer.*;
import java.util.Properties;

public class SimpleProducer {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
                  "org.apache.kafka.common.serialization.StringSerializer");
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
                  "org.apache.kafka.common.serialization.StringSerializer");

        try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {
            for (int i = 0; i < 10; i++) {
                ProducerRecord<String, String> record =
                        new ProducerRecord<>("order-events", "user-" + i % 3, "value-" + i);
                producer.send(record);
            }
            producer.flush();
        }
    }
}

try-with-resources 关闭时会隐式 flush(),确保缓冲区里的消息发出。注意:send() 是异步的,直接退出程序可能丢消息,这是新手最常踩的坑。

5.3 发送流程全景

                    main 线程                          sender IO 线程
┌──────────┐   ┌──────────────┐   ┌────────────┐   ┌──────────────┐
│ 拦截器    │ -> │ 序列化 key/value │ -> │ 分区器     │ -> │ RecordAccumulator │
│ (可选)    │   └──────────────┘   │ 选择分区    │   │ 按 (topic,分区)   │
└──────────┘                      │            │   │ 组织批次 batch     │
                                  └────────────┘   └───────┬──────┘
                                                           │ 满足条件
                                                           v
                                                    ┌──────────────┐
                                                    │ Sender 线程   │
                                                    │ 按 Broker 分组 │
                                                    │ 发送 Produce  │
                                                    └──────────────┘

几个要点:

5.4 异步回调与异常处理

生产环境必须写回调,否则发送失败你根本不知道:

producer.send(record, (metadata, exception) -> {
    if (exception == null) {
        System.out.printf("partition=%d offset=%d%n",
                metadata.partition(), metadata.offset());
    } else {
        // 记录日志、落补偿表、触发告警,绝不能静默吞掉
        exception.printStackTrace();
    }
});

同步发送(性能差,仅特定场景):

try {
    RecordMetadata md = producer.send(record).get();
} catch (Exception e) {
    // 处理执行异常
}

异常分两类:

5.5 分区策略

发送一条消息时,目标分区的决策顺序:

  1. ProducerRecord 指定了 partition,直接用;
  2. 否则若 key 不为空,对 key 做 murmur2 哈希 后对分区数取模: partition = Utils.toPositive(murmur2(keyBytes)) % numPartitions
  3. key 为空时使用粘性分区(sticky partitioning):先随机选一个分区尽量填满批次,满了再换下一个,让无 key 消息也能均匀分布且保持较高批效率。

自定义分区器示例——按用户等级路由:

public class UserLevelPartitioner implements Partitioner {
    @Override
    public int partition(String topic, Object key, byte[] keyBytes,
                         Object value, byte[] valueBytes, Cluster cluster) {
        int numPartitions = cluster.partitionsForTopic(topic).size();
        String k = key.toString();
        if (k.startsWith("vip-")) return 0;          // VIP 固定走 0 号分区
        return Math.abs(k.hashCode()) % (numPartitions - 1) + 1; // 其他走剩余分区
    }

    @Override
    public void close() {}

    @Override
    public void configure(Map<String, ?> configs) {}
}

注册:

props.put(ProducerConfig.PARTITIONER_CLASS_CONFIG, UserLevelPartitioner.class.getName());

注意:把 VIP 单独塞进一个分区会破坏负载均衡,这里只是演示语法。实际业务更常见的做法是“按实体 ID 做 key,均匀哈希”。

5.6 拦截器

public class TimingInterceptor implements ProducerInterceptor<String, String> {
    @Override
    public ProducerRecord<String, String> onSend(ProducerRecord<String, String> record) {
        record.headers().add("sent-at", String.valueOf(System.currentTimeMillis()).getBytes());
        return record;
    }

    @Override
    public void onAcknowledgement(RecordMetadata metadata, Exception exception) {
        // 统计成功率、耗时,注意这里在 IO 线程执行,必须快
    }

    @Override public void configure(Map<String, ?> configs) {}
    @Override public void close() {}
}

典型用途:打链路追踪头、统一加时间戳、统计指标。不适合做业务校验或慢 IO。

5.7 核心参数详解

参数 默认值 说明
bootstrap.servers 初始连接地址,写 2-3 台即可,客户端会自动发现全量 Broker
acks all 0 不等确认;1 Leader 写入即成功;all 等 ISR 确认,最可靠
retries Integer.MAX_VALUE 重试次数(配合 delivery.timeout.ms 兜底)
delivery.timeout.ms 120000 消息从发送到成功/失败的最终期限,含重试与发送时间
linger.ms 0 攒批等待时间。调大到 5-20ms 常能显著提升吞吐
batch.size 16384 单批字节数。过小则请求碎片化,过大占内存
buffer.memory 33554432 生产者总缓冲。满时 send() 会阻塞(最多 max.block.ms
compression.type none lz4/zstd 兼顾速度与压缩比,gzip CPU 换带宽
max.request.size 1048576 单请求上限,单条大消息需同时调大 Broker 端 message.max.bytes
enable.idempotence true 幂等生产者,防止重试导致的 broker 端重复
max.in.flight.requests.per.connection 5 单连接未确认请求数;幂等开启且 <=5 时重试仍保持分区内有序

acks 详解

注意:acks=all 在 ISR 只剩 Leader 自己时退化为 acks=1min.insync.replicas=2 才能真正保证至少两个副本落盘。

5.8 顺序性保证

Kafka 只保证分区内有序。要保证业务顺序:

  1. 需要顺序的实体使用相同 key(如同一订单号);
  2. 生产端开启幂等(默认已开),且 max.in.flight.requests.per.connection <= 5
  3. 不要在发送失败后自行乱序重试(例如把失败消息丢进另一个线程重新 send,而后续消息已发出);
  4. 增加分区会改变 key 的哈希落点,历史顺序与新增分区间可能断开(第 7 章)。

5.9 一个生产级配置模板

Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "broker1:9092,broker2:9092");
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());

props.put(ProducerConfig.ACKS_CONFIG, "all");
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
props.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, 120000);
props.put(ProducerConfig.LINGER_MS_CONFIG, 10);
props.put(ProducerConfig.BATCH_SIZE_CONFIG, 32 * 1024);
props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "zstd");

这套配置适合大多数“可靠性优先、吞吐要求高”的业务场景。日志类场景可以更激进地调大 linger.msbatch.size

本章小结

思考题

  1. 为什么 send() 返回成功不代表消息“安全”了?它与 acks 是什么关系?
  2. linger.ms=0linger.ms=20 在延迟、吞吐、CPU 上分别会带来什么变化?
  3. 如果业务要求“所有订单消息严格有序”,分区数必须是多少?这样做有什么代价?