大数据架构部署技术 下——搭建kafka(KRaft模式)集群

一、kafka简述

Kafka是由LinkedIn开发的一个分布式流处理平台,作为一种消息中间件,主要用于消费数据、解耦、冗余、异步通信、流量削峰、缓冲等

Kafka支持多种消息格式,如JSON、XML、Avro等,具有高吞吐量、可持久化、可水平扩展、高容错性、支持数据压缩、支持批量和流式处理等特性

在大数据架构中Zookeeper作为集群的分布式协调服务为kafka和Hadoop的高可用集群提供多Master节点选主、节点状态监控、动态配置变更、注册与发现等服务。所以在kafka2.8版本以前与Hadoop都依赖于Zookeeper,从2.8版本开始kafka引入KRaft模式一种自研的共识协议,作为预览/早期访问,但依然建议使用外部的Zookeeper集群,一直到3.3版本KRaft才被标记为生产可用,最终在3.5版本ZooKeeper模式被标记为deprecated而KRaft模式成为默认,现在kafka4.0是完全移除ZooKeeper代码的版本,4.x全线默认且仅支持KRaft模式

新组织模式下Hadoop生态仍然需要Zookeeper为其服务,但kafka已经通过KRaft实现元数据的自治,数据集成关系不变的情况下使 Kafka 运维更轻量、故障隔离更好、分区上限更高

kafka核心概念

Kafka主要由Producer、Broker、Consumer三部分组成

kafka角色组成

topic:是Kafka中消息的逻辑分类,类似于数据库中的"表"。生产者将消息发送到指定的topic,消费者从指定的topic读取消息。一个Kafka集群可以有多个topic,每个topic之间相互隔离

partition:Kafka中的消息分区,每个topic可以包含多个partition,每个partition可以存储多个消息,这是Kafka实现高吞吐和水平扩展的核心机制。消息被写入topic后会被分配到不同的partition中,每个partition是一个有序的、不可变的消息序列,消息在partition内部按写入顺序分配一个递增的偏移量(Offset)
分区的核心作用在于并行处理和水平扩展;多个partition可以分布在不同的broker上,消费者组内的多个消费者可以同时消费不同partition实现并行处理;当单个broker的读写能力不足时,增加partition数量并将它们分散到更多broker上即可提升吞吐量

partition是topic的物理实现单元,topic只是逻辑概念,真正的消息读写发生在partition上
默认参数中分区数量设置为3,分区数量决定了topic最大的消费并行度,一个消费者组内最多只能有与分区数相同数量的消费者同时工作,超出的消费者将处于空闲状态

replica(副本)与replication factor(副本因子):为了防止某个broker宕机导致数据丢失,Kafka为每个partition创建多个replica,分散存储在不同的broker上。默认每个partition有3个副本,

ISR(同步副本集):并非所有副本都有资格参与Leader选举。只有与Leader保持同步的副本才属于ISR。如果某个Follower的数据落后太多或网络异常导致心跳超时,它会被移出ISR

broker:Kafka集群中的节点会在Zookeeper注册并保持相关的元数据更新,负责存储消息、处理客户端请求、转发消息等

producer:生产者,负责将消息发送到Kafka集群

consumer:消费者,负责从Kafka集群中消费消息

consumer group(消费者组):消费者组是Kafka消费端的核心概念。多个消费者可以组成一个消费者组,共同消费一个topic的数据,实现负载均衡。一个partition在同一消费者组内只能被一个消费者消费,保证消息不会被重复处理,反之一个消费者可以消费多个partition,同时消费者组内消费者的数量不能超过partition数量,多余超出的消费者会处于空闲状态。不同消费者组之间互不影响,同一消息可以被不同消费者组各自消费一次。因此分区数决定了单个消费者组的最大消费并行度

Offset(偏移量):类似链表的指针,每个消息在partition内都有一个唯一且递增的编号,称为Offset(偏移量)。消费者通过Offset来记录自己消费到了哪个位置(committed offset,消费者提交的消费位点)。如果消费者宕机后重启,会从上次提交的 Offset 继续消费,避免消息丢失或重复,并且服务端不会记录消费者的消费进度,由消费者自行管理 Offset

