这是《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 │
└──────────────┘
几个要点:
- 消息先经过拦截器
onSend,再序列化,再选分区; 顺序是固定的,自定义拦截器不要在这里做重逻辑; - 同一个
(topic, partition)的消息在RecordAccumulator中攒成批次; 批次满了(batch.size)或等待超时(linger.ms)就发出; Sender线程按 Broker 维度合并请求,一个请求可以携带多个分区的批次;- 返回结果通过回调交给业务线程。
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) {
// 处理执行异常
}
异常分两类:
- 可重试:
RetriableException的子类,如TimeoutException、NotEnoughReplicasException。客户端按retries/delivery.timeout.ms自动重试; - 不可重试:如
RecordTooLargeException(单条超过max.request.size)、序列化失败、AuthorizationException。重试无意义,必须修代码或配置。
5.5 分区策略
发送一条消息时,目标分区的决策顺序:
ProducerRecord指定了partition,直接用;- 否则若 key 不为空,对 key 做 murmur2 哈希 后对分区数取模:
partition = Utils.toPositive(murmur2(keyBytes)) % numPartitions; - 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=0:发出去就算成功。吞吐最高,Broker 宕机、网络丢包都会直接丢消息;acks=1:Leader 本地写入成功即返回。Leader 刚确认就宕机且未同步到 Follower 时会丢;acks=all:等待min.insync.replicas个 ISR 副本确认。配合replication.factor>=3、min.insync.replicas>=2是可靠性的黄金组合。
注意:acks=all 在 ISR 只剩 Leader 自己时退化为 acks=1。min.insync.replicas=2 才能真正保证至少两个副本落盘。
5.8 顺序性保证
Kafka 只保证分区内有序。要保证业务顺序:
- 需要顺序的实体使用相同 key(如同一订单号);
- 生产端开启幂等(默认已开),且
max.in.flight.requests.per.connection <= 5; - 不要在发送失败后自行乱序重试(例如把失败消息丢进另一个线程重新 send,而后续消息已发出);
- 增加分区会改变 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.ms 和 batch.size。
本章小结
send()异步进入累加器,由 Sender 线程攒批发送;退出前必须flush()/close();- 分区决策顺序:显式指定 > key 哈希 > 粘性分区;
- 生产环境必须实现回调,区分可重试与不可重试异常;
acks=all + min.insync.replicas>=2是不丢消息的底线组合;- 顺序性的三要素:同 key 同分区、幂等开启、in-flight <= 5。
思考题
- 为什么
send()返回成功不代表消息“安全”了?它与acks是什么关系? linger.ms=0和linger.ms=20在延迟、吞吐、CPU 上分别会带来什么变化?- 如果业务要求“所有订单消息严格有序”,分区数必须是多少?这样做有什么代价?