Kafka的消息是怎么读写的?


一、写操作

1、Kafka存储模型=分区(Partition)+日志分段(Log Segment)

在Kafka中每个Topic都被分成多个分区,每个分区都是一个有序、不可变的消息日志流(只追加不修改
同时为了避免单个日志文件过大导致IO效率低,所以每个分区的日志都为拆分成多个日志分段,每个分段包含:

  • .log:核心数据文件,存储数据原始信息(二进制格式)
  • .index:偏移量索引文件,记录offset->.log文件物理位置,加速读取
  • .timeindex:时间戳索引文件,记录timestamp->offset,支持时间范围查询

分段规则

  • 大小阈值:超过阈值自动创建新分段(log.segment.bytes=1073741824 默认1GB)
  • 时间阈值:超时旧分段自动删除/归档(log.retention.hours=168 默认7天)

2、消息写入流程

a. 序列化消息:默认转成字节数组
b. 分区:默认轮询策略(尽可能确保每个分区数据差不多),也可以指定key
c. 请求路由:producer获取分区leader并发送
d. 内存缓存:leader收到消息先写入内存缓冲区,再根据策略刷盘(减少磁盘随机IO,提高吞吐)
e. 顺序追加写入内存缓冲区达到阈值就会刷盘到最新的.log文件
f. 同步副本:leader写入成功后,会通知当前分区的follower副本同步消息(fetch)
g. 提交确认:所有ISR副本同步完消息后,leader会更新分区的高水位线(HW),并返回ack(如需)
h. 刷盘持久化:异步刷盘,通过定期、触发阈值将数据刷到自盘

3、消息的物理结构

前面说到,每个消息在.log中都是二进制格式,主要字段如下:

  • offset:分区唯一消息id(单调递增)
  • timestamp:写入时间
  • KV:消息存储内容

二、写操作

1、消息读取流程

a. 分区分配:默认使用消息哈希值进行分区,特殊场景可指定(如同一个uid分配到同一个分区)
b. 获取消息信息:通过metadata获取leader副本位置,以及分区的日志分段信息
c. offset定位

通过分区记录的offset定位消息所在.log文件
再通过.index文件定位消息在.log中的物理偏移量(避免全文扫描)
如按时间戳查询,则通过.timeindex找到对应offset,则走上述流程

d. 数据读取:通过零拷贝直接拷贝数据,不通过用户态切换,提高吞吐
e. 消息返回
f. 提交offset:默认自动提交(enable.auto.commit

三、总结

1、Kafka读写有什么亮点?

  • 顺序写:通过顺序写到日志文件,读取时保证是顺序IO,效率远远高于随机IO
  • 零拷贝:读取数据会直接在系统层进行拷贝,不经过用户层,提高吞吐
  • 批量读写:支持批量发送和拉取,减少网络IO
  • 分段存储:避免单日志文件过大导致IO效率下降
  • 索引:通过额外索引文件,快速定位信息,避免全文扫描
  • 副本异步同步:根据业务场景配置合适的acks,平衡消息可靠性和吞吐

文章作者: GaryLee
版权声明: 本博客所有文章除特別声明外,均采用 CC BY 4.0 许可协议。转载请注明来源 GaryLee !
  目录