KRaft新模式下的角色变动

在KRaft模式下角色发生变动,其中producer、consumer、topic、partition等角色与老模式保持一致

partition:Leader选举由controller节点负责,不再依赖ZooKeeper

broker:在KRaft模式下的注册由Zookeeper被替换为controller

controller:负责管理集群中的所有broker,包括broker的注册、注销、选举等操作,并通过KRaft协议实现集群的共识保证稳定性和一致性

controller 专注元数据,broker 专注消息读写

二、kafka集群节点环境部署

系统规划

主机名硬件性能组件服务
192.168.8.20 kafka-ansible1C2G 20Gansible
192.168.8.21 kafka-node14C4G 20GKafka
192.168.8.22 kafka-node24C4G 20GKafka
192.168.8.23 kafka-node34C4G 20GKafka

案例部署三节点的kafka集群,同样使用上篇中hadoop集群的node1-3节点就不更改主机名了,如果是自行生成新节点可以直接参照上篇中的环境准备工作,重新部署环境、注意主机名的变动,之后再下载kafka,也可以加入ansible方便后续操控

三、部署kafka

安装包版本选择

案例使用KRaft模式的kafka做部署,可以从官网下载kafka4.1.2版本(目前较为稳定的版本),案例部署通过wget获取安装包

https://kafka.apache.org/downloads

kafka配置说明

kafka集群的配置文件只需要修改server.properties,均在/usr/local/kafka/config目录下,并且官方有提供单节点的快速启动文档

https://kafka.apache.org/quickstart/

配置文件说明

config目录下有五个配置文件分别对应各个角色

server.properties:kafka 主配置文件,包含节点标识、角色、监听器、存储等核心配置
producer.properties:对应producer角色
consumer.properties:对应consumer角色
broker.properties:对应broker角色
controller.properties:对应controller角色

配置参数说明

node.id:节点id,相当于节点的"身份证号",集群内不论什么节点的id都不相同

process.roles:节点角色broker、controller,标识节点在集群中扮演的角色,多角色用逗号隔开。小集群各节点既管元数据又处理消息,大集群隔离角色

controller.quorum.voters:定义仲裁组(controller节点)的成员列表,即哪些节点参与Raft共识投票,并且所有节点(包括纯broker节点)都需要配置都需要知道谁能投票,格式为node.id@host:port,多个节点用逗号隔开

注意:controller.quorum.voters和controller.quorum.bootstrap.servers两个参数的区别在于前者解决谁来投票后者解决群众去哪里找"党",前者是必须配置的,后者是可选配置,只有纯brocker节点才需要知道去哪里找controller节点并且无法提供Raft选举所需的完整成员列表;在扩缩容时如果是纯broker节点靠后者就可以自动发现注册集群所以不需要重启集群,但如果是controller节点,所有broker节点成员都需要知道这个新增的节点必须重启集群服务

listeners:监听器,broker监听地址来自客户端的链接,controller监听Raft协议的链接;不绑定具体主机名直接监听所有地址防止hosts解析失败导致服务无法启动

advertised.listeners:节点对外暴露的广播地址

log.dirs:存储目录

num.partitions:设置partition分区数量

default.replication.factor:每个partition分区的副本数量

min.insync.replicas:最少同步副本数

default.replication.factor:副本总数

server.properties配置

修改配置文件/usr/local/kafka/config/server.properties,修改以下配置

# 节点id,集群中每个节点需要唯一,案例中的node1-3节点分别配置为0、1、2
node.id=0
# 关键节点角色,同一个节点既管元数据又处理消息
process.roles=broker,controller
# controller投票节点列表
controller.quorum.voters=0@node1:9093,1@node2:9093,2@node3:9093
# 可选,提供冗余发现能力
# controller.quorum.bootstrap.servers=node1:9093,node2:9093,node3:9093
# 监听器配置,broker监听地址来自客户端的链接,controller监听Raft协议的链接
listeners=PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093
# 对外暴露的广播地址,注意同步其他节点的主机名
advertised.listeners=PLAINTEXT://node1:9092
# 存储目录
log.dirs=/var/kafka/data

