
消息路由策略
生产者是以 Record 为消息进行发布,Record 包含了 Key-value。其中主要是利用key 来做消息路由策略
若指定了某个和partition, 则写入到特定的partition
public ProducerRecord (String topic , Integer partition, K key, V value)若有指定key,则用key来做取模选出partition
public ProducerRecord(String topic, K key, V value)若以上都未指定,则采用轮询存入
public ProducerRecord(String topic, V value)
消息写入过程
Producer 向 kafka broker 集群提交连接请求, 任意节点的broker都会提供一个 broker Controller的地址, 有了 broker Controller的地址 就可以通过他来监听所有 broker的信息 当要写入消息到某个topic时, 可以通过 broker Controller 在zookeeper中获取当前topic 的所有 partition leader 的地址, 用于写入消息. partition leader 写入本地的log文件后, 会通知 partition follow 进行同步,同步完成向leader发送ACK 收到所有follow的ACK之后, 会增加相应的HW, 代表当前所有节点最高的同步点, 意味着消费者只能消费到这个点 假如follow回复ACK超时了, 那么就把这个follow在ISR列表中除去, 然后再增加HW
消费者消费过程解析
consumer 向 broke 集群提交连接请求,任意broker 会返回 broker controller 的地址 consumer 把需要消费的topic 告诉broke controller 后会分配相应数量的 partition leader , 并将该partition 的当前offset 发送给consumer 确定消费起点 consumer 根据offset和相应的 partition 进行消费,消费完成后更新offset缓存和远程的 _comsumer_offset ,并且一直重复本步骤,知道消费者停止请求消息 消费者可以重置offset 从而灵活的消费broker上的任意消息
HW HighWatermark,高水位机制
表示当前所有 partition 能同步到的最高消息,消费者最高能消费到这里.
保证了所有partition leader 和 follow 的消息一致性
LEO: Log End Offset. 消息被写入到当前partition log文件的最后偏移量, HW会停留在所有partition集群中最低LEO的位置.
当所有follow的LEO都同步完成, 代表最低LEO增加了, 此时HW 也会增加, 那么consumer才会消费

HW 消息截断机制
当发生leader A 宕机后, B当选新Leader , 当A再次重启会导致消息不一致, 此时需要截断掉A宕机前最后同步的消息, 保证旧leader消息不一致的情况.
当然如果截断会导致消息丢失 (A 的6号消息就丢失了)

Partition Leader 选举范围
当leader挂了之后,broker controller 会从 ISR 中重新选举出新的leader,但当ISR中已经无副本
此时通过unclean.leader.election.enable 设置是否从OSR中选出leader
true :允许任意的partition作为新的leader 万一选中了 OSR 的话,会出现消息不一致,但保证了可用性 false:需要等待所有的ISR都活过来才重新选举
重复消费的解决方案
同个consumer重复消费
一般是由于,消息的消费时长超时导致消息重发不同的consumer重复消费
consumer消费'消息A'后没来得及提交offset 宕机,下次消费交给其他的consumer由于没有读取到'消息A'的消费记录,所以重复消费了'消息A'
所以综上所余, kafka 无法保证消息的绝对可靠性,需要我们进行消息幂等的校验

文章转载自阿哲是哲学的哲,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。




