Learn
Kafka/05-producer-basics

Producer 原理与核心配置

producer.send() 一行代码背后,消息经历了拦截器、序列化、分区路由、内存攒批、网络发送五个阶段。理解这条链路,你才能解释"为什么消息没有立刻发出去"、"为什么内存爆了"、"batch.size 和 linger.ms 到底怎么配"。

1. 发送流程全景

Producer 内部是双线程模型:主线程负责把消息放进缓冲区,Sender 线程负责真正的网络 IO。

主线程 (你的业务线程)
  send(record)
    |
    v
  [拦截器 Interceptors]        onSend() 加工消息 (打标、埋点)
    |
    v
  [序列化器 Serializer]        key/value 对象 -> byte[]
    |
    v
  [分区器 Partitioner]         决定去哪个分区 (有 key 按哈希)
    |
    v
  [RecordAccumulator 缓冲区]   按 "topic-分区" 维度攒批
    |  buffer.memory=32MB       batch1[msg,msg,..] batch2[..]
    |                                 |
- - | - - - - - - - - - - - - - - - - | - - - - - -
    v                                 v
Sender 线程 (IO 线程)          批次满(batch.size) 或到时(linger.ms)
  把就绪批次组装成请求 ──> 网络发送到分区 leader 所在 broker
  等待 acks 响应 ──> 成功: 执行回调 / 失败: 重试或抛错

关键点:

  • send() 返回时消息只是进了本地缓冲区,并没有发出去。
  • 一个批次(batch)只属于一个分区,发往同一 broker 的多个批次会合并成一个请求。
  • Sender 线程异步工作,主线程几乎无阻塞——除非缓冲区满。

2. 攒批三参数

参数默认值含义
batch.size16384 (16KB)单个批次的字节上限,攒满立即可发
linger.ms0批次未满时最多等多久,到时也发
buffer.memory33554432 (32MB)整个 Producer 缓冲区总大小

发送时机 = 批次攒满 batch.size 或 等待达到 linger.ms,先到先发。

  • linger.ms=0(默认):消息尽快发出,延迟最低,但批次很小、吞吐差。
  • linger.ms=5~100:故意等一等换取更大的批次,吞吐大幅提升,延迟增加几毫秒——高吞吐场景必调。
  • 缓冲区满时 send() 会阻塞,最长 max.block.ms(默认 60 秒),超时抛异常。
linger.ms=0:    [m][m][m][m]         每条几乎单独一个请求, 请求数多
linger.ms=20:   [m m m m m m m][...] 攒成大批次, 请求数骤减, 吞吐up

3. 压缩

compression.type 对整个批次压缩,批次越大压缩比越高(这也是要攒批的原因之一):

算法压缩比CPU 开销建议
none-无默认
lz4中低通用首选
snappy中低与 lz4 接近
zstd高中带宽敏感首选(2.1+)
gzip高高不推荐新项目

4. Java 客户端实战

依赖 org.apache.kafka:kafka-clients:3.7.1。

import org.apache.kafka.clients.producer.*;
import org.apache.kafka.common.serialization.StringSerializer;
import java.util.Properties;
import java.util.concurrent.Future;
 
public class OrderProducer {
    public static void main(String[] args) throws Exception {
        Properties props = new Properties();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        // 攒批与压缩
        props.put(ProducerConfig.BATCH_SIZE_CONFIG, 32768);        // 32KB
        props.put(ProducerConfig.LINGER_MS_CONFIG, 20);            // 等 20ms
        props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "lz4");
        props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 64 * 1024 * 1024L);
 
        Producer<String, String> producer = new KafkaProducer<>(props);
        ProducerRecord<String, String> record =
            new ProducerRecord<>("orders", "user-1", "{\"orderId\":1001}");
 
        // 方式1: 发后即忘 —— 最快, 可能丢消息且无感知
        producer.send(record);
 
        // 方式2: 同步发送 —— 逐条确认, 吞吐最低
        RecordMetadata meta = producer.send(record).get();
        System.out.printf("sync -> partition=%d offset=%d%n",
            meta.partition(), meta.offset());
 
        // 方式3: 异步 + 回调 —— 生产推荐, 高吞吐且能感知失败
        producer.send(record, (metadata, exception) -> {
            if (exception != null) {
                // 记日志/告警/落补偿表, 不要吞掉
                exception.printStackTrace();
            } else {
                System.out.printf("async -> partition=%d offset=%d%n",
                    metadata.partition(), metadata.offset());
            }
        });
 
        producer.flush();   // 强制清空缓冲区
        producer.close();   // 会先 flush 再关闭
    }
}
⚠️Producer 是线程安全的, 全局一个就够

KafkaProducer 线程安全且很重——内部有缓冲区和 IO 线程。正确用法是整个应用共享一个实例。每次发送 new 一个 Producer 是最常见的性能反模式,会导致连接风暴且完全无法攒批。另外,进程退出前必须调用 close(),否则缓冲区里还没发出去的消息会直接丢失。

5. Python 对照(confluent-kafka)

from confluent_kafka import Producer
 
producer = Producer({
    "bootstrap.servers": "localhost:9092",
    "batch.size": 32768,
    "linger.ms": 20,
    "compression.type": "lz4",
})
 
def on_delivery(err, msg):
    if err is not None:
        print(f"delivery failed: {err}")
    else:
        print(f"delivered to {msg.topic()}[{msg.partition()}]@{msg.offset()}")
 
for i in range(10):
    producer.produce(
        topic="orders",
        key=f"user-{i % 3}",
        value=f'{{"orderId": {1000 + i}}}',
        on_delivery=on_delivery,
    )
    producer.poll(0)   # 触发已完成发送的回调
 
producer.flush()       # 退出前必须 flush

produce() 同样只是入队;poll(0) 驱动回调执行;flush() 阻塞直到全部送达。

6. 其他重要参数速查

参数默认说明
max.request.size1MB单个请求最大字节数,超大消息会被直接拒绝
request.timeout.ms30000等 broker 响应的超时
delivery.timeout.ms120000send() 到确认成功/失败的总时限(含重试)
max.block.ms60000缓冲区满或元数据不可用时 send() 最长阻塞时间
client.id空客户端标识,方便服务端日志与配额定位
💡吞吐调优三板斧

高吞吐场景先调这三个:linger.ms=20~100、batch.size=64KB~128KB、compression.type=lz4/zstd。通常能把吞吐提升 3-10 倍,代价只是几十毫秒的发送延迟。可靠性相关的 acks、retries、幂等在下一章专门讲。

小结

  • Producer 双线程模型:主线程写缓冲区,Sender 线程网络发送;send() 返回不等于已发送
  • 链路:拦截器 → 序列化 → 分区器 → RecordAccumulator 攒批 → Sender
  • 攒批由 batch.size 与 linger.ms 共同控制,先到先发;压缩作用于整个批次
  • 三种发送:发后即忘(慎用)、同步(低吞吐)、异步回调(生产推荐)
  • 全局共享一个 Producer 实例,退出前 close/flush
  • 下一章解决"怎么保证消息不丢不重" →
🎯练习
  1. 写一个循环发送 10 万条消息的程序,分别用 linger.ms=0 和 linger.ms=50 + lz4 压缩跑一遍,对比总耗时。
  2. 把 buffer.memory 调小到 1MB 并快速发送大消息,观察 send() 阻塞与 max.block.ms 超时异常。
  3. 用异步回调统计成功/失败条数,验证进程退出前不调 flush 会丢多少条。