自动化部署kafka集群

更改配置文件

可以参考上述配置内容手动修改,案例复用Hadoop集群中的ansible环境编写playbook剧本,使用master节点操控node1-3节点修改配置文件

cd /ansible
# 更新inventory,新增kafka组
echo "
[master]
master ansible_host=192.168.8.20

[kafka]
node1 ansible_host=192.168.8.21
node2 ansible_host=192.168.8.22
node3 ansible_host=192.168.8.23" > /ansible/inventory
server.properties配置

准备server.properties配置文件的模板文件

# 创建server.properties.j2模板文件
cat >> /ansible/template/server.properties.j2 << "EOF"
##############################################
# Kafka KRaft 集群配置
# 节点: {{ inventory_hostname }} (node.id={{ kafka_node_ids[inventory_hostname] }})
##############################################

#------------- 节点标识与角色 ---------------
node.id={{ kafka_node_ids[inventory_hostname] }}
process.roles=broker,controller

#------------- Controller Quorum -----------
controller.quorum.voters=0@node1:9093,1@node2:9093,2@node3:9093
controller.listener.names=CONTROLLER

#------------- 监听器配置 ------------------
listeners=PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093
advertised.listeners=PLAINTEXT://{{ inventory_hostname }}:9092

#------------- 数据存储 --------------------
log.dirs={{ kafka_data_dir }}

#------------- Topic 默认值 ----------------
num.partitions=3
default.replication.factor=3
min.insync.replicas=2

#------------- 日志保留 --------------------
log.retention.hours=168
log.segment.bytes=1073741824
log.retention.check.interval.ms=300000

#------------- 网络 ------------------------
num.network.threads=3
num.io.threads=8
socket.send.buffer.bytes=102400
socket.receive.buffer.bytes=102400
socket.request.max.bytes=104857600

#------------- Group Coordinator ----------
group.initial.rebalance.delay.ms=3
EOF
service服务文件

准备service服务文件的模板文件

# 创建kafka.service.j2模板文件,用于systemd管理kafka服务
cat >> /ansible/template/kafka.service.j2 << "EOF"
[Unit]
Description=Apache Kafka Broker (KRaft)
Documentation=https://kafka.apache.org/documentation/
After=network-online.target
Wants=network-online.target

[Service]
Type=simple
User={{ kafka_user }}
Group={{ kafka_group }}

Environment="JAVA_HOME={{ java_home }}"
Environment="KAFKA_HOME={{ kafka_home }}"
Environment="KAFKA_LOG4J_OPTS=-Dlog4j.configuration=file:{{ kafka_home }}/config/log4j2.properties"
Environment="KAFKA_HEAP_OPTS=-Xms1g -Xmx1g"

ExecStart={{ kafka_home }}/bin/kafka-server-start.sh {{ kafka_home }}/config/server.properties
ExecStop={{ kafka_home }}/bin/kafka-server-stop.sh

Restart=on-failure
RestartSec=10
LimitNOFILE=65536

