乐途乐途
主页
  • 计算机基础

    • TCP/IP
    • Linux
    • HTTP
  • 数据库

    • SQL
    • MySQL 5.7
  • 编程语言

    • C
    • C++
    • Java SE
    • Python2
    • Python3
  • 数据格式

    • JSON
    • XML
  • 认证与安全

    • JWT
  • 工具

    • Markdown
  • Git

    • GitFlow
  • Quartz

    • Quartz
  • Java

    • Maven 入门
    • Maven 进阶
    • MyBatis
    • Spring
    • Spring MVC
  • Java

    • Spring Boot
    • Spring Cloud
    • Spring Cloud Alibaba
    • Spring Security
    • Spring AI
    • Spring Batch
    • Kafka
    • Java 设计模式
  • 缓存

    • Redis
  • 搜索引擎

    • Elasticsearch
  • 分布式协调

    • ZooKeeper
联系
阿里云
主页
  • 计算机基础

    • TCP/IP
    • Linux
    • HTTP
  • 数据库

    • SQL
    • MySQL 5.7
  • 编程语言

    • C
    • C++
    • Java SE
    • Python2
    • Python3
  • 数据格式

    • JSON
    • XML
  • 认证与安全

    • JWT
  • 工具

    • Markdown
  • Git

    • GitFlow
  • Quartz

    • Quartz
  • Java

    • Maven 入门
    • Maven 进阶
    • MyBatis
    • Spring
    • Spring MVC
  • Java

    • Spring Boot
    • Spring Cloud
    • Spring Cloud Alibaba
    • Spring Security
    • Spring AI
    • Spring Batch
    • Kafka
    • Java 设计模式
  • 缓存

    • Redis
  • 搜索引擎

    • Elasticsearch
  • 分布式协调

    • ZooKeeper
联系
阿里云
  • 学习路径
  • 第1章 消息队列与 Kafka 概述

    • 章节导读
    • 消息队列基础
    • Kafka 概述
    • Kafka 为什么快
  • 第2章 快速上手:单机环境搭建

    • 章节导读
    • 环境准备与安装
    • 快速启动
    • Topic 管理
  • 第3章 核心概念:主题、分区与日志

    • 章节导读
    • Topic
    • Partition
    • Offset
    • Segment 与存储结构
  • 第4章 生产者详解

    • 章节导读
    • Producer 概述
    • Producer 发送机制
    • Producer 分区策略
    • Producer 幂等与事务
    • Producer 配置
  • 第5章 消费者与消费组

    • 章节导读
    • Consumer 概述
    • Consumer Group 消费组
    • Consumer 分区分配策略
    • Consumer Offset 提交
    • Consumer 多线程
    • Consumer 配置
  • 第6章 Broker 与控制器

    • 章节导读
    • Broker 概述
    • Broker 配置
    • Controller
    • KRaft 共识协议
  • 第7章 副本与数据可靠性

    • 章节导读
    • 副本机制
    • ISR 与副本同步
    • Leader 选举
    • ACK 与一致性保证
    • 高水位与 Leader Epoch
  • 第8章 存储与性能优化

    • 章节导读
    • 存储架构
    • 日志清理与压缩
    • Page Cache 与零拷贝
    • Producer 性能优化
    • Consumer 性能优化
    • Broker 性能优化
  • 第9章 生产环境运维与监控

    • 章节导读
    • Topic 管理
    • Kafka 运维工具
    • 监控
    • 常见故障排查
  • 第10章 Kafka生态与面试考点

    • 章节导读
    • Kafka 生态全景
    • 面试高频 30 题

Producer 配置

定义与作用

本节提供 Producer 关键配置参数的速查表与调优指南。不罗列所有参数,聚焦于影响吞吐、延迟、可靠性的核心参数及其组合策略。

核心原理

Producer 的配置参数关系着三个维度的权衡:

配置速查表

核心连接与发送

参数默认值说明调优建议
bootstrap.servers(必填)Broker 地址列表填写多个,防止单点连接失败
key.serializer(必填)Key 序列化器StringSerializer / ByteArraySerializer
value.serializer(必填)Value 序列化器同上
client.id""客户端标识建议设置,便于日志追踪

批处理与吞吐

