Kafka读取者读取数据流
发布时间
阅读量:
阅读量
目录
- 消费组与分区的再平衡机制
-
构建Kafka消费者实例
-
注册主题订阅
-
消息拉取流程循环
-
消费者参数配置
-
位移提交(commit)与偏移量(offset)管理
- 自动位移提交机制
- 当前偏移量的提交操作
- 异步方式的位移提交
-
重平衡事件监听器(Rebalance Listener)
-
通过KafkaConsumer实现对特定主题的消息订阅,从而获取并处理相关消息。Kafka消费者作为消费组中的成员,当多个消费者共同组成一个消费组以处理同一主题时,每个消费者将被分配到不同的分区,并接收对应分区内的消息。
消费组与分区重平衡机制
当有新的用户进入消费组时,其将接管一个或多个原本由其他用户处理的分区;同时,若某个用户退出消费组(例如发生重启或故障等情况),其之前负责的分区将被重新分配给其他用户。这一过程被称为重平衡(rebalance)。
为维持在消费组中的活跃状态,用户需定期向作为组协调者的broker发送心跳信号(heartbeat)。该broker并非固定不变,不同消费组可能对应不同的broker。每当用户从broker拉取数据或提交信息时,都会同步发送心跳信号。一旦用户在规定时间内未能发送心跳,其会话(session)将被视为超时,此时组协调者会判定该用户已失效
全部评论 (0)
还没有任何评论哟~
