Kafka提供了两种消息发送方式:同步发送和异步发送。同步发送通过producer.send(msg).get()实现,发送后会阻塞等待返回结果,能准确获知消息发送状态;异步发送则采用producer.send(msg, callback)方式,通过回调机制处理结果,显著提升吞吐量。

消息发送流程涉及两个核心线程和一个关键组件:

  1. 主线程(Main):负责初始化配置、创建Producer实例并执行发送逻辑,根据用户选择的发送方式处理消息
  2. 发送线程(Sender):负责实际网络通信,处理发送请求结果
  3. 消息累加器(RecordAccumulator):实现消息批量处理的关键组件

消息发送过程的具体步骤:

  1. 调用send方法后,消息依次经过:
    • 拦截器:支持消息预处理(修改/日志/统计等)
    • 序列化器:将键值对象转为字节数组
    • 分区器:确定目标Partition
  2. 消息进入RecordAccumulator缓冲区,根据配置参数(batch.size/linger.ms)进行批量处理
  3. 满足条件时,Sender线程通过NetworkClient和Selector组件将批次消息发送至目标Partition Leader
  4. Leader将消息写入本地日志,并通过复制机制同步到Follower副本
  5. Follower写入成功后向Leader发送ACK确认

Kafka提供三种ACK机制,通过request.required.acks参数配置:

  • 0:不等待ACK,可能丢失数据
  • 1:等待Leader确认,不保证Follower同步
  • -1:等待所有ISR副本确认,保证最高可靠性
Logo

腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。

更多推荐