[Install]
WantedBy=multi-user.target
EOF
编写playbook剧本
# 编写set_kafka_conf.yml主剧本
cat >> /ansible/set_kafka_conf.yml << "EOF"
---
##############################################
# Kafka KRaft 集群自动化部署剧本
##############################################
- name: 部署 Kafka KRaft 集群
  hosts: kafka
  become: yes
  gather_facts: yes

  vars:
    #------------- 版本与路径 ---------------
    kafka_version: "4.1.2"
    scala_version: "2.13"
    kafka_install_dir: "/usr/local"
    kafka_home: "/usr/local/kafka"
    kafka_data_dir: "/var/kafka/data"
    kafka_user: "root"
    kafka_group: "root"
    java_home: "/usr/lib/jvm/java-17-openjdk"  # 此处按实际java版本更改

    #------------- 节点 ID 映射 ------------
    kafka_node_ids:
      node1: 0
      node2: 1
      node3: 2

  tasks:
    # ======================================
    # 环境检查
    # ======================================
    - name: "验证 Java 版本"
      shell: "java -version"
      register: java_check
      changed_when: false

    - name: "显示 Java 版本"
      debug:
        msg: "{{ inventory_hostname }}: {{ java_check.stderr_lines[0] }}"

    - name: "验证主机名解析"
      command: "getent hosts {{ item }}"
      loop: "{{ kafka_node_ids.keys() | list }}"
      changed_when: false

    - name: "检查端口占用"
      shell: "ss -tlnp | grep -E ':9092|:9093' || true"
      register: port_check
      changed_when: false

    - name: "端口冲突警告"
      debug:
        msg: "警告: {{ inventory_hostname }} 上 9092 或 9093 端口已被占用: {{ port_check.stdout }}"
      when: port_check.stdout | length > 0

    # ======================================
    # 安装 Kafka
    # ======================================
    - name: "检查 Kafka 是否已安装"
      stat:
        path: "{{ kafka_home }}/bin/kafka-server-start.sh"
      register: kafka_installed

    - name: "下载 Kafka 二进制包"
      get_url:
        url: "https://downloads.apache.org/kafka/{{ kafka_version }}/kafka_{{ scala_version }}-{{ kafka_version }}.tgz"
        dest: "/tmp/kafka_{{ scala_version }}-{{ kafka_version }}.tgz"
        mode: '0644'
        timeout: 300
      when: not kafka_installed.stat.exists
      tags: [install]

    - name: "解压 Kafka"
      unarchive:
        src: "/tmp/kafka_{{ scala_version }}-{{ kafka_version }}.tgz"
        dest: "{{ kafka_install_dir }}"
        remote_src: yes
      when: not kafka_installed.stat.exists
      tags: [install]

    - name: "创建软链接"
      file:
        src: "{{ kafka_install_dir }}/kafka_{{ scala_version }}-{{ kafka_version }}"
        dest: "{{ kafka_home }}"
        state: link
        force: yes
      tags: [install]

    - name: "清理安装包"
      file:
        path: "/tmp/kafka_{{ scala_version }}-{{ kafka_version }}.tgz"
        state: absent
      when: not kafka_installed.stat.exists

    # ======================================
    # 生成集群 ID(只在 node1 执行一次)
    # ======================================
    - name: "生成 Cluster UUID"
      shell: "{{ kafka_home }}/bin/kafka-storage.sh random-uuid"
      register: uuid_result
      delegate_to: node1
      run_once: true
      changed_when: false

    - name: "存储 Cluster UUID 到所有节点"
      set_fact:
        cluster_uuid: "{{ hostvars['node1']['uuid_result'].stdout }}"

    - name: "显示 Cluster UUID"
      debug:
        msg: "集群 UUID: {{ cluster_uuid }}"
      run_once: true

    # ======================================
    # 配置
    # ======================================
    - name: "创建数据目录"
      file:
        path: "{{ kafka_data_dir }}"
        state: directory
        owner: "{{ kafka_user }}"
        group: "{{ kafka_group }}"
        mode: '0755'

    - name: "部署 server.properties"
      template:
        src: "/ansible/template/server.properties.j2"
        dest: "{{ kafka_home }}/config/server.properties"
        owner: "{{ kafka_user }}"
        group: "{{ kafka_group }}"
        mode: '0644'
        backup: yes
      notify: restart kafka

    - name: "部署 systemd 服务文件"
      template:
        src: "/ansible/template/kafka.service.j2"
        dest: "/etc/systemd/system/kafka.service"
        mode: '0644'

    - name: "重载 systemd"
      systemd:
        daemon_reload: yes

    # ======================================
    # 格式化存储(仅首次)
    # ======================================
    - name: "检查是否已格式化"
      stat:
        path: "{{ kafka_data_dir }}/meta.properties"
      register: formatted

    - name: "执行存储格式化"
      shell: >
        {{ kafka_home }}/bin/kafka-storage.sh format
        -t {{ cluster_uuid }}
        -c {{ kafka_home }}/config/server.properties
      when: not formatted.stat.exists

    # ======================================
    # 启动服务
    # ======================================
    - name: "启动并设置开机自启"
      systemd:
        name: kafka
        state: started
        enabled: yes

    - name: "等待 Kafka 启动"
      wait_for:
        port: 9092
        host: "{{ inventory_hostname }}"
        timeout: 60
        state: started

    - name: "等待 Controller 就绪"
      wait_for:
        port: 9093
        host: "{{ inventory_hostname }}"
        timeout: 30
        state: started

  # ======================================
  # 通知处理器
  # ======================================
  handlers:
    - name: restart kafka
      systemd:
        name: kafka
        state: restarted
