Advertisement

Kafka读取者读取数据流

阅读量:

目录

  • 消费组与分区的再平衡机制
    • 构建Kafka消费者实例

    • 注册主题订阅

    • 消息拉取流程循环

    • 消费者参数配置

    • 位移提交(commit)与偏移量(offset)管理

      • 自动位移提交机制
      • 当前偏移量的提交操作
      • 异步方式的位移提交
    • 重平衡事件监听器(Rebalance Listener)

通过KafkaConsumer实现对特定主题的消息订阅,从而获取并处理相关消息。Kafka消费者作为消费组中的成员,当多个消费者共同组成一个消费组以处理同一主题时,每个消费者将被分配到不同的分区,并接收对应分区内的消息。

消费组与分区重平衡机制

当有新的用户进入消费组时,其将接管一个或多个原本由其他用户处理的分区;同时,若某个用户退出消费组(例如发生重启或故障等情况),其之前负责的分区将被重新分配给其他用户。这一过程被称为重平衡(rebalance)。
为维持在消费组中的活跃状态,用户需定期向作为组协调者的broker发送心跳信号(heartbeat)。该broker并非固定不变,不同消费组可能对应不同的broker。每当用户从broker拉取数据或提交信息时,都会同步发送心跳信号。一旦用户在规定时间内未能发送心跳,其会话(session)将被视为超时,此时组协调者会判定该用户已失效

全部评论 (0)

还没有任何评论哟~