title: Kafka : 消息队列 date: 2026-03-01 authors: - name: 猫酒 email: 2738035238@qq.com categories: - Kafka - 消息队列 comments: true draft: false
Kafka(消息队列)¶
kafka的底层配置相关¶
生产者端重要配置¶
acks:应答确认机制,控制生产者发送消息后是否等待Broker确认acks=0:发送后不等待确认,吞吐量最高,但可能丢失数据acks=1:Leader写入成功后确认,Leader宕机可能丢失数据-
acks=-1/all:Leader和所有ISR中的Follower都写入成功才确认,最安全,默认值 -
retries:消息发送失败后的重试次数,默认值是Integer.MAX_VALUE retry.backoff.ms:两次重试之间的等待时间max.in.flight.requests.per.connection:客户端单连接上允许发送的未确认请求的最大数,默认5- 开启幂等性时可保持 <= 5 保证顺序
-
未开幂等性需要设置为1保证发送顺序
-
buffer.memory:生产者缓存未发送消息的内存大小,默认32MB batch.size:每个批次的大小,默认16KB,消息达到这个大小就会发送linger.ms:如果消息没填满batch,等待多久后发送,默认0mscompression.type:消息压缩类型,可选none、gzip、snappy、lz4,默认不压缩
Broker端重要配置¶
broker.id:每个Broker的唯一标识log.dirs:Kafka数据存储的目录zookeeper.connect:Zookeeper连接地址num.network.threads:处理网络请求的线程数num.io.threads:处理磁盘IO的线程数default.replication.factor:Topic默认副本数num.partitions:新建Topic默认分区数,默认1log.segment.bytes:单个Segment文件的大小,默认1GBlog.retention.hours:消息保留时间,默认7天log.retention.bytes:消息保留总大小,默认-1表示不限制min.insync.replicas:ISR中最小的同步副本数,配合acks=-1使用- 当ISR中副本数 < 这个值时,Broker拒绝写入消息,保证数据不丢失
消费者端重要配置¶
group.id:消费者组ID,同一个组内的消费者共同消费Topicbootstrap.servers:Kafka集群地址enable.auto.commit:是否自动提交offset,默认trueauto.commit.interval.ms:自动提交offset的间隔时间,默认5000msauto.offset.reset:当没有初始offset或offset超出范围时的行为latest:从最新消息开始消费(默认)earliest:从头开始消费none:抛出异常fetch.min.bytes:一次拉取请求最小拉取字节数,默认1Bfetch.max.wait.ms:如果没达到fetch.min.bytes,等待多久,默认500msmax.poll.records:一次poll拉取最大消息数,默认500session.timeout.ms:会话超时时间,默认10秒,超过这个时间没收到心跳就认为消费者宕机,触发再平衡heartbeat.interval.ms:发送心跳间隔,默认3秒,一般设置为 session.timeout.ms 的 ⅓
kafka的数据不丢失的保证¶
Kafka从生产者、Broker、消费者三个层面保证数据不丢失:
生产者端保证¶
- 合理配置acks应答级别
-
要求不丢失数据需要设置
acks=-1(即all),必须等待Leader和所有ISR中的Follower都确认写入成功才认为发送成功 -
开启重试机制
- 设置
retries > 0,发送失败自动重试,避免因网络瞬时故障导致数据丢失 -
配合幂等性避免消息重复
-
开启幂等性
enable.idempotence=true(0.11版本后默认开启),保证即使重试也只会持久化一条消息,不会重复
Broker端保证¶
- 多副本冗余存储
- 设置合理的副本数
replication.factor >= 2,每个分区存储在多个Broker上,Leader宕机后Follower可以接管 -
即使一台Broker宕机,数据也不会丢失
-
ISR同步机制
- 只有保持与Leader同步的Follower才会留在ISR中,Leader只等待ISR中所有副本同步完成才返回ack
-
配合
min.insync.replicas >= 2,当ISR中副本数小于这个值时,Broker拒绝写入,避免数据只写到一个副本就宕机导致丢失 -
数据持久化
- Kafka消息写入后直接追加到磁盘文件,Kafka重启后数据不会丢失
- 操作系统页缓存会定时刷盘,即使宕机也只会丢失极少量数据
消费者端保证¶
- 正确提交offset
- 消息消费完成后再提交offset,而不是消费前提交
- 如果使用自动提交(
enable.auto.commit=true),只要消费成功就会定期提交,一般不会丢失,但可能重复 -
手动提交(
enable.auto.commit=false),需要在消费完消息处理完业务逻辑后手动提交offset,保证处理完才提交,避免消费失败但offset已经提交导致数据丢失 -
消费失败处理
- 如果消费者处理消息失败,可以不提交offset,下次重启后会从上次失败的位置重新消费
- 对于重要消息,可以将失败消息存入死信队列后续处理,不会丢失
总结:不丢失配置最佳实践¶
| 层面 | 配置 |
|---|---|
| 生产者 | acks=-1 + enable.idempotence=true + retries > 0 |
| Broker | replication.factor >= 2 + min.insync.replicas >= 2 |
| 消费者 | 手动提交offset,消费完成后再提交 |
通过以上三个层面的配合,Kafka可以保证数据不丢失,实现at-least-once语义,配合幂等性可以实现exactly-once语义。
Kafka负载均衡策略¶
- 生产者方面,借助分区器实现负载均衡
- 分区策略:
- 直接指明 partition 的值
- key 不为 null:分区值=key 的 hash 值与 topic 的分区数取余
- key 为 null:消息将以轮询的方式,在所有可用分区中分别写入消息
- 消费者方面,通过 消费组协调器 与 消费者协调器,实现消费者再均衡操作。
- 触发再均衡操作的条件:
- 消费组中的消费者数量增加或减少
- 消费者宕机下线(不一定是真的下线,本质是消费者长时间未向消费组协调器发送心跳包);
- 消费组对应的 GroupCoordinator 节点发生了变更;
- 任意主题或主题分区数量发生变化
- 三种再均衡策略(即分区分配策略)
- Range(默认):按照消费者总数和分区总数进行整除运算,然后按顺序分配,如果数量有余数,就分配给前面几个消费者
- RoundRobin:所有 Topic 的所有分区进行字典排序,然后轮询分配
- sticky:尽量均匀的分配,分配尽量与上一次分配的相同
介绍一下Kafka的集群架构¶
- 一个 Kafka 集群由多个 broker 组成
- 一个 broker 下可以有多个 topic
- 一个 topic 又可以分为多个 partition
- 每个 partition 又有若干个副本,一个 leader,若干 follower
说一下你对kafka是怎么理解的¶
简述 kafka 的架构¶
Kafka 整体架构由以下几个核心组件组成:
1. Producer(生产者)¶
- 消息的产生方,负责将消息发布到 Kafka 的 Topic 中
- 支持异步批量发送,提高吞吐量
- 可以选择不同的分区策略将消息分发到不同分区
2. Consumer(消费者)¶
- 消息的消费方,从 Broker 拉取消息进行处理
- 消费者通过消费组(Consumer Group)的方式组织,同一个消费组内的消费者共同消费一个 Topic,每个分区只能被消费组内的一个消费者消费
- 支持拉模式(pull)消费,消费者根据自身能力控制消费速度
3. Broker(Kafka服务器)¶
- 一个 Kafka 集群由多个 Broker 组成,每个 Broker 存储一部分消息数据
- Broker 负责消息的存储、复制和读写请求处理
- 每个 Broker 可以容纳多个 Topic,不成为中心节点,扩展性好
4. Topic(主题)¶
- 消息的分类标识,生产者发送消息到特定 Topic,消费者订阅特定 Topic 消费
- Topic 是逻辑概念,物理上分为多个Partition(分区),分布在不同的 Broker 上
5. Partition(分区)¶
- Topic 的物理存储单元,每个 Topic 可以分为多个 Partition,分布在不同 Broker 上
- 每个 Partition 是一个有序的日志文件,生产者消息不断追加到末尾
- 每个分区有多个副本(Replica),其中一个是 Leader,其余是 Follower
- Leader:处理所有读写请求
- Follower:被动同步Leader数据,Leader宕机后从Follower中选举新的Leader
- ISR(In-Sync Replica):与Leader保持同步的副本集合,只有ISR中的副本完成同步后,Leader才会返回ack确认
6. Zookeeper / KRaft¶
- Zookeeper(老版本):负责存储集群元信息、选举控制器、监控Broker上下线、维护消费者偏移量(老版本)
- KRaft(Kafka 2.8+ 新版本):Kafka 自实现的元数据管理机制,不再依赖Zookeeper,性能更好
整体架构图¶
┌───────────┐ ┌─────────────────────────────────────────────┐ ┌───────────┐
│ Producer │────▶│ Kafka Cluster │────▶│ Consumer │
└───────────┘ │ ┌─────────┐ ┌─────────┐ ┌─────────┐ │ └───────────┘
│ │ Broker1 │ │ Broker2 │ │ Broker3 │ │
│ │ ┌─────┐ │ │ ┌─────┐ │ │ ┌─────┐ │ │
│ │ │ P0 │ │ │ │ P1 │ │ │ │ P2 │ │ │
│ │ └─────┘ │ │ └─────┘ │ │ └─────┘ │ │
│ └─────────┘ └─────────┘ └─────────┘ │
└─────────────────────────────────────────────┘
▲
│
┌───────┴───────┐
│ Zookeeper │
│ (KRaft) │
└───────────────┘
核心特点¶
- 分布式架构:分区分布在多个Broker,支持水平扩展
- 多副本冗余:通过ISR机制保证数据可靠性和高可用
- 高吞吐低延迟:顺序写入、零拷贝、批量发送等优化
※※简述kafka消息的存储机制¶
kafka中的消息就是topic,但topic 只是逻辑上的概念,而partition 才是物理上的概念,每个partition会对应一个log文件,它存储的就是producer 生产的数据。
- kafka的数据是放在本地磁盘的log文件上
生产者生产的数据会不断追加到log文件中,如果log文件很大了,就会导致定位数据变慢,因此kafka会将大的log文件分为多个segment,每个segment 会对应.log文件和.index文件和.timeindex文件,.log 存储数据,.index存储偏移量索引信息,.timeindex存储时间戳索引信息。
【log.segment.bytes=1g:指log 日志划分成segment的大小】
【log.index.interval.bytes=4kb】.index为稀疏索引,大约每往log文件写入 4kb 数据,会往 .index 文件写入一条索引,.index文件中保存的 offset为相对 offset,这样能确保 offset 的值所占空间不会过大
【log.retention.hours】Kafka 中默认的日志保存时间为7天
Kafka为什么具备高吞吐量?【速度快的原因】¶
- kafka 本身是分布式集群,并且采用分区技术,并行度高
- 读数据采用稀疏索引,可以快速定位要消费的数据
- 顺序写入log文件,一直是追加到文件末端。官网中说了,同样的磁盘,顺序写能达到600Ms,而随机写只有100K/S
- 实现了零拷贝技术。只用将磁盘文件的数据复制到页面缓冲区一次,然后将数据从页面缓冲区直接发送到网络中,这样就避免了在内核空间和用户空间之间的拷贝
补充:传统的读取数据发送到网络中的步骤?
- 操作系统将数据从磁盘文件读取到内核空间的页面进行缓存
- 应用程序将数据从内存空间读入用户空间缓冲区【x】
- 应用程序将读到的数据写回到内核空间并放入socket缓冲区【x】
- 操作系统将数据从 socket缓冲区复制到网卡接口,此时数据才能通过网络进行发送
简述 kafka 的分区策略¶
- 分区好处:便于合理使用存储资源;提高并行度
- 生产者分区策略
- 直接指明partition的值
- 没有指明partition的值但有key,那么
分区值=key的hash值与topic的分区数取余 - 没有指明partition也没有key,kafka采用sticky partition(黏性分区器), 会随机选择一个分区,并尽可能一直使用该分区,待该分区达到了 batchsize 大小或者达到了默认发送时间,kafka就会再次选择一个分区使用,只要与上一次的分区不同就可以了
- 自定义分区器:实现Partitioner接口,重写partition方法
- 消费者分区策略
- 【一个消费者组有多个消费者,一个topic有多个分区,所 以就会出现到底由哪个consumer来消费哪个partition的数据】【
partition.assignment.strategy】- Range(默认):【针对一个topic】首先对同一个topic里面的分区按照 序号进行排序,并对消费者按照字母顺序进行排序,然后用分区数除以 消费者数,得到每个消费者消费几个partition,然后按分区顺序连续分 配若干partition,除不尽的话 前面几个消费者会多分配一个分区的数据。
- 问题:如果只是针对1个topic,消费者0多消费一个分区影响不大; 但是如果有N个topic,那么消费者0就会多消费N个分区,那么就 容易发生数据倾斜
- 再平衡:挂掉了一个消费者之后,45秒以内重新发送消息,此时剩余的消费者暂时不能消费到挂掉的消费者应该消费的分区,等到了 45 秒以后,消费者就真正的挂掉了,此时会把它应该消费的分区数都分配给消费者1或者消费者2
- RoundRobin:【针对所有topic】首先将所有partition和consumer按照 一定顺序排列,然后按照 consumer依次分配排好序的partition,若该 consumer 没有订阅即将要分配的主题,那么直接跳过,继续向下分配
- 再平衡:也会进行轮询
- sticky:尽量均匀的分配分区给消费者(随机),黏性体现在在执行新的分配之前,考虑上一次的分配结果,尽量少的变动,这样就可以节省大量的开销
- 再平衡:均匀分配
kafka的消费方式¶
从拉取方式角度:pull vs push¶
Kafka 采用 pull(拉) 模式,而不是 push(推) 模式:
- pull模式(拉模式,Kafka采用)
- 消费者主动从 Broker 拉取数据,消费者控制消费速率
- 优点:可以根据消费者自身的消费能力以适当的速率消费消息,消费者处理快就拉得快,处理慢就拉得慢,不会压垮消费者
- 缺点:如果 Kafka 中没有数据,消费者可能会频繁循环拉取,一直返回空数据,浪费CPU资源
-
解决:Kafka 采用长轮询(Long Polling)机制,消费者拉取时如果没有数据,会阻塞等待一段时间(由
fetch.max.wait.ms配置),有数据到达或超时再返回,避免空轮询 -
push模式(推模式,Kafka不采用)
- 由 Broker 主动推送消息给消费者,Broker 决定消息发送速率
- 缺点:很难适应所有消费者的消费速率,如果消费者处理慢,Broker推送太快会把消费者压垮
- 适用场景:消息队列推送给多个订阅者,典型如 MQTT
从消费者组角度:集群消费 vs 广播消费¶
- 集群消费(默认)
- 同一个消费组内的多个消费者共同消费一个 Topic,每条消息只会被同一个消费组中的一个消费者消费
- 实现了负载均衡,提高消费能力
-
适用场景:需要提高消费吞吐量,消息只需处理一次
-
广播消费
- 同一个 Topic 的消息会被发送给所有消费者,每个消费者都能收到全量消息
- 适用场景:一条消息需要多个业务场景各自处理(比如订单创建后,通知库存服务和通知积分服务)
- Kafka 默认不直接支持广播消费,可以通过每个消费者使用不同的消费组ID来实现广播效果
从offset提交角度:自动提交 vs 手动提交¶
- 自动提交
- 配置:
enable.auto.commit=true(默认开启) - 消费者周期性自动提交offset,周期由
auto.commit.interval.ms配置,默认5秒 - 优点:使用简单,不需要手动管理
-
缺点:容易出现重复消费或漏消费
- 如果消费者处理消息过程中宕机,offset已经自动提交了,但消息没处理完,就会丢数据
- 如果消费者还没处理完下一批消息,自动提交时间到了,offset提前提交了,消费者宕机重启后从下一条开始消费,上一批没处理完的消息就丢了
-
手动提交
- 配置:
enable.auto.commit=false - 消费者处理完消息后手动调用API提交offset
- 优点:保证消息处理完成才提交offset,不会丢数据
- 分类:
- 同步提交:
consumer.commitSync(),阻塞等待提交完成,直到成功才继续拉取下一批 - 异步提交:
consumer.commitAsync(),不阻塞,提高吞吐量,但需要处理提交失败的情况
- 同步提交:
- 适用场景:对数据可靠性要求高的场景,需要保证消息处理完成才提交
从消息顺序角度:顺序消费 vs 并发消费¶
- 顺序消费
- Kafka 只能保证单个分区内消息有序,不能保证全局有序
- 如果需要全局有序,只能将 Topic 分区数设置为1
- 如果只需要局部有序(比如同一个订单的消息有序),可以通过将相同key的消息发送到同一个分区实现
-
同一个分区只能被同一个消费组中的一个消费者消费,保证分区内消费顺序
-
并发消费
- 多个消费者并发消费不同分区,提高消费吞吐量
- 不能保证消息全局有序,但每个分区内部仍然有序
- 这是Kafka最常用的消费方式
kafka 中的数据是有序的吗,如何保证有序的呢¶
kafka核心原则:只能保证单个partition内是有序的,但是跨partition的有序天然无法保证
- 单个partition内有序性的保证:(kafka 如何实现消息的有序的?)
- 生产者
- 写入顺序
- 分区的leader副本负责让数据以先进先出(FIFO)的顺序写入,来保证消息顺序性
- 重试机制
- 1.x版本之前:将允许最多没有返回
ack的次数参数max.in.flight.requests.per.connection设置为1 - 1.x版本之后:如果开启了幂等性,那么只要设置
max.in.flight小于等于5就可以了- 原因:启用幂等后,kafka服务端会缓存producer发来的最近5个request的元数据,因此无论如何,都可以保证最近5个request的数据都是有序的
- 消费者
- 同一个分区内的消息只能被一个group里的一个消费者消费,保证分区内消费有序
- 跨partition的消息顺序性的解决办法:
- 设置 topic 有且只有一个 partition:
partition=1- 优点:全局有序性保证
- 缺点:牺牲横向扩展能力,成为性能瓶颈
- 指定相同的分区号,或者使用相同的消息key
- 适用场景:局部有序的业务需求(如相同订单号的交易流水)
- 注意事项:需合理设计 Key 的分布以避免数据倾斜
kafka怎么保证数据一致性【如果kafka挂掉怎么保证数据不丢失,kafka如何保证数据不重复】¶
分别从如何保证数据不丢失以及数据不重复两个方面来回答
- 如何保证数据不丢失
- 生产者端:producer 发送数据到kafka的时候,当kafka接收到数据之后,需要向producer 发送 ack 确认收到,如果producer接收到ack,才会进行下一轮的发送,否则重新发送数据
- 什么时候发送ack呢?
- kafka提供了3种ack应答级别:
- ack=0,生产者发送过来的 数据,不需要等数据落盘就会应答,一般不会使用
- ack=1,生产者发送过来的数据,Leader收到数据后就会应答,因为这个级别也会丢失数据,所以一般用于传输普通日志
- ack=-1(默认级别),生产者发送过来的数据,Leader和ISR队列里面的所有节点收到数据后才会应答,不会丢失数据,一般用于传输和钱相关的数据
- 为什么提出了ISR队列?
- 因为如果要等待所有的follower 都同步完成,才发送ack,假设有一个follower迟迟不能同步,那怎么办呢?难道要一直等吗?因此就出现了ISR队列,这里面会存放和leader保持同步的follower集合,如果长时间(30s)未和leader通信或者同步数据,就会被踢出去
- 消费者端:消费者消费数据的时候会不断提交
offset【消费数据的偏移量】, 就是为了当系统挂掉后,下次可以从上次消费结束的位置继续消费- offset 在0.9 版本之前,保存在 Zookeeper 中;从0.9 版本开始,consumer 将 offset 保存在 Kafka一个内置的 topic 中【
__consumer_offsets】
- offset 在0.9 版本之前,保存在 Zookeeper 中;从0.9 版本开始,consumer 将 offset 保存在 Kafka一个内置的 topic 中【
- broker端:每个partition都会有多个副本
- 每个 broker 中的 partition 我们一般都会设置有 replication(副本) 的个数,生产者写入的时候首先根据分区分配策略(有 partition 按 partition,有 key 按 key,都没有轮询)写入到 leader 中,follower (副本)再跟 leader 同步数据,这样有了备份,也可以保证消息数 据的不丢失。
- 如何保证数据不重复【如何保证exactly once语义】
-
kafka提供了3种ack应答级别:ack=0,ack=1,ack=-1(默认级别)
-
生产者端(kafka已实现)
-
重复消费-ack应答失败:设置系统为
ack=-1后,假设leader收到数据并且同步 ISR 队列之后,在返回ack应答给producer端(生产者)之前leader挂掉了,那么producer端就会认为数据发送失败,再次重新发送,那么此时集群就会收到重复的数据,这样在生产环境中显然是有问题的 -
0.11版本之后,kafka提出了一个非常重要的特性:幂等性(默认是开启 的)。
-
幂等性的意思是无论producer发送多少次重复的数据,kafka只会持久化一条数据
-
将幂等性和至少一次语义【
ack=-1+副本数>=2+ISR 最小副本数>=2】结合在一起,就可以实现精确一次性(既 不丢失又不重复) -
底层原理:在producer刚启动的时候会分配一个PID,然后发送到同一个分区的消息都会携带一个单调自增的SequenceNum(单调递增保证不重复)【各分区独立】,broker会持久化PID+SequenceNum,也就是把它当做主键。如果有相同主键(也就是相同的SequenceNum)的消息提交时,broker只会持久化一条数据。但是这个机制只能保证单会话的精准一次性,如果想要保证跨会话的精准一次性,那么就需要事务的机制来进行保证【producer 在使用事务功能之前,必须先自定义一个唯一的事务id,这样,即使客户端重启,也能继续处理未完成的事务;并且这个事务的信息会持久化到一个特殊的主题topic当中】
-
消费者端(取决于消费者的消费策略)
-
重复消费-自动提交offset: ;consumer每5s自动提交offset,如果提交后的2s,consumer挂掉了,再次重启consumer,则从上一次提交 offset 处继续消费,导致重复消费
-
漏消费-手动提交offset:消费者消费的数据还在内存中,消费者挂掉了,导致漏消费
-
解决方式:手动提交offset + 采用消费者事务,比如mysql。也就是说kafka下游的消费者必须支持事务(能够回滚)
kafka的精准一次性¶
精准一次 = 至少一次 + 幂等性写入
至少一次需要将 ack 级别设置为-1 + 副本数>=2 + ISR 最小副本数>=2
幂等性在0.11版本后默认开启
幂等性原理: producer 刚启动的时候会分配一个 PID,发送到同一个分区的消息都会携带一个SequenceNum(单调自增的),broker 会对
0.11 版本之后,kafka 提出了一个非常重要的特性,幂等性(默认是开启的),也就是说无论 producer 发送多少次重复的数据,kafka 只会持久化一条数据,把这个特性和至少一次语义(ack 级别设置为-1+副本数>=2+ISR 最小副本数>=2)结合在一起,就可以实现精确一次性(既不丢失又不重复)。我大致介绍一下它的底层原理:在 producer 刚启动的时候会分配一个 PID,然后发送到同一个分区的消息都会携带一个SequenceNum(单调自增的),broker 会对
做缓存,也就是把它当做主键,如果有相同主键的消息提交时,broker 只会持久化一条数据。但是这个机制只能保证单会话的精准一次性,如果想要保证跨会话的精准一次性,那么就需要事务的机制来进行保证
Kafka 为什么可以扛住这么高的qps¶
Kafka能够处理高qps(每秒查询率)的原因主要有以下几点:
- 分布式架构 Kafka是一个分布式流处理平台,可以将数据和负载分散到多个服务器(broker)上,实现水平扩展。每个broker都可以独立处理效据和请求,通过增加broker数量可以线性提高系统整体的吞吐量
- 顺序读写 Kafka将消息追加到日志文件(logsegment)的末尾,实现了顺序写操作,顺序写的性能远离于随机写。在消费端,Kafka也支持顺序读磁盘,通过零拷贝技术(Zer0-Copy)直接将数据从磁盘发送到网络,减少了数据拷贝和上下文切换的开销,理高了数据输效率。
- 批量处理 Kafka在生产者和消费者端都支持批量发送和接收消息,可以减少网络传输的次数和开铜,提高整体吞吐量。生产者可以将多条消息打包成一个批次发送,消费者也可以一次性接收多个批次的清息进行处理
- 零拷贝技术 Kafka利用零拷贝技术减少了数据在内存中的拷贝次数,直接从磁盘读取数据并发送到同络,提高了数据俦输效率。
- 高效的IO结构 Kafka使用内存映射文件(Memory-Mapped Flle)来存储消息,将磁盘文件换射到内存中,可以直接访问内存中的致据,减少了感盘I/O的开销。Kafka的内部数据结构,如索引和偏移量(offset)管理,都经过优化,能够高效地定位和读取消息
- 异步处理 Kafka的生产者和消费者都支持异步处理,生产者可以异步发送消息,消费害可以异步处理消息,提高了系统的并发性和响应速度
- 分区和并行处理 Kafka的主题(topic)可以划分为多个分区(partition),每个分区可以独立处理数据和请求,实现了并行处理。生产者和消费者可以根据分区数量进行并行操作,进一步提高系统的吞吐量
- 消息压缩 Kafka支持消息压缩,可以将多个消息压缩成一个批次进行传输和存储,减少了网络传输和磁盘存储的开销。提高了数据传输和存储效率
※4个消费者消费3个分区,3个消费者消费4个分区,分别是什么现象?¶
在kafka中,消费者通过消费分区来并行处理消息,且每个分区只能被一个消费者消费,以避免消息重复处理。
-
4个消费者消费3个分区 由于每个分区只会被一个消费者消费,因此会有1个消费者处于空闲状态,不会消费任何消息,其他3个消费者会各自消费一个分区,实现消息的并行处理
-
3个消费者消费4个分区 由于每个分区只能被一个消费者消费,因此会有1个分区没有被消费,导致该分区的消息无法被及时处理。其他3个消责者会各自消费一个分区,实现消息的并行处理。
总结来说。消费者数量与分区数量不匹配时,会导致部分消费密空闲或分区来被消费,影消息处理的效率和及时性。
Kafka与 Rocketmq区别是什么?¶
Kafka和RocketMQ都是分布式消息队列系统,但它们在设计理念、应用场和特性上有所不同
| Kafka | RocketMQ | |
|---|---|---|
| 设计理念 | 最初由Linkedin开发,设计目标是高吞吐量、低延迟和分布式,Kafka强调消息的持久化和分区特性,适用于大规模数据处理和实时数据管道 | 由阿里巴巴开发,设计目标是高可用、高性能和可置性。RocketMQ强调消息的秩序性及事务性。适用于金融、电商等对消息可靠性要求较高的场景 |
| 消息模型 | 基于发布/订阅模型,消息发布到主题(Topic),消费者订阅主题并消费消息 | 支持发布/订阅模型和队列模型,消息可以被发送到主题或队列,消费者可以订阅主题或从队列中获取消息 |
| 消息顺序 | 保证分区内的消息顺序,但不保证跨分区的消息顺序 | 保证主题和队列内的消息顺序 |
| 消息可靠性 | 通过多副本机制(ISR集合)实现高可用,支持事务性消息(自0.11版本起),但需配合幂等性生产者使用 | 提供原生分布式事务消息支持(通过两阶段提交) |
| 性能 | 在高吞吐量场景下表现优异,适台大规模数据处理 | 在高并发、低延迟场景下表现优异,适合金融、电商等对消息实时性要求较高的场景 |
※说一下Rocketmq的消费方式¶
RocketMQ支持两种消费方式:
- 集群消费
多个消费者实例组成一个消费组,共同消费同一个主题的消息,每条消息只会被消费组中的一个消费者实例消费。适用于消息需要被多个消费者并行处理的场景
- 广播消费
每个消费者实例都会消费主题中的所有消息。适用于消息需要被所有消费者实例处理的场景。
※什么是死信队列?¶
死信队列是RocketMQ中的一种特殊队列,用于存储消费失败的消息。
当消息在消费过程中出现异常,达到最大重试次数后仍未被成功消费,该消息会被发送到死信队列。
死信队列可以帮助开发者定位和排查消息消费失败的原因,确保消息不被丢失。
※kafka auto.offset.reset 参数什么意思¶
用于控制消费者在没有初始偏移量或偏移量超出范围时的行为。 这个参数有三个可选值:latest、earliest 和 none。默认值是latest
三者均有共同定义: 对于同一个消费者组,若已有提交的offset,则从提交的offset开始接着消费
不同的点为:
- latest(默认):对于同一个消费者组,若没有提交过offset,则只消费`消费者连接topic后,新产生的数据
- earliest:对于同一个消费者组,若没有提交过offset,则从头开始消费
- none:对于同一个消费者组,若没有提交过offset,会抛异常
※kafka 的rebalance机制什么时候会触发¶
四种情况的时候,kafka会触发Rebalance机制:
- 消费组成员发生变更,比如有新的消费者加入了消费组组或者有消费者宕机
- 消费超时,消费者无法在指定的时间之内完成消息的消费
- 消费组订阅的 Topic 发生了变化
- 订阅的 Topic 的 partition 发生了变化
kafka的消费原则,消费者和分区的关系?¶
一个消费者可以消费多个分区的数据,每个分区只能被同一个消费组中的一个消费者消费
Kafka消费者如何实现“至少一次”语义?写出关键配置参数及代码片段。¶
项目¶
Kafka在项目中起到的作用¶
有的,Kafka在项目里可以说是核心的数据总线,贯穿了整个实时数仓的上下游。可以毫不夸张地说,它起到了“承上启下、解耦分层”的关键作用,我分几点给您讲讲它的具体角色:
第一,作为 ODS 层(操作数据存储层)的数据底座。 我们所有原始数据都会先打到 Kafka 里。用户行为日志通过 Flume 写到 topiclog,业务数据库的 Binlog 通过 Maxwell 写到 topicdb。Kafka 充当了数据的“削峰填谷”缓冲区,这样就算业务高峰期瞬间流量很大,也不会直接把下游 Flink 任务冲垮,同时给了我们数据重播回溯的能力。
第二,作为 DWD 层(明细数据层)的“解耦神器”。 我们 DWD 层有很多 Flink 任务,比如流量域分流、交易域下单、支付等。上游的 ODS 层写完 Kafka 后,下游不同的 Flink 任务可以按照自己的消费速度去消费同一个 Topic,互不影响。而且,DWD 层处理完的结果,比如 dwdtradeorder_detail,我们也是写回 Kafka 的。这样 DWS 层(汇总层)再去消费这个明细主题,实现了计算层的解耦,哪怕 DWS 层挂了重启,只要 Kafka 里的数据还在,就能重新消费,不会影响上游。
第三,利用 Kafka 的多分区特性保证数据有序性。 这是个细节,但我们特意做了优化。项目里我们将 Kafka 的主题分区数统一设置为 4,同时 Flink 任务的并行度也设置为 4,这样每个 Flink 并行度就固定消费一个分区。Kafka 单分区内是严格有序的,这为我们后续做订单去重、状态修复提供了非常便利的前提,不需要在 Flink 内部再额外做复杂的排序缓存。
第四,配合 Checkpoint 实现端到端的精确一次(Exactly-Once)。 我们在写 Kafka 的时候,用了 Flink 的 KafkaSink 并开启了 DeliveryGuarantee.EXACTLY_ONCE(精确一次语义),配合 setTransactionalIdPrefix 设置了事务 ID。这样当 Flink 任务发生故障重启时,通过事务机制和状态回滚,保证了数据既不会丢也不会重复,这对于实时数仓的报表准确性至关重要。
总之,Kafka 在我们项目里绝对不只是个消息管道,它既是大数据的临时仓库,也是连接各个计算层的桥梁,更是保证数据一致性和可靠性的基石。