EOF

在这里插入图片描述

启动集群服务

# 执行playbook剧本
ansible-playbook set_kafka_conf.yml

# 验证kafka集群
ansible kafka -m shell -a "jps"

# 创建一个topic
/usr/local/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --create --partitions 1 --replication-factor 1 --topic msgtest

# 在节点上使用生产者
/usr/local/kafka/bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic msgtest

# 在另一节点上使用消费者
/usr/local/kafka/bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic msgtest

# 查看topic列表
/usr/local/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --list

# 查看topic详情(分区分布、Leader 所在节点、ISR 列表)
/usr/local/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 \
  --describe --topic msgtest

# 查看消费者组
/usr/local/kafka/bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --list

# 查看消费者组的消费进度
/usr/local/kafka/bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --describe --group <group-id>

在这里插入图片描述

执行测试数据消费

在这里插入图片描述

关闭集群服务

# 关闭服务并停止开机自启
ansible kafka -m shell -a "systemctl disable kafka --now"

四、手动部署(替代方案)

下载安装包

从官网下载kafka4.1.2版本tar包并解压到/usr/local目录下

# 解压到/usr/local目录下
tar -zxf kafka_2.13-4.1.2.tgz -C /usr/local/ --transform='s/kafka_2.13-4.1.2/kafka/'

更改配置文件

在node1-3节点各修改一遍配置文件

# 节点id,集群中每个节点需要唯一,案例中的node1-3节点分别配置为0、1、2
node.id=0
# 关键节点角色,同一个节点既管元数据又处理消息
process.roles=broker,controller
# controller投票节点列表
controller.quorum.voters=0@node1:9093,1@node2:9093,2@node3:9093
# 可选,提供冗余发现能力
# controller.quorum.bootstrap.servers=node1:9093,node2:9093,node3:9093
# 监听器配置,broker监听地址来自客户端的链接,controller监听Raft协议的链接
listeners=PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093
# 对外暴露的广播地址,注意同步其他节点的主机名
advertised.listeners=PLAINTEXT://node1:9092
# 存储目录
log.dirs=/var/kafka/data

生成集群UUID

首先生成集群的UUID,是KRaft模式引入的设计目的是标识集群的唯一ID保存到环境变量中,防止不同集群之间的broker节点互相干扰

# 生成集群UUID,进入kafka的根目录下执行
KAFKA_CLUSTER_ID="$(bin/kafka-storage.sh random-uuid)"
# 验证集群ID
echo $KAFKA_CLUSTER_ID

格式化节点存储目录

使用生成的UUID格式化存储目录

## 在node1-3所有节点都执行一遍
# 创建存储目录
mkdir -p /var/kafka/data

# 格式化节点存储目录
/usr/local/kafka/bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c /usr/local/kafka/config/server.properties

注意:官方文档中的命令包含standalone参数表示单机节点,集群中需要删除参数

启动集群服务

在所有节点执行命令

## 在node1-3所有节点都执行一遍
/usr/local/kafka/bin/kafka-server-start.sh -daemon /usr/local/kafka/config/server.properties

更多推荐