调优 Kafka Producer 吞吐量

Producer 吞吐由批次形成、网络往返、压缩成本和确认要求共同决定。等待更多记录合批可以减少请求数,但增加延迟;降低确认数可能提速,却会削弱故障时的数据保证。参数应同时按吞吐、尾延迟、内存、CPU 与容灾目标评估。

Producer 将记录聚成批次、压缩发送到 broker,acks 决定等待的确认级别。
未完纪原创流程图。

四个关键参数

  • batch.size:每分区批次目标上限。原文建议试验 100,000–200,000 字节(示例默认 16,384);大批次可减少请求但占用缓冲,低流量分区也许无法填满。
  • linger.ms:批次未达到上限时最多等待多久。源教程建议 10–100 毫秒(页面记录默认 0)。等待可能提高吞吐,也会增加排队延迟;按实际客户端版本核对默认。
  • compression.type:示例采用 lz4,可减小网络负载但消耗 CPU;受数据可压缩性与客户端/broker 版本影响。
  • acks:决定何时确认写入。教程设置 acks=1,只等 leader 确认;相较 acks=all,leader 故障时尚未复制的数据可能丢失。不能只为跑分快速套用。

本地 Docker 流程

原文用 Confluent 教程仓库的 Compose 环境启动 broker,运行 kafka-producer-perf-test:对 topic-perf 写 10,000 条、每条 8,000 字节、throughput -1 表示不限速。页面报告基线约 23.58 MB/s;设批次 200,000、linger.ms=100、LZ4、acks=1 后约 94.89 MB/s。

这两个数字是源文特定机器、容器网络与 broker 配置的观测,不是本稿复测或性能保证。可重复评估须维持相同负载,多次运行,记录延迟分位数、资源和磁盘指标,逐项调参;不要把降低确认保证与压缩效果合并成一个结论。

Confluent Cloud 流程

教程先在 Cloud 控制台创建集群与 topic-perf,生成含 API key 的 Java 配置,保存在本地 cloud.properties,再挂载给容器性能工具。源文容器镜像是 confluentinc/cp-server:7.5.1,报告基线约 9.28 MB/s、调优后 15.24 MB/s。版本和服务负载可能变化,这些数字不是 Cloud SLA。

密钥与版本安全

cloud.properties 含 API key/secret,应限制权限、加入忽略规则、不得提交 Git、写日志或打进镜像;用只读挂载更稳妥,测试后按组织策略轮换/撤销临时 key。教程中 linger 默认值、镜像和吞吐均有版本边界,部署前查目标客户端实际配置。本文没有克隆仓库、运行 Docker、接触 Cloud 或复测数字。

来源:Confluent Developer, How to optimize a Kafka producer for throughput。

本文补充:核对与维护细节

把吞吐实验做成可解释的比较

本地教程先从 Confluent tutorials 仓库启动 Compose broker 环境,再使用 kafka-producer-perf-test 生成固定数量和尺寸的记录。参数 --num-records 10000、--record-size 8000 和 --throughput -1 固定负载以比较两轮;配置文件分别包含 broker 与 Producer 参数。必须确认每轮写入相同 topic、partition 数和消息尺寸,否则 MB/s 不能作直接对比。

Producer 吞吐受分区并行度、broker 磁盘、网络与副本策略限制。若调大 batch 后吞吐没升,可能因为缓冲/分区或 broker 已饱和;若 linger 增长后 p99 延迟上升,则需调低等待或改异步负载。压缩率和 CPU 使用要一起看;高压缩并不一定适合 CPU 紧张的 broker。基准结果应包含消息丢失容忍度与延迟目标,而不只看单个吞吐指标。

Cloud 配置文件的生命周期

教程要求从控制台生成客户端配置,并将密钥放到 cloud.properties 再挂载给容器。运行前检查容器挂载是主机文件到容器路径,宿主机路径正确且仅当前测试可读;使用结束后删除文件、清理终端历史和日志中的可能回显,并撤销临时凭据。容器映像版本 7.5.1 应作为示例来源信息保留,而不是盲目升级/降级到与 Cloud 集群不兼容的 tag。

复现教程中的四轮基准

