Producer 概述
定义与作用
Kafka Producer 是向 Kafka Topic 发送消息的客户端程序。它的职责不限于"把数据发出去"——它承担了分区路由、序列化、压缩、批处理攒批、重试容错等一系列复杂职责,而将这些复杂性对应用开发者透明化。
Producer 在 Kafka 架构中的位置:数据入口,是所有数据管道的第一步。
核心原理
Producer 的内部架构
双线程设计:
| 组件 | 线程 | 职责 |
|---|---|---|
KafkaProducer(主线程) | 应用程序线程 | 拦截 → 序列化 → 分区 → 放入 RecordAccumulator |
Sender(后台线程) | 独立 I/O 线程 | 从 RecordAccumulator 取批次 → 发送到 Broker → 处理 ACK |
ProducerRecord 消息结构
ProducerRecord<String, String> record = new ProducerRecord<>(
"orders", // topic: 目标 Topic
0, // partition: 指定分区(可选)
1718208000000L, // timestamp: 时间戳(可选)
"order-1001", // key: 消息 Key(可选,用于分区路由)
"{\"amount\":99}" // value: 消息体
);
| 字段 | 必选 | 用途 |
|---|---|---|
topic | 是 | 目标 Topic |
value | 是 | 消息负载(payload) |
key | 否 | 用于分区路由,相同 Key 的消息进入同一分区 |
partition | 否 | 显式指定分区,优先级高于 Key 路由 |
timestamp | 否 | 消息时间戳,不指定则使用 System.currentTimeMillis() |
headers | 否 | 键值对形式的自定义元数据(可用于链路追踪等) |
完整示例
示例一:最小可用的 Producer
场景:向 orders Topic 发送一条订单消息。
操作前环境:Kafka 已启动,Topic orders 已创建(1 分区)。
import org.apache.kafka.clients.producer.*;
import java.util.Properties;
public class SimpleProducer {
public static void main(String[] args) throws Exception {
// 1. 配置 Producer
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
// 2. 创建 Producer 实例
try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {
// 3. 构造消息
ProducerRecord<String, String> record = new ProducerRecord<>(
"orders", "order-1001", "{\"amount\": 99.90, \"item\": \"book\"}"
);
// 4. 发送(异步)
producer.send(record, (metadata, exception) -> {
if (exception == null) {
System.out.printf("发送成功 → topic=%s, partition=%d, offset=%d%n",
metadata.topic(), metadata.partition(), metadata.offset());
} else {
exception.printStackTrace();
}
});
}
// 5. try-with-resources 自动调用 close(),等待缓冲区消息全部发送
}
}
执行结果:
发送成功 → topic=orders, partition=0, offset=0
操作后状态:
| 指标 | 发送前 | 发送后 |
|---|---|---|
orders Partition 0 的 LEO | 0 | 1 |
| 消息内容 | — | {"amount": 99.90, "item": "book"} |
示例二:同步发送 vs 异步发送
场景:对比 get() 同步阻塞和 Callback 异步两种方式。
// ===== 方式一:同步发送(阻塞等待结果) =====
try {
RecordMetadata metadata = producer.send(record).get();
System.out.println("同步发送成功, offset=" + metadata.offset());
} catch (Exception e) {
System.err.println("同步发送失败: " + e.getMessage());
}
// ===== 方式二:异步发送(Callback 非阻塞) =====
producer.send(record, new Callback() {
@Override
public void onCompletion(RecordMetadata metadata, Exception e) {
if (e != null) {
System.err.println("异步发送失败: " + e.getMessage());
} else {
System.out.println("异步发送成功, offset=" + metadata.offset());
}
}
});
操作前后对比:
| 方式 | 吞吐量 | 延迟 | 适用场景 |
|---|---|---|---|
同步 get() | 低(~100 QPS) | 每条 ~1-5ms | 强顺序要求、少量消息 |
| 异步 Callback | 高(~100K QPS) | 应用层几乎零等待 | 高吞吐场景(默认选择) |
注意:异步方式下,producer.close() 会阻塞等待所有未完成的发送确认,所以上例中 try-with-resources 保证了消息不会因进程退出而丢失。
易错场景
易错 1:忘记调用 close()
现象:main 方法结束后 JVM 退出,发现 Producer 已发送的消息在 Kafka 中丢失。
原因:消息在 RecordAccumulator 的缓冲区中尚未被 Sender 线程实际发送到 Broker。
解决:始终使用 try-with-resources 或 finally { producer.close(); },close() 会等待缓冲区清空。
易错 2:每条消息创建新的 Producer 实例
// 错误:每条消息创建 Producer(开销巨大)
for (int i = 0; i < 10000; i++) {
KafkaProducer<String, String> p = new KafkaProducer<>(props);
p.send(new ProducerRecord<>("topic", "msg" + i));
p.close();
}
后果:每次 new KafkaProducer 都会创建 Sender 线程、建立 TCP 连接、获取元数据,性能极差。
正确做法:Producer 是线程安全的,一个进程通常只需一个 Producer 实例。
面试高频考点
Q:Kafka Producer 为什么设计为异步发送?
A:异步发送 + 批处理是 Kafka 高吞吐的核心设计之一:
- 解耦应用线程和 I/O:应用线程放入缓冲区即返回,Sender 线程独立管理网络 I/O
- 批处理攒批:
RecordAccumulator按分区聚合消息,将小消息合并为大网络包,减少网络往返 - 压缩:批次级别的压缩远优于逐条压缩
- 容错:Sender 线程可独立处理重试,不阻塞业务线程