阶段四:消息中间件之Kafka
学习目标
- [ ] 了解MQ消息队列是什么?
- [ ] 了解Kafka的三大主体、三大内容;
- [ ] 知道分片(Partition)和副本(Replica)的作用
- [ ] 熟练掌握Kafka集群的安装部署;
- [x] 了解Python调用Kafka的相关模块使用方法;
Kafka 背景
Kafka是由LinkedIn开发并开源的分布式消息系统,因其分布式及高吞吐率而被广泛使用,现已与Cloudera Hadoop,Apache Storm,Apache Spark集成。
kafka名字灵感来自卡夫卡小说(《变形记》作者),寓意:数据像小说情节一样有序流动但错综复杂。

2010年,Kafka诞生。2011年,正式开源。2012年:成为Apache顶级项目。LinkedIn用Kafka解决了自己的数据洪流问题,却意外创造了改变互联网基础设施的工具。
2015,国内开始引进;2019,巅峰;2023年,AI崛起。
MQ 消息队列
MQ(Message Queue,消息队列)就像是一个“快递站”,负责把消息从一个系统传递到另一个系统,确保数据不丢失、不重复,还能应对高并发流量。
没有MQ:
大量用户点击 → 服务器(崩溃)
有MQ:
用户点击 → MQ(缓冲) → 每秒放100个请求给服务器想象你是一个卖家(生产者),每天要发很多快递给买家(消费者)。 没有驿站(MQ): 你每次都要亲自联系快递员,等他上门取件。 如果快递员忙不过来,你的包裹就会堆积,甚至丢失。 有驿站(MQ): 你把包裹送到驿站,驿站帮你暂存、排队,快递员有空时来取。 好处:你不用等快递员,快递员也不用立刻接单,双方更高效。
因此,我们总结一下,关于MQ的核心功能。
<sheet sheet-id="fCeT58" token="VF7fscIFthTIBBtJYbCcExobnmb"></sheet>
点对点消息系统
在点对点系统中,消息被保留在队列中。 一个或多个消费者可以消耗队列中的消息,但是特定消息只能由最多一个消费者消费。 一旦消费者读取队列中的消息,它就从该队列中消失。

发布-订阅消息系统
在发布 - 订阅系统中,消息被保留在主题中。 与点对点系统不同,消费者可以订阅一个或多个主题并使用该主题中的所有消息。
在发布 - 订阅系统中,消息生产者称为发布者,消息使用者称为订阅者。
一个现实生活的例子是Dish电视,它发布不同的渠道,如运动,电影,音乐等,任何人都可以订阅自己的频道集,并获得他们订阅的频道时可用。

Kafka 概述
Kafka 介绍
Kafka是一种分布式的,基于发布/订阅的消息系统。主要设计目标如下:
- 以时间复杂度为O(1)的方式提供消息持久化能力,即使对TB级以上数据也能保证常量级复杂度的性能
- 高吞吐率。即使在非常廉价的商用机器上也能做到单机支持每秒100K条以上消息的传输
- 支持Kafka Server间的消息分区,及分布式消费,同时保证每个Partition内的消息顺序传输
- 同时支持离线数据处理和实时数据处理
- Scale out:支持在线水平扩展
用人体来类比,消息队列(MQ)就是数字世界的血液输送网络,而Kafka则是其中的"超级主动脉"。
<sheet sheet-id="NILF1U" token="VF7fscIFthTIBBtJYbCcExobnmb"></sheet>
Kafka的官网:http://kafka.apache.org

相关术语

为什么会有 Zookeeper Cluster 集群? 在传统 Kafka 架构中,ZooKeeper 是 “集群大脑”,负责元数据存储、控制器选举、Leader 选举、消费组协调等核心分布式协调工作。
三大主体
Brokers(kafka服务)
- Kafka 集群中的服务器节点,负责存储和处理消息。
Producers(生产者)
- 向 Topics 发送消息的客户端。
- 可指定消息发送到哪个分区。(es.yml、jvm)
Consumers(消费者)
- 从 Topics 读取消息的客户端。
- 通过消费者组(Consumer Group)实现多个消费者并行处理。
三大内容
Topics(主题)
- 消息的分类单元,类似于数据库中的表。
- 每个 Topic 可分为多个分区(Partitions),支持并行处理。
Partitions(分区)
- 消息保存在 Topic 中,而为了能够实现大数据的存储,一个 topic 划分为多个分区,每个分区对应一个文件,可以分别存储到不同的机器上,以实现分布式的集群存储。
- 总结起来就是,一个 topic 对应的多个 partition 分散存储到集群中的多个 broker 上,存储方式是一个 partition 对应一个文件,每个 broker 负责存储在自己机器上的 partition 中的消息读写。
- 另外,每个 partition 可以有一定的副本,备份到多台机器上,以提高可用性。
为什么要分区呢? 为什么要进行分区呢?最根本的原因就是:kafka基于文件进行存储,当文件内容大到一定程度时,很容易达到单个磁盘的上限。
Replica(副本)
- kafka 还可以配置 partitions 需要备份的个数(replicas),每个 partition 将会被备份到多台机器上,以提高可用性,备份的数量可以通过配置文件指定。
- Kafka 定义了两类副本:领导者副本(Leader Replica) 和 追随者副本(Follower Replica);前者对外提供服务,后者只是被动跟随。
- 这种冗余备份的方式在分布式系统中是很常见的,那么既然有副本,就涉及到对同一个文件的多个备份如何进行管理和调度。kafka 采取的方案是:每个 partition 选举一个 server 作为“leader”,由 leader 负责所有对该分区的读写,其他 server 作为 follower 只需要简单的与 leader 同步,保持跟进即可。如果原来的 leader 失效,会重新选举由其他的 follower 来成为新的 leader。
- 至于如何选取 leader,实际上如果我们了解 ZooKeeper,就会发现其实这正是 Zookeeper 所擅长的,Kafka 使用 ZK 在 Broker 中选出一个 Controller,用于 Partition 分配和 Leader 选举。
- 另外,这里我们可以看到,实际上作为 leader 的 server 承担了该分区所有的读写请求,因此其压力是比较大的,从整体考虑,有多少个 partition 就意味着会有多少个leader,kafka 会将 leader 分散到不同的 broker 上,确保整体的负载均衡。