以下列出来源中的本地演示命令和结果,供核对;本稿没有运行这些命令。它们会启动 Docker 并写入 Kafka topic,执行前须检查本机资源与 topic 状态。仓库入口为 git clone git@github.com:confluentinc/tutorials.git,随后进入 tutorials:

docker compose -f ./docker/docker-compose-kafka.yml up -d
docker cp optimize-producer-throughput/kafka/local.properties broker:/etc/producer.properties
docker exec broker /usr/bin/kafka-producer-perf-test --topic topic-perf --num-records 10000 --record-size 8000 --throughput -1 --producer.config /etc/producer.properties
docker exec broker /usr/bin/kafka-producer-perf-test --topic topic-perf --num-records 10000 --record-size 8000 --throughput -1 --producer.config /etc/producer.properties --producer-props batch.size=200000 linger.ms=100 compression.type=lz4 acks=1
源文轮次 记录/秒 MB/s 平均延迟 p95 / p99
本地基线 3,091.19 23.58 927.27 ms 1,302 / 1,352 ms
本地调优 12,437.81 94.89 4.92 ms 16 / 38 ms

源文还报告本地基线最大延迟 1,362 ms、p50 949 ms、p99.9 1,360 ms;调优轮最大 378 ms、p50 3 ms、p99.9 43 ms。特定容器、硬件和 broker 的一次结果不能代表所有集群,特别要把吞吐与 p99、CPU、压缩率、副本确认保证一起评价。

Confluent Cloud 原始轮次

教程先在控制台建立集群与 topic-perf,生成 Java 客户端配置存入 optimize-producer-throughput/kafka/cloud.properties。文件含 API key/secret,不要提交 Git、打印或烘焙进镜像;本文未取得任何真实密钥。源文使用镜像 confluentinc/cp-server:7.5.1:

docker run -v ./optimize-producer-throughput/kafka/cloud.properties:/etc/producer.properties confluentinc/cp-server:7.5.1 /usr/bin/kafka-producer-perf-test --topic topic-perf --num-records 10000 --record-size 8000 --throughput -1 --producer.config /etc/producer.properties
docker run -v ./optimize-producer-throughput/kafka/cloud.properties:/etc/producer.properties confluentinc/cp-server:7.5.1 /usr/bin/kafka-producer-perf-test --topic topic-perf --num-records 10000 --record-size 8000 --throughput -1 --producer.config /etc/producer.properties --producer-props batch.size=200000 linger.ms=100 compression.type=lz4 acks=1
源文轮次 记录/秒 MB/s 平均延迟 p95 / p99
Cloud 基线 1,216.25 9.28 2,213.16 ms 3,550 / 4,344 ms
Cloud 调优 1,997.60 15.24 1,172.76 ms 1,547 / 1,701 ms

Cloud 两轮的最大延迟分别为 4,665/1,901 ms,p50 为 2,098/1,246 ms,p99.9 为 4,640/1,901 ms。服务区、配额、网络与后台负载均不同于本地 Docker,不能横向比较绝对数值;教程报告不是本稿测试。压缩、linger 与 acks=1 应分开测量,且需按实际客户端版本检查默认值与服务可靠性目标。

环境、控制台步骤与原版权

原文要求可用的Docker Desktop或Docker Engine;本地路径还需要Docker Compose,并以 docker compose version 检查安装。克隆后进入 tutorials 目录。源文记录compression.type默认none;acks默认all自Apache Kafka 3.0起。以上为源文版本信息,仍须按实际客户端核对。

Cloud分支需要Confluent Cloud账户:在环境页面选择Add cluster;集群运行后,从左侧Topics创建采用默认配置的topic-perf;在Cluster Overview的Clients中选择Java,生成含API keys的配置并保存到指定cloud.properties路径。创建服务及写入压测可能产生费用。本文没有创建集群或生成凭据。

Copyright © Confluent, Inc. 2014–2026. Apache、Apache Kafka、Kafka、Apache Flink、Flink、Apache Iceberg、Iceberg及相关开源项目名称是Apache Software Foundation的商标。中文翻译及风险注释:未完纪;原示例观测以表格呈现并适当四舍五入,不是重新测试结果。

© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容