
1. 为什么需要consumer group 好处是什么1.实际上consumer group是用于实现高伸缩性、高容错性的consumer机制。2.组内多个 consumer实例可以同时读取Kafka消息而且一旦有某个 consumer“挂”了consumer group会立即将已崩溃 consumer负责的分区转交给其他 consumer来负责从而保证整个group可以继续工作不会丢失数据——这个过程被称为重平衡rebalance。2 消费组和消息的顺序性关系另外由于 Kafka目前只提供单个分区内的消息顺序而不会维护全局的消息顺序因此如果用户要实现 topic 全局的消息读取顺序就只能通过让每个 consumer group 下只包含一个consumer实例的方式来间接实现。3 consumer offset1.consumer端的 offset与分区日志中的 offset是不同的含义。2.每个 consumer 实例都会为它消费的分区维护属于自己的位置信息来记录当前消费了多少条消息。很多消息引擎都把消费端的 offset 保存在服务器端broker,这样做的好处当然是实现简单但会有以下3个方面的问题。1.broker从此变成了有状态的增加了同步成本影响伸缩性2.需要引入应答机制acknowledgement来确认消费成功3.由于要保存许多 consumer 的 offset故必然引入复杂的数据结构从而造成不必要的资源浪费。Kafka则选择了不同的方式让 consumer group保存 offset那么只需要简单地保存一个长整型数据就可以了同时 Kafka consumer 还引入了检查点机制checkpointing定期对offset进行持久化从而简化了应答机制的实现。从下图中我们可以看到当前Kafka consumer在内部使用一个 map来保存其订阅topic所属分区的offset。4 offset提交consumer客户端需要定期地向Kafka集群汇报自己消费数据的进度这一过程被称为位移提交offset commit。新版本和旧版本 consumer提交位移的方式截然不同旧版本 consumer会定期将位移信息提交到ZooKeeper下的固定节点上把位移提交到ZooKeeper的做法并不合适。ZooKeeper本质上只是一个协调服务组件它并不适合作为位移信息的存储组件毕竟频繁高并发的读/写操作并不是 ZooKeeper擅长的事情。新版本consumer把位移提交到 Kafka 的一个内部 topic__consumer_offsets上通常不能直接操作该topic就可以了特别是注意不要擅自删除或搬移该topic的日志文件。5 _consumer_offsets_consumer_offsets是Kafka自行创建的因此用户不可擅自删除该 topic的所有信息。通常情况下这样的文件夹应该有50个编号从0到49。打开图中的任意一个文件夹会发现它就是一个正常的Kafkatopic日志文件目录里面至少有一个日志文件.log和两个索引文件.index 和.timeindex。该日志中保存的消息都是 Kafka 集群上consumer特别是 consumer group的位移信息罢了。_consumer_offsets的每条消息格式大致如图你可以把它想象成一个KV格式的消息key就是一个三元组group.id topic 分区号而value就是offset的值。每当更新同一个 key的最新offset值时该topic就会写入一条含有最新 offset的消息同时 Kafka会定期地对该 topic执行压实操作compact即为每个消息key 只保存含有最新offset的消息。这样既避免了对分区日志消息的修改也控制住了_consumer_offsets topic总体的日志容量同时还能实时反映最新的消费进度。考虑到一个Kafka生产环境中可能有很多consumer或consumer group如果这些consumer同时提交位移则必将加重__consumer_offsets的写入负载因此社区特意为该topic创建了50个分区并且对每个group.id做哈希求模运算从而将负载分散到不同的__consumer_offsets分区上。这就是说每个consumer group保存的offset都有极大的概率分别出现在该topic的不同分区上。6 消费者组重平衡如果使用的是 standalone consumer则压根就没有rebalance的概念即rebalance只对consumer group有效。何为 rebalance它本质上是一种协议规定了一个 consumer group下所有 consumer如何达成一致来分配订阅 topic的所有分区。举个例子假设我们有一个 consumer group它有20个 consumer实例。该 group订阅了一个具有100个分区的 topic。那么正常情况下consumer group平均会为每个consumer分配5个分区即每个 consumer负责读取5个分区的数据。这个分配过程就被称作rebalance。7 consumer主要参数7.1 session.timeout.mssession.timeout.ms是consumer group检测组内成员发送崩溃的时间。假设你设置该参数为5分钟那么当某个group成员突然崩溃了比如被kill-9或宕机管理 group的 Kafka 组件即消费者组协调者也称 group coordinator有可能需要 5 分钟才能感知到这个崩溃。显然我们想要缩短这个时间让coordinator 能够更快地检测到 consumer 失败。遗憾的是这个参数还有另外一重含义consumer消息处理逻辑的最大时间——倘若consumer两次poll之间的间隔超过了该参数所设置的阈值那么coordinator 就会认为这个 consumer 已经追不上组内其他成员的消费进度了因此会将该consumer实例“踢出”组该consumer负责的分区也会被分配给其他consumer。在最好的情况下这会导致不必要的rebalance因为consumer需要重新加入group。更糟的是对于那些在被踢出group后处理的消息consumer都无法提交位移——这就意味着这些消息在rebalance之后会被重新消费一遍。如果一条消息或一组消息总是需要花费很长的时间处理那么consumer甚至无法执行任何消费除非用户重新调整参数。鉴于以上的“窘境”,Kafka社区于0.10.1.0版本对该参数的含义进行了拆分。在该版本及以后的版本中session.timeout.ms 参数被明确为“coordinator 检测失败的时间”。因此在实际使用中用户可以为该参数设置一个比较小的值让 coordinator能够更快地检测 consumer崩溃的情况从而更快地开启 rebalance避免造成更大的消费滞后consumer lag。目前该参数的默认值是10秒。7.2 max.poll.interval.ms如前所述session.timeout.ms 中“consumer 处理逻辑最大时间”的含义被剥离出来了Kafka为这部分含义单独开放了一个参数——max.poll.interval.ms。通过将该参数设置成实际的逻辑处理时间再结合较低的session.timeout.ms 参数值consumer group既实现了快速的consumer崩溃检测也保证了复杂的事件处理逻辑不会造成不必要的rebalance。7.3 auto.offset.reset指定了无位移信息或位移越界即 consumer 要消费的消息的位移不在当前消息日志的合理区间范围时 Kafka的应对策略。特别要注意这里的无位移信息或位移越界只有满足这两个条件中的任何一个时该参数才有效果。举例说明假设你首次运行一个consumer group并且指定从头消费。显然该group会从头消费所有数据因为此时该 group 还没有任何位移信息。一旦该 group 成功提交位移后你重启了 group依然指定从头消费。此时你会发现该 group并不会真的从头消费——因为Kafka已经保存了该group的位移信息因此它会无视auto.offset.reset的设置。该参数有如下3个可能的取值1。earliest指定从最早的位移开始消费。注意这里最早的位移不一定就是0。2.latest指定从最新处位移开始消费3.none指定如果未发现位移信息或位移越界则抛出异常。在实际使用过程中几乎从未见过将该参数设置为none的用法因此该值在真实业务场景中使用甚少。7.4 enable.auto.commit该参数指定 consumer是否自动提交位移。若设置为 true则 consumer在后台自动提交位移否则用户需要手动提交位移。7.5 fetch.max.bytes指定了 consumer 端单次获取数据的最大字节数。若实际业务消息很大则必须要设置该参数为一个较大的值否则consumer将无法消费这些消息。7.6 max.poll.records该参数控制单次 poll调用返回的最大消息数。比较极端的做法是设置该参数为1那么每次 poll只会返回1条消息。如果用户发现 consumer端的瓶颈在 poll速度太慢可以适当地增加该参数的值。如果用户的消息处理逻辑很轻量默认的500条消息通常不能满足实际的消息处理速度。7.7 heartbeat.interval.ms要搞清楚consumergroup的其他成员如何得知要开启新一轮rebalance——当coordinator决定开启新一轮rebalance时它会将这个决定以REBALANCE_IN_PROGRESS异常的形式“塞进”consumer心跳请求的response中这样其他成员拿到response后才能知道它需要重新加入group。显然这个过程越快越好而heartbeat.interval.ms就是用来做这件事情的。比较推荐的做法是设置一个比较低的值让 group 下的其他 consumer成员能够更快地感知新一轮rebalance开启了。注意该值必须小于session.timeout.ms这很容易理解毕竟如果consumer在session.timeout.ms这段时间内都不发送心跳coordinator就会认为它已经dead因此也就没有必要让它知晓coordinator的决定了。7.8 connections.max.idle.ms经常有用户抱怨在生产环境下周期性地观测到请求平均处理时间在飙升这很有可能是因为 Kafka会定期地关闭空闲Socket连接导致下次consumer处理请求时需要重新创建连向broker的Socket连接。当前默认值是9分钟如果用户实际环境中不在乎这些Socket资源开销比较推荐设置该参数值为-1即不要关闭这些空闲连接。在zk的bin目录下启动客户端脚本查看节点信息:./zkCli.sh执行命令查看节点数据:ls /查看kafka集群配置信息: