Kafka
Kafka 是一个分布式的消息队列 / 流平台。它解决的核心问题是:把数据的生产方和消费方解耦——生产方只管往里写,消费方按自己的节奏来读,两边的速度、上下线互不影响。在数仓体系里,它常被放在数据入口处充当缓冲与总线:业务日志、埋点、CDC 变更先打进 Kafka,再由离线/实时链路各取所需。
核心概念
- Topic(主题):一类消息的逻辑名字,比如
order_log。生产和消费都以 Topic 为单位。 - Partition(分区):一个 Topic 被切成多个分区,分区是并行和有序的最小单位。消息在单个分区内有序,跨分区不保证全局有序。
- Replica(副本):每个分区有多个副本分布在不同 Broker 上,一个是 Leader(读写都走它),其余是 Follower(同步数据)。Leader 挂了从 Follower 里选新 Leader,这是高可用的基础。
- Producer / Consumer:写入方和读取方。
- Consumer Group(消费组):一组消费者共同消费一个 Topic,一个分区同一时刻只会被组内一个消费者消费。所以消费的并行度上限就是分区数。
- Offset(位点):消费者读到哪了,用 offset 记录。位点由消费者自己提交,因此「重复消费」和「消息丢失」本质上都是位点提交时机的问题。
为什么 Kafka 这么快
不是因为用了什么黑魔法,而是几条工程选择叠加:
- 顺序写磁盘:消息追加到分区文件末尾,顺序 IO 的速度接近内存随机 IO,远快于随机写。
- 页缓存(Page Cache):读写都先走操作系统页缓存,热数据几乎不落到实际磁盘读。
- 零拷贝(Zero-Copy):发送消息时用
sendfile直接从页缓存送到网卡,省掉内核态到用户态的多次拷贝。 - 批量 + 压缩:生产端按批发送,配合压缩,摊薄网络和 IO 开销。
可靠性:ISR 与 ack
- ISR(In-Sync Replicas):与 Leader 保持同步的副本集合。只有 ISR 里的副本才有资格被选为新 Leader。
- acks 参数决定生产端认为「写成功」的标准:
acks=0:发出去就算成功,最快但可能丢。acks=1:Leader 落盘即成功,Leader 切换瞬间可能丢。acks=all:ISR 全部确认才成功,最稳,代价是延迟更高。
exactly-once(精确一次) 靠幂等 Producer(给消息编序号去重)+ 事务(跨分区原子写)实现,常用于「消费—处理—再写回 Kafka」的流处理场景。
常见追问方向
- 如何保证消息顺序:把需要有序的消息用同一个 key 路由到同一分区。
- 如何避免重复消费 / 消息丢失:核心在位点提交时机与 ack 配置的权衡,以及下游做幂等。
- 消息积压怎么办:扩分区 + 扩消费者(但消费并行度受分区数限制),或提升单条处理效率。
- Kafka 和 RabbitMQ 的区别:Kafka 面向高吞吐的日志/流场景、消息可回溯重放;传统 MQ 更偏重复杂路由和逐条投递语义。
在数仓里的位置
Kafka 是实时数仓的「主动脉」:ODS 层的实时接入、维表变更(CDC)分发、多个下游作业共享同一份数据流,都依赖它做削峰填谷和解耦。离线链路也常从 Kafka 落一份到 HDFS/对象存储作为 T+1 的数据源。
登录后可以选中正文添加批注(仅自己可见)。
评论 (0)
登录后参与评论。
还没有评论,来做第一个。