工业数字化—IoT建设学习笔记(4):大数据应用——Kafka
目录
1、kafka概述
Kafka 起初是 由 LinkedIn 公司采用 Scala 语言开发的一个多分区、多副本且基于 ZooKeeper 协调的分布式消息系统,现已被捐献给 Apache 基金会。目前 Kafka 已经定位为一个分布式流式处理平台。它以高吞吐、可持久化、可水平扩展、支持流数据处理等多种特性而被广泛使用,主要是由 Scala 和 Java 编写。

它以一种高吞吐量的发布订阅机制,可以处理事件流数据。
Kafka 的基本术语
- 消息:Kafka 中的数据单元被称为消息,也被称为记录,可以把它看作数据库表中某一行的记录。
- 主题:消息的种类称为主题(Topic),相当于是对消息进行分类,可以说一个主题代表了一类消息。主题就像是数据库中的表。
- 分区:主题可以被分为若干个分区(partition),同一个主题中的分区可以在不同服务器上,由此来实现 kafka 的伸缩性。分区是有序的
-
生产者: 向主题发布消息的客户端应用程序称为生产者(Producer)
-
消费者:订阅主题消息的客户端程序称为消费者(Consumer),生产者可以与消费者相互转换。
-
消费者群组:生产者与消费者的关系就如同餐厅中的厨师和顾客之间的关系一样,一个厨师对应多个顾客,也就是一个生产者对应多个消费者,消费者群组(Consumer Group)指的就是由多个消费者组成的群体。
通过 Kafka 你可以非常方便通过“发布-订阅”的方式,把想要发布的消息分发给任何想要订阅该消息的接收者。生产者只需要把消息输入到 Kafka 指定 Topic ,消费者只要订阅该 Topic ,就能低延时、高吞吐量的接收到上游的消息;Kafka 还支持同一个 Topic 同时被多个下游消费者消费,且不同消费者之间数据处理进度互不干扰。

- 偏移量:偏移量(Consumer Offset)是一种元数据,它是一个不断递增的整数值,用来记录消费者发生重平衡时的位置,以便用来恢复数据。
- broker: 一个独立的 Kafka 服务器就被称为 broker,broker 接收来自生产者的消息,为消息设置偏移量,并提交消息到磁盘保存。
- broker 集群:broker 是集群 的组成部分,broker 集群由一个或多个 broker 组成,每个集群都有一个 broker 同时充当了集群控制器的角色(自动从集群的活跃成员中选举出来)。
- 副本:Kafka 中消息的备份又叫做 副本(Replica),副本的数量是可以配置的,Kafka 定义了两类副本:领导者副本(Leader Replica) 和 追随者副本(Follower Replica),前者对外提供服务,后者只是被动跟随,一个topic的每个分区都有若干个副本。
- 重平衡:Rebalance。消费者组内某个消费者实例挂掉后,其他消费者实例自动重新分配订阅主题分区的过程。Rebalance 是 Kafka 消费者端实现高可用的重要手段。

2、kafka基础架构

如上图所示,一个典型的 Kafka 集群中包含若干Producer(可以是web前端产生的Page View,或者是服务器日志,系统CPU、Memory等),若干broker(Kafka支持水平扩展,一般broker数量越多,集群吞吐率越高),若干Consumer Group,以及一个Zookeeper集群。Kafka通过Zookeeper管理集群配置,选举leader,以及在Consumer Group发生变化时进行rebalance。Producer使用push模式将消息发布到broker,Consumer使用pull模式从broker订阅并消费消息。
3、Kafka的Producer
在 Kafka 中,我们把产生消息的那一方称为生产者,比如我们经常回去淘宝购物,你打开淘宝的那一刻,你的登陆信息,登陆次数都会作为消息传输到 Kafka 后台,当你浏览购物的时候,你的浏览信息,你的搜索指数,你的购物爱好都会作为一个个消息传递给 Kafka 后台,然后淘宝会根据你的爱好做智能推荐,致使你的钱包从来都禁不住诱惑,那么这些生产者产生的消息是怎么传到 Kafka 应用程序的呢?发送过程是怎么样的呢?
尽管消息的产生非常简单,但是消息的发送过程还是比较复杂的,如图

