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

    • 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 达到 100K-1M 条/秒的吞吐。

核心原理

Producer 内部流水线

两个关键缓冲:

  • RecordAccumulator(buffer.memory,默认 32MB):暂存待发送的消息
  • 每分区 Deque<ProducerBatch>:等待攒满或超时

性能三要素

配置速查表

参数默认值吞吐优先延迟优先说明
batch.size16384 (16KB)65536-1310720-512攒批大小
linger.ms010-1000等多久再发
buffer.memory33554432 (32MB)67108864+—总缓冲区
compression.typenonelz4 / zstdnone压缩算法
max.in.flight.requests.per.connection551并发请求数
acksall (3.x)10确认级别
enable.idempotencetrue (3.x)truefalse幂等

完整示例

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

Properties props = new Properties();
props.put("bootstrap.servers", "kafka1:9092,kafka2:9092");
props.put("key.serializer", "...");
props.put("value.serializer", "...");

// 高吞吐核心
props.put("batch.size", 131072);           // 128KB 攒批
props.put("linger.ms", 50);                // 等 50ms
props.put("compression.type", "lz4");      // LZ4 压缩
props.put("buffer.memory", 134217728);     // 128MB
props.put("max.in.flight.requests.per.connection", 5);

测试结果:

bin/kafka-producer-perf-test.sh --topic perf-test --num-records 5000000 \
  --record-size 100 --throughput -1 --producer-props \
  bootstrap.servers=localhost:9092 batch.size=131072 linger.ms=50 compression.type=lz4

# 结果
# 5000000 records sent, 250000 records/sec, 23.84 MB/sec
# avg latency: 15 ms, max latency: 120 ms

示例二:低延迟配置(实时风控)

Properties props = new Properties();
props.put("bootstrap.servers", "kafka1:9092,kafka2:9092");

// 低延迟核心
props.put("batch.size", 0);                // 不攒批
props.put("linger.ms", 0);                 // 立即发送
props.put("acks", "1");                    // 快速确认
props.put("max.in.flight.requests.per.connection", 1); // 保序

测试结果:

# 低延迟配置
bin/kafka-producer-perf-test.sh --topic perf-test --num-records 100000 \
  --record-size 100 --throughput -1 --producer-props \
  bootstrap.servers=localhost:9092 batch.size=0 linger.ms=0 acks=1

# 结果
# 100000 records sent, 12000 records/sec, 1.14 MB/sec
# avg latency: 2 ms, max latency: 8 ms

操作前后对比:

指标默认配置高吞吐配置低延迟配置
每秒消息数~50K~250K~12K
平均延迟~30ms~15ms~2ms
P99 延迟~200ms~120ms~8ms
网络效率中高(压缩+攒批)低(频繁小包)

易错场景

易错 1:linger.ms=0 + 极高写入速率导致"小包轰炸"

场景:linger.ms=0 但每秒写入 100 万条。

后果:每条消息一个独立的 ProduceRequest → 大量小数据包 → 网络带宽利用率 < 30%,CPU 大量用于网络中断处理。

解决:吞吐场景永远设置 linger.ms=5 以上。

易错 2:buffer.memory 太小导致 send() 阻塞

场景:buffer.memory=32MB(默认),写入速率 > Broker 处理速率。

后果:RecordAccumulator 被写满 → send() 阻塞最长达 max.block.ms(默认 60s)→ 应用线程卡住。

诊断:

// 监控缓冲区使用率
float usage = (float) producer.metrics()
    .get("buffer-total-bytes-used").metricValue()
    / producer.metrics().get("buffer-total-bytes-max").metricValue();

面试高频考点

Q:batch.size 和 linger.ms 的关系是什么?谁先触发就发送?

A:

  • batch.size 是空间触发:攒够指定字节就发
  • linger.ms 是时间触发:超过指定毫秒就发
  • 先满足谁就触发:如果 5ms 攒够了 128KB → 以 batch.size 触发;如果 10ms 还没攒够 128KB → 以 linger.ms 触发

实际调优中,linger.ms 是更主要的控制变量——它保证了延迟的确定性上限。

上一页
Page Cache 与零拷贝
下一页
Consumer 性能优化