先有Topic,数据量越来越大,有了Partition的概念,集群规模越大,又有了Replica副本的概念。
消息传递机制
Kafka 支持 3 种消息投递语义。在业务中,常常都是使用 At least once 的模型。
- At most once:最多一次,消息可能会丢失,但不会重复。
- At least once:最少一次,消息不会丢失,可能会重复。
- Exactly once:只且一次,消息不丢失不重复,只且消费一次。
应用场景
Kafka 的使用场景
- 活动跟踪:Kafka 可以用来跟踪用户行为,比如我们经常回去淘宝购物,你打开淘宝的那一刻,你的登陆信息,登陆次数都会作为消息传输到 Kafka ,当你浏览购物的时候,你的浏览信息,你的搜索指数,你的购物爱好都会作为一个个消息传递给 Kafka ,这样就可以生成报告,可以做智能推荐,购买喜好等。
- 传递消息:Kafka 另外一个基本用途是传递消息,应用程序向用户发送通知就是通过传递消息来实现的,这些应用组件可以生成消息,而不需要关心消息的格式,也不需要关心消息是如何发送的。
- 度量指标:Kafka也经常用来记录运营监控数据。包括收集各种分布式应用的数据,生产各种操作的集中反馈,比如报警和报告。
- 日志记录:Kafka 的基本概念来源于提交日志,比如我们可以把数据库的更新发送到 Kafka 上,用来记录数据库的更新时间,通过kafka以统一接口服务的方式开放给各种consumer,例如hadoop、Hbase、Solr等。
- 流式处理:流式处理是有一个能够提供多种应用程序的领域。
- 限流削峰:Kafka 多用于互联网领域某一时刻请求特别多的情况下,可以把请求写入Kafka 中,避免直接请求后端程序导致服务崩溃。
Kafka 单机上手
环境准备
准备一台虚拟机,CentOS Stream 9
| 序号 | 主机IP | 主机名 | 备注 |
|---|---|---|---|
| 1 | 192.168.88.156 | kafka1.itcast.cn | 单机模式 |
#!/bin/bash
IP_ADDRESS="192.168.88.156" # 改成自己的
NEW_HOSTNAME="kafka1.itcast.cn" # 改成自己的
hostnamectl set-hostname $NEW_HOSTNAME
source ~/.bashrc
echo "$IP_ADDRESS $NEW_HOSTNAME" >> /etc/hosts
# 关闭防火墙(测试环境)
systemctl disable --now firewalld
systemctl stop firewalld
# 临时关闭
setenforce 0
# 永久关闭 SELinux 需要重启系统才能完全生效
sed -i 's/SELINUX=enforcing/SELINUX=disabled/g' /etc/selinux/config
# 配置时间同步
dnf install -y chrony
systemctl enable --now chronyd
chronyc makestep
# 安装必备软件
yum install -y vim wget rsync net-tools bash-completion命令作用:
- 把命令后的文本或变量展开结果写到标准输出;配合
>或>>时会改为写入文件,前者覆盖、后者追加。
- 使用 systemctl 对
--now服务执行disable操作:启动、停止、重启、查看状态或设置开机自启。
- 使用 systemctl 对
firewalld服务执行stop操作:启动、停止、重启、查看状态或设置开机自启。
- 这条命令用于批量替换文件内容。这里常见用途是把默认镜像地址、配置值或路径替换成当前环境可用的值。
- 使用
dnf包管理器执行install操作,对软件包进行安装、更新、删除或查询。
- 使用 systemctl 对
--now服务执行enable操作:启动、停止、重启、查看状态或设置开机自启。
- 使用
yum包管理器执行install操作,对软件包进行安装、更新、删除或查询。
执行结果:
- 成功后文本或变量展开结果会输出到终端;带重定向时会写入目标文件。
- 成功后会对
--now服务执行disable操作;restart会造成短暂中断,enable影响开机自启。
- 成功后会对
firewalld服务执行stop操作;restart会造成短暂中断,enable影响开机自启。
- 成功后应以该命令实际输出、生成的文件或目标服务状态为准;失败时先阅读终端报错,再核对命令参数和运行环境。
- 成功后软件包会按指定动作完成安装、更新、删除或查询;失败时检查仓库、网络和权限。
- 成功后会对
--now服务执行enable操作;restart会造成短暂中断,enable影响开机自启。
关键参数:
注意事项:
第一步,在其中一台主机上,执行
**注意事项:**
>
注意事项:
22 25
**注意事项:**
每台节点不同
注意事项:
所有节点依次执行
**注意事项:**
注意事项:
结果如下:
WATCHER::

注意事项:


单词拼接错误,应该是kafka。
**注意事项:**



分片有什么用呢?
**注意事项:**
hadoop就是这样的:
比如:有10台服务器,建议是3个副本,但是,能不能将数据压缩一下
hdfs存储数据的时候,会选择使用压缩来进行数据存储
当副本数是1的时候,就代表,没有副本注意事项:
本文由飞书云文档同步生成。涉及命令、SQL、配置示例时,请以飞书源文档和实际环境执行结果为准。