我们从创建一个ProducerRecord 对象开始,ProducerRecord 是 Kafka 中的一个核心类,它代表了一组 Kafka 需要发送的 key/value 键值对,它由记录要发送到的主题名称(Topic Name),可选的分区号(Partition Number)以及可选的键值对构成。
在发送 ProducerRecord 时,我们需要将键值对对象由序列化器转换为字节数组,这样它们才能够在网络上传输。然后消息到达了分区器。
如果发送过程中指定了有效的分区号,那么在发送记录时将使用该分区。如果发送过程中未指定分区,则将使用key 的 hash 函数映射指定一个分区。如果发送的过程中既没有分区号也没有,则将以循环的方式分配一个分区。选好分区后,生产者就知道向哪个主题和分区发送数据了。
ProducerRecord 还有关联的时间戳,如果用户没有提供时间戳,那么生产者将会在记录中使用当前的时间作为时间戳。Kafka 最终使用的时间戳取决于 topic 主题配置的时间戳类型。
- 如果将主题配置为使用 CreateTime,则生产者记录中的时间戳将由 broker 使用。
- 如果将主题配置为使用LogAppendTime,则生产者记录中的时间戳在将消息添加到其日志中时,将由 broker 重写。
然后,这条消息被存放在一个记录批次里,这个批次里的所有消息会被发送到相同的主题和分区上。由一个独立的线程负责把它们发到 Kafka Broker 上。
Kafka Broker 在收到消息时会返回一个响应,如果写入成功,会返回一个 RecordMetaData 对象,它包含了主题和分区信息,以及记录在分区里的偏移量,上面两种的时间戳类型也会返回给用户。如果写入失败,会返回一个错误。生产者在收到错误之后会尝试重新发送消息,几次之后如果还是失败的话,就返回错误消息。
4、Kafka的Consumer
应用程序使用 KafkaConsumer 从 Kafka 中订阅主题并接收来自这些主题的消息,然后再把他们保存起来。应用程序首先需要创建一个 KafkaConsumer 对象,订阅主题并开始接受消息,验证消息并保存结果。一段时间后,生产者往主题写入的速度超过了应用程序验证数据的速度,这时候该如何处理?如果只使用单个消费者的话,应用程序会跟不上消息生成的速度,就像多个生产者像相同的主题写入消息一样,这时候就需要多个消费者共同参与消费主题中的消息,对消息进行分流处理。
Kafka 消费者从属于消费者群组。一个群组中的消费者订阅的都是相同的主题,每个消费者接收主题一部分分区的消息。下面是一个 Kafka 分区消费示意图

上图中的主题 T1 有四个分区,分别是分区0、分区1、分区2、分区3,我们创建一个消费者群组1,消费者群组中只有一个消费者,它订阅主题T1,接收到 T1 中的全部消息。由于一个消费者处理四个生产者发送到分区的消息,压力有些大,需要帮手来帮忙分担任务,于是就演变为下图

这样一来,消费者的消费能力就大大提高了,但是在某些环境下比如用户产生消息特别多的时候,生产者产生的消息仍旧让消费者吃不消,那就继续增加消费者。
- 一个分区只能被组内的一个消费者消费,不会出现多个消费者同时消费一个分区的情况。
- 消费者组会遵循 分区重平衡(Rebalance) 机制,将主题的所有分区均匀分配给组内的消费者。

如上图所示,每个分区所产生的消息能够被每个消费者群组中的消费者消费,如果向消费者群组中增加更多的消费者,那么多余的消费者将会闲置,如下图所示

向群组中增加消费者是横向伸缩消费能力的主要方式。总而言之,我们可以通过增加消费组的消费者来进行水平扩展提升消费能力。这也是为什么建议创建主题时使用比较多的分区数,这样可以在消费负载高的情况下增加消费者来提升性能。另外,消费者的数量不应该比分区数多,因为多出来的消费者是空闲的,没有任何帮助。
Kafka 一个很重要的特性就是,只需写入一次消息,可以支持任意多的应用读取这个消息。换句话说,每个应用都可以读到全量的消息。为了使得每个应用都能读到全量消息,应用需要有不同的消费组。对于上面的例子,假如我们新增了一个新的消费组 G2,而这个消费组有两个消费者,那么就演变为下图这样

在这个场景中,消费组 G1 和消费组 G2 都能收到 T1 主题的全量消息,在逻辑意义上来说它们属于不同的应用。
那么问题来了,同一消费组里的不同消费者不可以消费同一分区,而不同消费组里不同消费者可以去消费同一分区,这是为啥?他们有啥区别?
原因主要为以下两点:
1. 同一消费者组:共享 Offset,分区独占
- 消费者组的核心是 “一个分区只能被组内一个消费者独占”,背后依赖 Offset 共享机制:
- 整个消费者组共用一套针对分区的 Offset 记录,这个 Offset 标记了整个组消费到了分区的哪个位置。
- 如果两个消费者同时消费一个分区,就会出现 Offset 冲突:比如消费者 A 消费了消息 M1 并提交 Offset=1,消费者 B 又去消费 M1 并提交 Offset=1,会导致 Offset 混乱,重复消费或消息丢失。
- 同时,为了保证分区内消息的消费顺序,必须让一个消费者从头到尾处理该分区的消息,避免乱序。
Kafka 的数据消费,其底线是一定是保障消息消费的有序性!!!同一消费组里的不同消费者消费同一分区会导致消息顺序混乱
2. 不同消费者组:独立 Offset,并行读取
- 不同消费者组之间 Offset 完全隔离,每个组都有自己的一套 Offset 记录:
- 消费者组 A 消费分区 P0 时,更新的是组 A 专属的 Offset;消费者组 B 消费 P0 时,更新的是组 B 专属的 Offset。
- 两个组的 Offset 互不干扰,Broker 只需要根据每个组的 Offset 位置,返回对应的数据即可,不存在冲突问题。
- 这种设计的本质是:Kafka 分区的数据是 “只读” 的多副本数据源,不同业务方可以独立读取,就像多个人同时读同一本书一样,彼此不影响。
5、消费者组和分区重平衡
5.1 消费者是什么?
消费者组(Consumer Group)是由一个或多个消费者实例(Consumer Instance)组成的群组,具有可扩展性和可容错性的一种机制。消费者组内的消费者共享一个消费者组ID,这个ID 也叫做 Group ID,组内的消费者共同对一个主题进行订阅和消费,同一个组中的消费者只能消费一个分区的消息,多余的消费者会闲置,派不上用场。
我们在上面提到了两种消费方式
- 一个消费者群组消费一个主题中的消息,这种消费模式又称为点对点的消费方式,点对点的消费方式又被称为消息队列
- 一个主题中的消息被多个消费者群组共同消费,这种消费模式又称为发布-订阅模式
5.2 消费者重平衡
我们从上面的消费者演变图中可以知道这么一个过程:最初是一个消费者订阅一个主题并消费其全部分区的消息,后来有一个消费者加入群组,随后又有更多的消费者加入群组,而新加入的消费者实例分摊了最初消费者的部分消息,这种把分区的所有权通过一个消费者转到其他消费者的行为称为重平衡,英文名也叫做 Rebalance 。如下图所示

