为什么选用kafka作为你们的消息队列
- 因为我们使用ClickHouse存储网络报文解析数据,经过调研ClickHouse有现成的开发方案可以与kakfa进行无缝接入,从而实现大数据的存储与清洗
- kakfa作为老牌开源的消息队列,相较于其他消息队列,各类问题都可以在开源社区找到找到相应的解决方案和技术文档
- kafka拥有顺序读写、零拷贝、topic分区等特性,在大数据量的情况下依然有很好的承载能力和吞吐能力,可以稳定的支撑我们的平台进行数据处理
- kafka 支持将消息持久化到磁盘,并保留可配置的时间。这样在我们的程序消费能力不够或者出现故障的情况下,依然可以保证数据不会直接丢失,通过重置offset可以在业务程序恢复后重新消费数据保证业务的连续性与稳定性
- kakfa的topic可以由多个Partition组成,Partition分布在不同的Broker上。只需增加 Broker 节点并迁移分区,就能线性扩展写入、存储、读取的能力,在运维侧可以根据业务实际情况快速进行扩容,方便运维
- kafka支持同一个消费组的多个消费者消费同一个topic,在微服务场景下,可以通过简单的配置实现多个消费者高效消费数据,为我们减小了业务服务开发的复杂性
kafka为什么性能比较好
- kafka顺序读写磁盘,消息写入和读取速度接近磁盘IO速度。kafka消息只追加不修改,所有消息以日志形式,按顺序追加到分区文件的末尾。这完全规避了磁盘性能最差的随机写入。
- kafka采用零拷贝(Zero-Copy)网络文件的传输。传统网络传输需要经历多次的内核态和用户态数据拷贝。而 Kafka 从磁盘传输消息给消费者时,利用了 Linux 的 sendfile 系统调用,在内核态完成文件的零拷贝操作
- 不依赖 Java 堆内存做缓存,消息直接写入操作系统的 Page Cache,由 OS 负责 I/O 调度和缓存,避免了大对象在 JVM 堆中造成 GC 压力。消费操作同样顺序化,消费者基于偏移量(offset)顺序读取,磁盘磁头几乎不用频繁移动。
- kakfa支持批量生产消息和消息压缩,提升了消息传输效率。
- 支持topic分区,一个 Topic 可切分为多个 Partition,分布在不同 Broker 上。生产者可并发写入多个分区,消费者组内,一个分区只由一个消费者处理,但可多个消费者并行消费不同分区,可以线性的扩容分区和消费者,极大的提升了消息的吞吐量
- Kafka Broker 不需要向消费者主动推送消息,只是简单记录消费者对应的offset信息,专注消息的存储与写入。
kafka的ACK参数解析
- ACK = 0 生产者生产消息时无需得到Broker的确认写入的响应,如果写入过程中broker出现问题,那么消息会直接丢失
- ACK = 1 生产者生产消息时只需要得到Broker中leader的确认写入的响应即可,如果写入leader的过程中出现了leader崩溃、leader重新选举的情况消息则会重新生产。如果是副本同步出现leader崩溃、leader重新选举的情况,则这条消息就会丢失,因为新选举的leader没有这条消息
- ACK = -1 在生产者发送消息后,只有等服务端ISR中的所有副本全部写入数据,才会返回成功的响应给生产者
在生产中如何实现高效的生产和消费kafka的消息
生产者侧
- 使用
sarama.AsyncProducer异步生产者,直接将消息推送到channel,避免阻塞 - Sarama 的异步生产者(AsyncProducer)有一个硬性要求:它会在后台不断往
Successes(成功)和Errors(失败)这两个通道里塞消息回执。启动一个后台协程去读取并清空这两个通道里的数据。确保通道永远畅通,从而保证我们的生产者能持续、高效地发送消息,不会因为没人处理回执而引发死锁
消费者侧
- 换更快的底层库:抛弃 Go 自带的
encoding/json,改用字节跳动的 Sonic。它底层用了汇编和 SIMD 指令,速度是原生库的好几倍。 - 对象池复用(sync.Pool):JSON 解析会产生很多临时对象,导致频繁 GC(垃圾回收)。用
sync.Pool把缓冲区“借”出来用,用完再“还”回去,极大减少内存分配和 GC 卡顿。
生产过程中如何保证kafka的消息不会被重复消费
生产端:防止重复发送
- 开启幂等性(Idempotence):将
enable.idempotence设为true。Kafka 会为生产者分配 PID 和递增序列号,Broker 端会自动识别并丢弃重复发送的消息,保证单分区内的消息不重复。 - 开启事务(Transaction):对于跨分区的复杂场景,配置
transactional.id,利用 Kafka 事务保证跨分区写入的原子性,避免部分成功导致的重复。
消费端:减少重复概率
- 手动提交 Offset:关闭自动提交(
enable.auto.commit=false),在业务逻辑真正执行成功后,再手动提交 Offset。这样即使消费者宕机,重启后也只会从上次成功的位置继续消费,将重复的窗口缩到最小。
业务端:最终兜底
由于网络抖动或消费者在处理完业务但还没来得及提交 Offset 时宕机,重复消费无法绝对避免,因此必须保证业务逻辑的幂等性:
- 数据库唯一索引:利用业务唯一键(如订单号)做唯一约束,重复插入直接忽略。
- Redis 去重:消费前用
SETNX写入消息 ID,存在则跳过。 - 状态机校验:处理前检查业务状态,已完成则直接跳过。
kafka集群是如何为消费组的消费者分配分区的
如果在一个消费组中,有 6 个分区但部署了 8 个消费者实例,Kafka 的分配结果会非常明确:会有 6 个消费者实例各自分到 1 个分区,而剩下的 2 个消费者实例将分配不到任何分区,处于空闲(空跑)状态。
Kafka 在同一个消费组内遵循严格的“排他性规则”:在任何时刻,一个分区恰好由消费者组中的一个 KafkaConsumer 实例拥有。 同一组中的两个实例绝对不会同时读取同一个分区。
kafka集群如果出现了某个节点宕机会发生什么或者新增一个节点会发生什么
节点宕机:
- 自动切换 Leader:宕机节点上的分区会自动在其他副本中选出新的 Leader,业务几乎无感。
- 副本数变少:该节点的副本会掉线,集群处于“未完全同步”状态,等它恢复后会自动追赶数据。
新增节点:
- 不会自动分担负载:新节点加入后是“空”的,不会自动迁移旧数据。
- 需要手动重平衡:必须手动执行数据迁移(Rebalance),把旧分区挪一部分过去,新节点才能真正分担压力。