阶段四:消息中间件之Kafka

学习目标

Kafka 背景

Kafka是由LinkedIn开发并开源的分布式消息系统,因其分布式及高吞吐率而被广泛使用,现已与Cloudera Hadoop,Apache Storm,Apache Spark集成。

kafka名字灵感来自卡夫卡小说(《变形记》作者),寓意:数据像小说情节一样有序流动但错综复杂

image
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>

点对点消息系统

在点对点系统中,消息被保留在队列中。 一个或多个消费者可以消耗队列中的消息,但是特定消息只能由最多一个消费者消费。 一旦消费者读取队列中的消息,它就从该队列中消失。

image

发布-订阅消息系统

在发布 - 订阅系统中,消息被保留在主题中。 与点对点系统不同,消费者可以订阅一个或多个主题并使用该主题中的所有消息。

在发布 - 订阅系统中,消息生产者称为发布者,消息使用者称为订阅者。

一个现实生活的例子是Dish电视,它发布不同的渠道,如运动,电影,音乐等,任何人都可以订阅自己的频道集,并获得他们订阅的频道时可用。

image

Kafka 概述

Kafka 介绍

Kafka是一种分布式的,基于发布/订阅的消息系统。主要设计目标如下:

用人体来类比,消息队列(MQ)就是数字世界的血液输送网络​,而Kafka则是其中的"超级主动脉"。

<sheet sheet-id="NILF1U" token="VF7fscIFthTIBBtJYbCcExobnmb"></sheet>

Kafka的官网:http://kafka.apache.org

image

相关术语

image

为什么会有 Zookeeper Cluster 集群? 在传统 Kafka 架构中,ZooKeeper 是 “集群大脑”,负责元数据存储、控制器选举、Leader 选举、消费组协调等核心分布式协调工作。

三大主体

Brokers(kafka服务)

Producers(生产者)

Consumers(消费者)

三大内容

Topics(主题)

Partitions(分区)

为什么要分区呢? 为什么要进行分区呢?最根本的原因就是:kafka基于文件进行存储,当文件内容大到一定程度时,很容易达到单个磁盘的上限。

Replica(副本)

image

先有Topic,数据量越来越大,有了Partition的概念,集群规模越大,又有了Replica副本的概念。

消息传递机制

Kafka 支持 3 种消息投递语义。在业务中,常常都是使用 At least once 的模型。

应用场景

Kafka 的使用场景

Kafka 单机上手

环境准备

准备一台虚拟机,CentOS Stream 9

序号主机IP主机名备注
1192.168.88.156kafka1.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

命令作用:

执行结果:

关键参数:

注意事项:

第一步,在其中一台主机上,执行


**注意事项:**







> 

注意事项:

22 25


**注意事项:**



每台节点不同

注意事项:

所有节点依次执行


**注意事项:**


注意事项:

结果如下:

WATCHER::


![](/assets/feishu-images/7906d7e6efddcd1ee3437ae7.png)



注意事项:

image
image

单词拼接错误,应该是kafka。


**注意事项:**



image
image
image

分片有什么用呢?


**注意事项:**

hadoop就是这样的:
    比如:有10台服务器,建议是3个副本,但是,能不能将数据压缩一下
    hdfs存储数据的时候,会选择使用压缩来进行数据存储

    当副本数是1的时候,就代表,没有副本

注意事项: