Kafka 消息的发送过程
·
Kafka提供了两种消息发送方式:同步发送和异步发送。同步发送通过producer.send(msg).get()实现,发送后会阻塞等待返回结果,能准确获知消息发送状态;异步发送则采用producer.send(msg, callback)方式,通过回调机制处理结果,显著提升吞吐量。
消息发送流程涉及两个核心线程和一个关键组件:
- 主线程(Main):负责初始化配置、创建Producer实例并执行发送逻辑,根据用户选择的发送方式处理消息
- 发送线程(Sender):负责实际网络通信,处理发送请求结果
- 消息累加器(RecordAccumulator):实现消息批量处理的关键组件
消息发送过程的具体步骤:

- 调用send方法后,消息依次经过:
- 拦截器:支持消息预处理(修改/日志/统计等)
- 序列化器:将键值对象转为字节数组
- 分区器:确定目标Partition
- 消息进入RecordAccumulator缓冲区,根据配置参数(batch.size/linger.ms)进行批量处理
- 满足条件时,Sender线程通过NetworkClient和Selector组件将批次消息发送至目标Partition Leader
- Leader将消息写入本地日志,并通过复制机制同步到Follower副本
- Follower写入成功后向Leader发送ACK确认
Kafka提供三种ACK机制,通过request.required.acks参数配置:
- 0:不等待ACK,可能丢失数据
- 1:等待Leader确认,不保证Follower同步
- -1:等待所有ISR副本确认,保证最高可靠性
更多推荐
所有评论(0)