参数默认值说明调优建议
batch.size16384 (16KB)每个分区的批次大小吞吐优化往上调至 32KB-512KB
linger.ms0批次等待时间(ms)吞吐优化设 5-100;延迟敏感保持 0
buffer.memory33554432 (32MB)缓冲区总大小高吞吐场景调至 64-256MB
compression.typenone压缩类型(none/gzip/snappy/lz4/zstd)推荐 lz4(吞吐/CPU 比最优)或 zstd(高压缩率)
max.request.size1048576 (1MB)单次请求最大字节大数据场景上调,需同步调整 Broker message.max.bytes

可靠性

参数默认值说明调优建议
acksall(3.0+)确认级别关键数据必须 all
enable.idempotencetrue(3.0+)幂等生产者建议保持开启
retries2147483647重试次数与 delivery.timeout.ms 配合
delivery.timeout.ms120000 (2min)发送超时(含重试)根据 SLA 调整
retry.backoff.ms100重试间隔(ms)网络抖动频繁可适当增大
max.in.flight.requests.per.connection5未确认请求数幂等时自动 5;不幂等需有序时设为 1

事务

参数默认值说明
transactional.idnull事务 ID(非空即开启事务模式)
transaction.timeout.ms60000 (1min)事务超时

完整示例

示例一:高吞吐日志收集配置

场景:日志收集系统,每秒 50 万条小消息(每条 ~200B),可容忍少量丢失。

Properties props = new Properties();
props.put("bootstrap.servers", "kafka1:9092,kafka2:9092,kafka3:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");

// 吞吐优先
props.put("acks", "1");
props.put("linger.ms", 10);
props.put("batch.size", 131072);       // 128KB
props.put("buffer.memory", 134217728); // 128MB
props.put("compression.type", "lz4");
props.put("max.in.flight.requests.per.connection", 5);

操作前后对比:

配置默认配置高吞吐配置
linger.ms010
batch.size16KB128KB
compression.typenonelz4
实测吞吐~50K/s~200K/s

示例二:交易系统可靠性优先配置

场景:交易系统,要求消息绝对不丢失、不重复、不乱序。

Properties props = new Properties();
props.put("bootstrap.servers", "kafka1:9092,kafka2:9092,kafka3:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");

// 可靠性优先
props.put("enable.idempotence", true);          // acks=all, retries=MAX, max.in.flight=5
props.put("delivery.timeout.ms", 120000);       // 2 分钟超时
props.put("transactional.id", "tx-trade-001");  // 开启事务(如需跨 Topic 原子性)

注意:可靠性优先配置配合 Broker 端:

  • Topic 级别:min.insync.replicas >= 2
  • Broker 级别:unclean.leader.election.enable=false

易错场景

易错 1:compression.type 设为 gzip 导致延迟飙升

场景:延迟敏感系统(P99 < 10ms)使用 gzip 压缩。

后果:gzip CPU 消耗高,压缩/解压时间可能超过网络传输时间,延迟不降反升。

建议:

  • 延迟优先:snappy 或 none
  • 吞吐优先:lz4(CPU/压缩比最佳平衡)
  • 存储优先:zstd(最高压缩比)

易错 2:buffer.memory 太小导致 block.on.buffer.full

场景:高吞吐场景下 buffer.memory=32MB(默认)。

后果:缓冲区满后 send() 阻塞,等待空间释放(受 max.block.ms 控制,默认 60s)。

诊断:

// 监控 Producer 指标
// bufferpool-wait-time: 等待缓冲区的时间(应为 0)
// buffer-available-bytes: 缓冲区可用字节(不应经常为 0)

面试高频考点

Q:linger.ms=0 时 batch.size 还有意义吗?

A:有意义。linger.ms 和 batch.size 是"或"的关系——哪个条件先达到就发送:

  1. linger.ms=0 意味着到达即发
  2. 但如果前一个请求还在传输中,新消息会继续积累
  3. 积累到 batch.size 上限时,即使 linger.ms 未到也会发送

所以 linger.ms=0 + 高并发场景下,批次仍然会填满 batch.size(因为总有消息在排队等待前一个请求完成)。linger.ms 的主要作用是在低负载场景下多等一会让批次填满。

上一页
Producer 幂等与事务