重平衡非常重要,它为消费者群组带来了高可用性 和 伸缩性,我们可以放心的添加消费者或移除消费者,不过在正常情况下我们并不希望发生这样的行为。在重平衡期间,消费者无法读取消息,造成整个消费者组在重平衡的期间都不可用。另外,当分区被重新分配给另一个消费者时,消息当前的读取状态会丢失,它有可能还需要去刷新缓存,在它重新恢复状态之前会拖慢应用程序。
消费者通过向组织协调者(Kafka Broker)发送心跳来维护自己是消费者组的一员并确认其拥有的分区。对于不同不的消费群体来说,其组织协调者可以是不同的。只要消费者定期发送心跳,就会认为消费者是存活的并处理其分区中的消息。当消费者检索记录或者提交它所消费的记录时就会发送心跳。
如果过了一段时间 Kafka 停止发送心跳了,会话(Session)就会过期,组织协调者就会认为这个 Consumer 已经死亡,就会触发一次重平衡。如果消费者宕机并且停止发送消息,组织协调者会等待几秒钟,确认它死亡了才会触发重平衡。在这段时间里,死亡的消费者将不处理任何消息。在清理消费者时,消费者将通知协调者它要离开群组,组织协调者会触发一次重平衡,尽量降低处理停顿。
重平衡是一把双刃剑,它为消费者群组带来高可用性和伸缩性的同时,还有有一些明显的缺点(bug),而这些 bug 到现在社区还无法修改。
重平衡的过程对消费者组有极大的影响。因为每次重平衡过程中都会导致万物静止,参考 JVM 中的垃圾回收机制,也就是 Stop The World ,STW,(引用自《深入理解 Java 虚拟机》中 p76 关于 Serial 收集器的描述):
更重要的是它在进行垃圾收集时,必须暂停其他所有的工作线程。直到它收集结束。Stop The World 这个名字听起来很帅,但这项工作实际上是由虚拟机在后台自动发起并完成的,在用户不可见的情况下把用户正常工作的线程全部停掉,这对很多应用来说都是难以接受的。
也就是说,在重平衡期间,消费者组中的消费者实例都会停止消费,等待重平衡的完成。而且重平衡这个过程很慢......
6、Kafka幂等性
6.1 什么是 Kafka 的幂等性
简单来说,Kafka 的幂等性是指生产者(Producer)向 Kafka 发送消息时,即使因为网络重试等原因重复发送了同一条消息,Kafka 也只会持久化一条,不会出现重复消息的情况。
你可以把它理解为:无论调用多少次 “发送消息” 这个操作,结果都和只调用一次完全一样,不会因为重复操作产生副作用(比如重复数据)。
6.2 幂等性的实现原理
Kafka 幂等性的核心是通过 Producer ID(PID) + 序列号(Sequence Number) 来实现的,具体逻辑如下:
- Producer ID(PID):每个生产者实例启动时,Kafka 会分配一个唯一的 PID(对用户透明),生产者重启后 PID 会变化。
- 序列号(Sequence Number):对于生产者向同一个分区(Partition)发送的消息,会按顺序分配一个递增的序列号(从 0 开始)。
- Broker 端校验:Broker 会为每个 <PID, 主题 - 分区> 维护一个最新的序列号:
- 当收到消息时,若消息的序列号 = 最新序列号 + 1 → 正常写入,并更新最新序列号;
- 若消息的序列号 ≤ 最新序列号 → 判定为重复消息,直接丢弃,不写入。
更多推荐
所有评论(0)