ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Kafka消息发送过程全解析:从分区器到副本确认的完整链路

Kafka消息发送过程全解析:从分区器到副本确认的完整链路 Kafka 消息的发送过程这个题目我一向不把它当面试八股看。这些年线上排查 Kafka 问题十次有八次最后都绕回到发送链路消息延迟高、发送超时、InvalidReceiveException、顺序错乱根因基本都藏在从业务线程到 broker 落盘这条线上。很多同学把分区器、批次、ACK、副本同步这些名词背得滚瓜烂熟可真的遇到故障时仍然一头雾水就是因为缺少完整的链路视角。这篇文章我按一条真实消息的发送路线把每个环节、关键参数、设计权衡和我踩过的坑一起讲清楚。入门的读者可以当成系统化的发送链路原理课工作过一两年的老手也可以拿它来对照检查自己的参数设置和排查思路。1. 选型与前置认知发送链路为什么值得单独拆开讲1.1 发送过程的可靠性承诺先看消息队列的选型差异很多项目在做技术选型时习惯把 Kafka、RabbitMQ、RocketMQ 摆在一起比吞吐量、比社区活跃度但很少人注意到三个产品的发送链路上层结构非常相似都是客户端攒数据 - 发到 broker - broker 确认真正拉开差距的是它们对可靠性承诺和堆积处理的设计目标。维度KafkaRabbitMQRocketMQ定位分布式流平台核心是分区日志轻量消息代理AMQP 协议路由灵活分布式消息中间件Java 生态完善典型吞吐极高适合海量日志、事件流单 broker 可达数十万条/秒量级中高适合业务系统解耦消息积压后性能下降明显高常见万级到十万级/秒堆积能力强发送可靠性通过 acks、ISR、副本复制保证3.x 客户端默认开启幂等Publisher Confirm 机制可精确确认事务消息、半消息机制适合分布式事务场景堆积与回溯消息默认保留 7 天可按 offset 偏移回溯适合流式计算堆积过多会显著降低性能一般按队列需求清理堆积能力强支持定时/延迟消息典型场景日志采集、监控指标、用户行为事件、大数据管道、实时数仓订单通知、任务调度、微服务间常规异步解耦电商交易、金融场景、事务消息、延迟消息选型经验我得说一句如果核心诉求是微服务之间异步解耦、低延迟、运维轻RabbitMQ 更省心如果消息量大、要长期保留、后面还要接 Flink/Spark 或实时数仓Kafka 几乎不用犹豫RocketMQ 则更适合深度 Java 技术栈、有事务消息和定时消息诉求的团队。但无论选哪个发送链路底层的追问是一致的谁负责缓冲、谁负责重试、什么时机才算确认成功。这几个问题搞不明白换哪个 MQ 都会踩坑。1.2 先认识消息本体即将被 send() 的那条 Record 长什么样讨论发送过程之前得先把消息本身看透。Kafka 里的每一条消息本质是一个 ProducerRecord它的核心字段如下字段含义由谁决定topic消息进哪个主题发送时指定partition分区号可选没指定时由分区器计算key分区路由依据发送时指定可以为 nullvalue消息内容发送时指定timestamp时间戳默认客户端当前时间headers自定义 KV 元数据发送时设置可以把这个结构理解成快递单topic 是城市partition 是小区key 是分拣规则offset 是楼层和房间号。发送前我们通常只填了 topic、key、value 和 headerspartition 和 offset 都是后续过程中被 Kafka 决定的。如果你发送时手动指定 partition分区器就不参与如果指定 key就走哈希路由如果 key 也是 null走的是粘性分区策略。整个发送链路的逻辑就是从这几个字段开始展开的。另外提醒一句单条消息大小要参考客户端 max.request.size默认只有 1MB超了会直接报 RecordTooLargeException。真要传大对象客户端、broker 的 message.max.bytes、副本拉取的 replica.fetch.max.bytes 三层得同步放量这个细节到第 4 章再展开。2. 一次 send() 的完整旅程从业务线程到磁盘落盘2.1 初始化 KafkaProducer 时到底发生了什么先看一段最常见的初始化代码Properties props new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, 10.0.0.1:9092,10.0.0.2:9092,10.0.0.3:9092); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.ACKS_CONFIG, all); KafkaProducerString, String producer new KafkaProducer(props); producer.send(new ProducerRecord(order_events, orderId, payload));new KafkaProducer 的时候客户端并不会主动连接任何 broker它只是把内部组件搭好序列化器、分区器、拦截器、RecordAccumulator发送缓冲池、Metadata元数据缓存以及一个后台运行的 Sender 线程。真正的联网发生在第一次 send() 之后客户端需要先拉取 topic 的元数据搞清楚每个分区的 leader 在哪台 broker 上。这个细节决定了初始化成功不等于集群连通。我见过不少人生产环境里 producer 创建没报错就以为集群状态正常结果第一次 send 直接超时。所以看到 TimeoutException 先别急着怀疑代码先从元数据有没有拉到查起。2.2 分区器Partitioner把消息分到哪条道上去send() 调用后消息先经过序列化把 key 和 value 转成字节数组然后交给分区器做路由。路由规则分三种情况发送时指定了 partition直接用指定值分区器不介入发送时指定了 key对 key 做 murmur2 哈希再对分区数取模。同一个 key 永远落进同一个分区这是按 key 保序的基础key 为 nullKafka 2.4 之后默认走粘性分区Sticky Partition。它不会每次随机选一个分区而是先挑一个分区然后拼命往这个分区对应的批次里攒消息直到批次满或者 linger 超时才换下一个分区。这么设计是为了让批次更满减少网络小请求数量。这里有一个很容易被忽略的坑按 key 取模路由的结果是时变的一旦 topic 分区数扩大原 key 的哈希取模结果就会变老数据和新数据可能落到不同分区基于分区的按 key 有序从扩容那一刻就失效了。所以业务上对消息 key 与顺序强相关的 topic创建时就要把分区数规划好尽量避免后续扩容。2.3 RecordAccumulator缓冲区到底是怎么蓄水的分区器算完分区之后消息进入 RecordAccumulator——整个发送流程的蓄水池。每个分区维护一个双端队列队列里是若干个批次ProducerBatch。新消息到达后会尝试追加到当前批次的末尾如果当前批次已经达到了 batch.size 的上限就新起一个批次。这个蓄水池的总容量由 buffer.memory 控制默认只有 32MB。当池子被打满、所有批次又都还没发出去时send() 调用会被阻塞等待最多等 max.block.ms默认 60 秒超时后直接抛 TimeoutException。线上常见的send 卡住几十秒后报错消息一堆堆压在本地多半就是缓冲池打满了而不是集群真的挂了。把发送链路比作机场行李传送带RecordAccumulator 就是安检前的暂存区。暂存区堆不下时后面排队的人只能原地等。这也是为什么生产者侧的 buffer.memory 设置太小会表现为发送超时单纯调大 broker 分区数量根本解决不了。2.4 Sender 线程攒批、按 Broker 分组、发出真正的发送引擎是后台 Sender 线程它循环做四件事从 Metadata 里确认每个分区当前 leader 在哪个 broker把条件满足的批次挑出来批次满了的、linger.ms 超时的、或者其他原因需要 flush 的把发往同一个 broker 的批次合并成一个 ProduceRequest交给网络层非阻塞发送。比如 topic 有 12 个分区、3 个 broker 各做 4 个分区的 leader理论上这几个分区的批次可以合并成一个请求大幅减少网络往返收到响应后把对应批次标记为完成触发回调逻辑。Sender 是单线程它的处理能力就是整个 producer 的硬上限。排障时如果发现队列缓冲字节数一直很高、但 Send 速率上不去大概率就是 Sender 线程或 broker 侧网络已经到瓶颈了这时候加再多的业务线程并发调用 send() 也不会改善反而会加剧缓冲区堆积。2.5 Broker 落盘从 ProduceRequest 到响应客户端把 ProduceRequest 发到 broker 后broker 也不是一个线程包打天下。SocketServer 的网络线程Processor先读取请求丢进请求队列KafkaRequestHandler 线程再从请求队列取出处理根据请求里的分区信息定位到 leader 副本的日志段把消息顺序追加到页缓存Page Cache。这一步是典型的顺序写速度快是 Kafka 高吞吐的根基。关键在响应的时机acks0broker 不做任何响应消息发出去就当成功了acks1leader 把消息写入页缓存后立即响应acksallleader 要等 ISR同步中的副本集合里所有副本都追上这条消息才给客户端返回。所以用 acksall 时响应延迟 同步副本拉取 确认的耗时。它不是一定比 acks1 慢一大截那么吓人但副本之间网络抖动、ISR 里混入慢副本时延迟就会被放大。这里必须强调一个冷知识Kafka 默认的落盘其实只到页缓存不保证 fsync 到物理磁盘。页缓存换来的是吞吐代价是节点断电时可能丢数据。真正强一致的场景需要增加副本数或者调整 log.flush.interval.messages 和 log.flush.interval.ms但吞吐会明显下降。高吞吐和强一致之间没有免费午餐。2.6 回调执行时机比你以为的晚一步send() 返回的 Future 和注册的 Callback看起来顺理成章但回调到底在哪个线程执行很多人想当然。实际机制是回调不会由 Sender 网络线程立刻执行而是先进入客户端的回调队列等到下一次 send() 或 poll() 时由用户线程触发执行。换句话说回调是寄生在业务线程上的。这个特性带来一个很实在的经验不要在回调里做重活不要同步查数据库、不要同步调远程接口。一旦回调很慢下一次 send() 就会跟着慢。有些团队排查发送慢排查了半天最后发现是回调函数里的 RPC 超时拖垮了发送主线程。回调里只做轻量标记重活丢到独立线程池去异步处理能避开大部分这类问题。3. 决定发送快慢的几个参数背靠背不是孤立的3.1 acks三个档位对应三种数据安全承诺acks 是发送链路最重要的可靠性开关三个档位的语义经常被背错acks 值语义丢数据场景适用场景0fire-and-forget不等任何确认发送即丢网络抖动、broker 挂掉都可能丢日志、监控等可丢数据追求极致吞吐1leader 写入本地页缓存即算成功leader 写完后崩溃且副本还没同步大多数业务吞吐和可靠性折中Kafka 3.0 前默认all等 ISR 所有副本同步完再返回配合 min.insync.replicas 后只有同时宕掉足够多副本才会丢交易、订单、对账等强一致场景Kafka 3.0 默认用一句话总结acksall 不是消息绝对不丢而是ISR 里所有副本都认了账才算数。如果 ISR 里只剩 Leader 一个副本其他副本都掉线了acksall 和 acks1 没有本质区别。所以 acksall 必须和 min.insync.replicas 配合起来看否则可靠性语义会被悄悄稀释。这个认知在集群运维决策时非常关键。3.2 retries 与幂等重试为什么没那么可怕了发送失败后客户端会重试retries 控制次数retry.backoff.ms 控制重试间隔。在没有幂等的老版本里重试可能带来两个副作用重复投递和乱序。比如消息 A、B 连续发往同一分区A 的第一次请求失败B 却成功了A 重试成功后反而排在 B 后面顺序就颠倒了。Kafka 3.x 客户端默认开启 enable.idempotencetrue原理是给每个分区维护一个自增 sequencebroker 按照 PID 分区号 sequence 做去重一旦发现旧 sequence 就能识别为重复批次直接丢弃。幂等开启后retries 默认会变成一个很大的数、max.in.flight.requests.per.connection 被限制在 5 以内目的就是给重试留出恢复空间同时保证单分区内请求有序。所以面试常问的Kafka 如何保证不重复第一层答案是幂等生产者保证单分区内不重复第二层是靠消费端做去重。要注意幂等保障的边界是单分区内不是跨分区的全局有序或全局不重复跨分区一致得靠事务 API那是另一套机制。3.3 batch.size、linger.ms、buffer.memory吞吐与延迟的三元平衡这三个参数决定发送链路蓄水和排水的行为是最容易被拍脑袋乱改的一组。batch.size 默认 16KB。批次越大能塞的消息越多网络请求越少吞吐越高但攒满一个批次的时间也更长延迟跟着上升broker 端请求解析的内存消耗也变大。linger.ms 默认 0。设置之后批次即使没满也会最多等 linger.ms 再发。想要极低延迟就保持 0想明显提升吞吐5ms 到 20ms 是高性价比区间。buffer.memory 默认 32MB。它是整个缓冲池占用的堆内内存太小会在高 TPS 下频繁阻塞调大之前要先确认 JVM 堆够不够。粗略估算公式高峰期每秒生产的字节数乘以你能容忍的最大阻塞秒数。比如峰值 10MB/s、想扛 10 秒buffer.memory 至少得给到 100MB。很多同学只调 batch.size忽略了 linger.ms 的默认值是 0结果批次永远攒不满消息一到就发吞吐毫无变化。这三个参数必须一起调而且要结合你们最高峰一秒要推进多少数据来算而不是从 16KB 改成 64KB 就完事。3.4 max.in.flight.requests.per.connection顺序和吞吐之间的暗线这个参数控制单个连接上同时最多有几个未确认请求默认 5。它和顺序性的关系很微妙如果 retries 大于 0、max.in.flight 又大于 1某个请求失败后重试就可能跟后续请求交错导致乱序。老版本里有人为了保顺序把值改成 1代价是单连接吞吐下降。Kafka 3.x 幂等机制出现后这个参数在幂等开启时默认被限制在 5 以内客户端靠 sequence 序号在 broker 端做排序既保顺序又保吞吐把这道两难解开了。所以我的建议是不要轻易调大或调小这个值更不要为了吞吐调到十几幂等约束下它起不到预期作用反而可能引入乱序风险。3.5 延伸消费端多线程怎么保住发送端好不容易保住的顺序Kafka 的顺序保证边界要记清楚单分区内生产者按 sequence 写入有序同一个分区也只会发给同一个消费组内的一个消费者实例。但如果你在消费端开线程池把消息丢给多个线程并发处理发送端做的所有顺序努力就白费了。要既并发又保序常见做法是按 key 分片 有序消费消费端拿到消息后把同一个 key 的消息通过哈希投递到固定的内部队列每个队列由一个线程消费队列数量等于并发线程数。这样同一个订单永远排在同一个线程里处理不同订单之间可以并行。这个模式我在实际项目里验证过是性价比最高的方案而且和发送端同一个 key 永远进同一个分区的思路是一一对应的。4. 发送链路上的坑和排查链路4.1 消息延迟高先从三张时间账单查起现象是 produce 延迟从正常 10ms 涨到几百毫秒甚至秒级。我一般的排查思路是倒着看时间花在了哪里看等待批次时间linger.ms 设置太大、批次始终不满消息在本地泡着。这个最容易发现但也最容易误判因为它是配置造成的故意延迟不是故障。看请求排队时间检查 producer 指标 queue.buffered.total.bytes如果持续上涨、接近 buffer.memory说明客户端侧蓄水池快满了。要么 Sender 线程处理不过来要么 broker 响应慢导致 in-flight 请求积压。看broker 处理时间broker 的 request handler idle ratio 和请求日志里的耗时。如果 handler idle 接近 0说明 broker CPU 忙不过来如果请求耗时高但 handler 不忙大概率是磁盘慢或者副本同步慢。用 acksall 时尤其要查 ISR 里有没有滞后副本。排查动作按这个顺序走通常十几分钟能定位。别一上来就改参数先弄清楚延迟发生在哪一段不然就是瞎调一个参数然后碰运气。这个是发送链路排障的基本功。4.2 InvalidReceiveException看到这个报错先别慌检查五件事org.apache.kafka.common.network.InvalidReceiveException: Invalid receive (size ... larger than ...) 在社区里很常见本质是 broker 侧网络线程在解析请求头时发现收到的长度字段异常认为这不是一个合法的 Kafka 协议请求。原因通常不在业务代码而在链路或参数请求体超过了 broker 的 socket.request.max.bytes默认 100MB或者客户端 max.request.size 与 broker 不一致客户端与 broker 大版本差距过大协议头部解析错位中间经过了四层负载均衡或端口转发TCP 流被截断或重组出畸形请求有非 Kafka 客户端往 9092 端口发了请求比如 HTTP 探测、监控脚本之类broker 自然解析失败。排障做法先看 broker 日志里该异常的来源 IP确认是不是 Kafka 客户端再确认 max.request.size 和 socket.request.max.bytes 的关系如果中间有 LB对比直连和走 LB 的行为差异必要时用 tcpdump 抓包看 TCP 段的前几个字节和 Kafka 请求协议魔数是否匹配。这个报错表面吓人实际排查起来路径很清晰关键是不要陷在业务代码里找。4.3 分区 Leader 切换与 Metadata 更新发送链路最隐蔽的抖动源发送端知道消息该发给哪台 broker依赖的是客户端 Metadata 缓存。如果某个分区 leader 发生变化比如 broker 重启、磁盘故障、副本被踢出 ISR发送端会收到 NotLeaderOrFollower 之类的异常然后重新拉取 Metadata。这个过程至少多一次 RTT表现就是发送延迟瞬间抖一下。我见过最典型的场景集群里一台 broker 反复重启发送端周期性出现几十毫秒到几百毫秒的尖刺查 topic 才发现那个分区的 leader 在几台 broker 之间来回切。所以发送链路稳定不只是调客户端参数集群侧稳定性同样重要。减少这类抖动的手段包括保证 broker 机器硬件和磁盘健康、不要频繁 kill -9 节点、合理设置 replica.lag.time.max.ms避免正常偏慢的副本被误踢出 ISR。建议每次排查发送抖动都顺手跑一下 kafka-topics.sh --describe看一眼分区 leader 和 ISR 是否完整。这一条命令成本很低收益很高。4.4 发送侧常被忽略的三个小坑单条消息过大。默认 max.request.size 是 1MB超了直接报 RecordTooLargeException。要放大不能只改客户端一个参数broker 的 message.max.bytes、replica.fetch.max.bytes 必须同步放量否则消息就算发出去了副本同步阶段也会被拒绝。这个三层一致是经典坑。序列化器配错。KafkaProducer 报 ConfigException 或 ClassNotFoundException 的常见原因就是 key.serializer 或 value.serializer 没配或者把反序列化器配到了生产端。问题本身很低级但发生频率并不低。auto.create.topics.enable 关闭时往一个不存在的 topic 发消息发送端会一直尝试拉元数据直到超时。现象是send 永远不报错但永远没有结果。遇到这种卡死先确认 topic 到底存不存在别在参数上浪费时间。5. 从发送链路延伸到集群部署配置、排查工具与硬件关系5.1 3节点集群部署时哪些配置直接决定发送成败单看发送链路集群侧最相关的配置集中在 server.properties 和 topic 创建环节。我常用的 3 节点最小合理配置是broker.id0 num.network.threads8 num.io.threads16 socket.request.max.bytes209715200 socket.send.buffer.bytes1048576 socket.receive.buffer.bytes1048576 default.replication.factor3 min.insync.replicas2每台 broker 的 broker.id 必须全局唯一num.network.threads 是 SocketServer 读请求的线程数num.io.threads 是真正处理 produce/fetch 请求的线程数CPU 核数多可以再调大。topic 创建时副本因子设 3、min.insync.replicas 设 2配合客户端 acksall允许一台 broker 挂掉而发送仍然可用但两台同时挂掉就会报 NOT_ENOUGH_REPLICAS。这个权衡要在部署时讲清楚一致性要求越高可用性边界越窄。另外 message.max.bytes 必须和客户端 max.request.size 对齐。很多默认配置只有 1MB如果业务要传 5MB 的文件或大日志不改这两个参数发送端迟早会踩到大消息被拒的报错。5.2 AdminClient发送链路排障时的后台数据透视Kafka 自带的 AdminClientkafka-topics.sh、kafka-configs.sh 底层就是它是排查发送问题的第一工具。几个最常用的操作# 查看 topic 明细分区、leader、replicas、isr kafka-topics.sh --bootstrap-server broker1:9092 --describe --topic order_events # 查消费者组 lag kafka-consumer-groups.sh --bootstrap-server broker1:9092 --describe --group order_consumer_group # 查 broker 配置确认 max.request 相关 kafka-configs.sh --bootstrap-server broker1:9092 --entity-type brokers --entity-name 0 --describe从发送链路视角我排障时最关心三个字段Leader 表示分区当前由哪个 broker 提供服务Replicas 表示副本分布Isr 表示同步副本是否完整。如果 Isr 列表明显比 Replicas 短说明有副本在 lagacksall 的发送延迟会受影响。这些命令几乎没有成本但能省下大量靠猜的时间。5.3 可视化工具直观看到消息到底有没有发进去、发到哪个分区命令行虽好日常巡检和团队协作还是需要可视化工具。我常用的有三类Offset Explorer原 Kafka Tool老牌桌面工具免费看 topic、分区、offset、消费组 lag 很方便还能手动查看消息内容适合个人快速验证。kafka-uiProvectusWeb 版现代界面能直接展示生产速率、消费速率、consumer lag、分区分布实时曲线对定位发送抖动很有帮助适合团队共用。Kafdrop轻量 Web 工具侧重快速浏览消息内容适合 Demo 和测试环境。我自己的习惯是线上排障先用 kafka-topics.sh 确认分区和 ISR再用可视化工具看生产速率和消费速率曲线。工具只是辅助关键是你能把生产速率下降对应到发送链路的哪个环节。工具用得再多链路本身不理解看到的只是表象数字。5.4 读写最大值和硬件的关系Kafka 的高吞吐到底高在哪Kafka 单 broker 能扛多高的读写上限本质由四个硬资源决定磁盘顺序写速度。消息先写页缓存、异步刷盘磁盘只要顺序写就不会太差SSD 在高并发下优势明显机械盘主要怕随机写。网卡带宽。千兆网卡理论 125MB/s实际收发并行时要打折扣万兆网卡能明显拉开吞吐上限。内存页缓存。热数据能留在页缓存里消费端走零拷贝读取读放大很小内存不够就退化到磁盘读吞吐断崖式下跌。CPU 与批次大小。压缩、序列化、副本同步都吃 CPU批次越大单位 CPU 处理的请求越少所以批次大小是影响最高吞吐率的最重要软件变量。给个量级参考单条 1KB、千兆网卡、SSD、批次调优后单 broker 生产 5万到 10万条/秒约 50~100MB/s已经不错万兆网卡加多分区并行的组合单 broker 更高也能看到。但最大读写从来不是一个固定数字它和可靠性配置强相关——acksall 加多副本确认必然会吃掉一部分吞吐上限。所以讨论 Kafka 极限时先问一句你要不要丢数据要等几个副本确认最后说点个人体会。我这些年跟 Kafka 打交道最值钱的一个习惯就是线上发送抖动出现时永远先画一遍生产链路时间账——消息在本地缓冲区等了多久在网络上跑了多久broker 处理花了多久ack 等副本同步又等了多久。很多时候排查到最后原因都不是什么高深故障而是某个参数在特定量级下暴露了瓶颈。这套思路要是你也能养成再回头看Kafka 消息的发送过程这九个字就不再是面试八股而是你定位所有线上问题的第一张地图。
返回列表