Kafka简介和集群搭建

Kafka是LinkedIn公司于2010年开发的一款高性能消息中间件,目前是由Apache软件基金会维护的开源项目。Kafka主要采用Scala语言编写,它最初是为解决超大量数据的实时传输问题(尤其是日志收集)而设计的,Kafka通过分区日志结构和高效磁盘存储能实现每秒数百万级别的消息处理,而且它实现了完善的高可用分布式架构,支持分布式水平扩展,完美符合企业级的消息中间件要求,因此目前也是许多大型企业的首选实时数据管道建设方案。

项目主页:https://kafka.apache.org/

Kafka的特性和应用场景

Kafka作为一款高性能中间件,相比传统的JavaEE JMS实现(例如ActiveMQ)和以功能强大和高实时性著称的RabbitMQ,Kafka的设计目标其实与它们略有不同,Kafka也未遵循JMS标准或AMQP规范。

Kafka在最初其实就是为了日志采集而设计的,而日志采集的特点是高吞吐量、中间件高可用性,且允许消息少量丢失。Kafka的写入依赖顺序磁盘IO和批量发送,在普通硬件上单机即可达到百万级消息每秒的吞吐,消息不会在消费后被立即删除,而是按配置的保留策略保存在磁盘上,这意味着消费者可以重放历史消息,也意味着Kafka天然适合作为数据管道的中间层,而且通过增加Broker节点和Topic分区数,Kafka集群可以水平扩展,线性提升吞吐能力。此外Kafka目前仅实现了发布订阅模式,多个消费者可以独立消费同一个Topic的消息彼此互不影响,这和传统的JMS标准有些区别。

注:Kafka其实也能实现“at least once”和“exactly once”语义,只是在默认配置下更偏向高吞吐而非强一致性

总的来说,Kafka特别适合以下场景:

日志采集与实时计算:Kafka的高吞吐特性天然契合应用日志收集、埋点数据收集、用户行为事件收集等高频写入场景。

异步解耦:异步解耦也是消息中间件的主流用途,以订单处理为例,订单服务下单后,用户看到的界面上可无需等待短信、库存、积分等服务同步处理,服务端直接将消息投递到Kafka,各下游服务独立消费。

削峰填谷:在流量洪峰时,上游系统可将请求写入Kafka,下游系统按自身处理能力匀速消费,避免整个系统被突发流量压垮。

事件驱动架构:服务间通过事件通信,Kafka可作为事件总线,各服务订阅自己关心的事件并独立响应。

当然,Kafka也并不是所有场景的最优解,了解它擅长什么有助于你做出正确的技术选型。如果你的场景需要极低延迟(毫秒级RPC调用)或复杂的消息路由规则,Kafka在这方面较为薄弱,此时RabbitMQ可能是更合适的选择。

Kafka中的几个核心概念

在具体搭建部署Kafka之前,我们得先了解几个之后会反复出现的核心概念。

Broker 消息中间件服务节点:Kafka的服务节点被叫做Broker,一个Kafka集群由一个或多个Broker组成,每个Broker是负责存储和转发消息的Kafka服务节点。Broker原本其实是经纪人或中介的意思,在消息系统中,中间件的服务节点需要负责接收消息 -> 存储与缓冲 -> 分发消息,这与Broker很像,因此大家给服务节点起了这样一个形象的名字。

Topic 主题:Topic是Kafka中消息的分类单元,可以近似理解为数据库中的表,我们常说的“建一个Topic”就类似于在数据库中“建一张表”,但不同的是Kafka中处理的是消息,生产者向指定Topic发送消息,消费者订阅Topic接收消息。

Partition 分区:Kafka是一个分布式系统,分布式系统中常用分区的概念。Kafka中一个Topic可以拆分成多个Partition分散存储在不同Broker上,这是Kafka实现水平扩展的关键设计。举例来说,如果数据有20TB,Kafka集群有3个节点,然而每个Kafka节点的硬盘只有8TB,那么我们就需要创建至少3个分区让数据落在不同节点上,否则在物理上单节点不可能存下所有数据。

Producer 消息生产者:Producer负责生产消息,并将消息写入到某个Topic下。不过一般Topic的Partition都不是1Kafka客户端会负责根据消息Key自动计算写入消息到某个Partition以均衡负载,或者由程序员手动指定分区号。

Consumer Group 消息消费者组:和JMS不同,Kafka中没有单独的消费者,消费者是基于groupId定义成组的。消费者必须指定groupId参数,即使只有一个消费者,只要它有唯一的groupId,也会自动被识别为一个消费者组。Kafka的消息是发布/订阅模式的,一条消息会被广播给所有消费者组,即不同组之间则相互独立,各自都可获得完整的消息副本,但一个组内只有一个消费者会收到消息。

Consumer 消息消费者:Consumer负责从Kafka中读取消息,是数据的“下游处理方”。它会从指定的Topic的某些Partition中拉取数据进行消费处理。Kafka的消费模式是pull(主动拉取)模式,也就是说Consumer会不断向Broker发起请求获取数据,而不是Broker主动推送。在一个Consumer Group内部,每个Partition同一时间只能被一个Consumer消费,但一个Consumer可能被对接到了多个Partition上。

Kafka安装部署

单机环境搭建

如果你只是想本地启动Kafka用于开发集成或测试用途,建议直接使用Docker搭建,我们执行以下命令创建所需的网络、Kafka实例和Kafka UI。Kafka早期版本依赖Zookeeper集群进行分布式管理,但新版本可以不使用Zookeeper而使用内置的Kraft模式,不仅性能更优也能让我们省去了搭建维护Zookeeper的麻烦。

由于我们要搭建Kafka和Kafka UI,其中Kafka UI需要能访问到Kafka,因此首先创建一个容器网络,后面会将两个容器都挂载到这个网络上。

docker network create kafka-net

然后我们创建Kafka容器,其中9092是对外部客户端使用的端口,9093是仅供容器网络内部的Kafka UI(其实也是作为客户端)使用的端口,9094是Kraft集群协调内部通信使用的端口,虽然我们创建的是单节点Kafka,但配置中也需要声明集群通信端口,不过容器对外只暴露9092就够用了。对客户端的端口我们创建了2个,这主要是因为Docker容器网络的限制以及为了顾及Kafka UI使用造成的,如果你不使用Kafka UI,那么除了集群内部通信端口以外监听1个端口就行。此外我们这里由于是测试用途也没有挂载任何数据卷,如果你在生产环境使用Docker部署需要注意这一点。

docker run -d \
  --name kafka \
  --network kafka-net \
  -p 9092:9092 \
  -e KAFKA_NODE_ID=1 \
  -e KAFKA_PROCESS_ROLES=broker,controller \
  -e KAFKA_LISTENERS="PLAINTEXT://0.0.0.0:9092,PLAINTEXT_INTERNAL://0.0.0.0:9093,CONTROLLER://0.0.0.0:9094" \
  -e KAFKA_ADVERTISED_LISTENERS="PLAINTEXT://localhost:9092,PLAINTEXT_INTERNAL://kafka:9093" \
  -e KAFKA_LISTENER_SECURITY_PROTOCOL_MAP="PLAINTEXT:PLAINTEXT,PLAINTEXT_INTERNAL:PLAINTEXT,CONTROLLER:PLAINTEXT" \
  -e KAFKA_INTER_BROKER_LISTENER_NAME=PLAINTEXT_INTERNAL \
  -e KAFKA_CONTROLLER_LISTENER_NAMES=CONTROLLER \
  -e KAFKA_CONTROLLER_QUORUM_VOTERS="1@kafka:9094" \
  -e KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1 \
  -e KAFKA_AUTO_CREATE_TOPICS_ENABLE=true \
  apache/kafka:3.9.2

执行下面命令启动Kafka UI。

docker run -d \
  --name kafka-ui \
  --network kafka-net \
  -p 8080:8080 \
  -e KAFKA_CLUSTERS_0_NAME=local \
  -e KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS="kafka:9093" \
  provectuslabs/kafka-ui:latest

注意:Kafka本身没有官方的UI界面,Kafka UI是一个非常流行的第三方Kafka图形界面项目,它的界面设计非常精美,推荐使用。

Kafka UI启动后,使用浏览器访问http://localhost:8080/即可打开Kafka UI。

Kafka的安装包中提供了一系列的管理脚本。如果你想用Kafka提供的脚本来访问,可以使用对外的9092端口访问,下面测试命令输出了Kafka服务的API列表。

./bin/kafka-broker-api-versions.sh --bootstrap-server localhost:9092

注意:不建议在Windows下使用Kafka提供的CMD管理脚本,Kafka虽然提供了Windows版本的CMD脚本,但执行后会报错CMD命令超长,实际不可用。Git Bash能在Windows下勉强运行Bash脚本,但运行非常缓慢也不建议使用。在Windows下如果一定要用管理脚本,可以尝试使用WSL2。

集群环境搭建

如果你打算搭建一个更接近实际生产情况的Kafka实例来研究,可以按照以下指南搭建Kafka集群模式。我们这里准备了3台Linux主机,每台主机都需要安装JDK21,至于Kafka的安装包可以在官网下载。

192.168.1.164
192.168.1.165
192.168.1.166

我们这里仍使用Kraft模式搭建,将压缩包解压缩后,我们可以找到config/kraft/server.properties配置文件,基于Kraft搭建Kafka集群可以基于该配置文件进行修改。在3个节点中,我们分别将其复制到一个位置,对于这个配置文件我们习惯将其根据节点ID命名为server0.propertiesserver1.properties等以此类推。

cp ./config/kraft/server.properties ./config/server0.properties

下面是一个例子配置。

# 当前节点的角色,如下配置即使用Kraft模式
process.roles=broker,controller
# 当前节点ID,集群中需要唯一,通常从0开始顺序指定即可
node.id=0
# 集群节点配置
controller.quorum.voters=0@192.168.1.164:9093,1@192.168.1.165:9093,2@192.168.1.166:9093

# 监听地址和端口
listeners=PLAINTEXT://:9092,CONTROLLER://:9093
inter.broker.listener.name=PLAINTEXT
advertised.listeners=PLAINTEXT://192.168.1.164:9092
controller.listener.names=CONTROLLER
listener.security.protocol.map=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,SSL:SSL,SASL_PLAINTEXT:SASL_PLAINTEXT,SASL_SSL:SASL_SSL
num.network.threads=3
num.io.threads=8
socket.send.buffer.bytes=102400
socket.receive.buffer.bytes=102400
socket.request.max.bytes=104857600

# 数据存储目录
log.dirs=/home/ubuntu/data
num.partitions=1
num.recovery.threads.per.data.dir=1

offsets.topic.replication.factor=1
transaction.state.log.replication.factor=1
transaction.state.log.min.isr=1

log.retention.hours=168
log.segment.bytes=1073741824
log.retention.check.interval.ms=300000

group.initial.rebalance.delay.ms=0

其中重要的配置都已经标出,Kafka默认使用9092端口提供服务,9093端口用于集群Kraft协调通信。我们这里将如上配置文件复制3份,注意要对应修改node.id和监听地址IP,修改完成后分发到三台主机上即可。

修改好配置文件后,我们还需要生成数据存储目录。Kafka默认使用/tmp分区下目录存储数据,然而通常我们是需要将数据持久化的,使用临时目录肯定不行。我们在上面配置文件中将其修改为了/home/ubuntu/data,但光有这个配置还不够,还需要在其中重新创建必要的文件。

Kafka安装包中的bin目录包含了很多管理脚本,我们先执行以下命令生成随机UUID。

./bin/kafka-storage.sh random-uuid

然后用同一个UUID在3台主机上分别执行类似如下命令,该命令会根据我们之前编辑的配置文件初始化本地存储目录内的文件结构,注意修改其中的-c参数。

./bin/kafka-storage.sh format -t FvMlnvveS5iw8p4d1nyYtw -c ./config/server0.properties

最后我们启动所有Kafka节点,在3台主机上执行以下启动命令,注意修改其中的配置文件名。

./bin/kafka-server-start.sh -daemon ./config/server0.properties

执行后,我们的Kafka即启动成功。

创建Topic

Kafka UI界面上可以直接创建Topic,不过我们这里还是先采用执行命令行脚本的方式。其中kafka-topics.sh可用于创建Topic,下面是一个创建Topic的例子。

./bin/kafka-topics.sh --create --bootstrap-server localhost:9092 --topic test.topic --partitions 1 --replication-factor 1
  • --bootstrap-server:Kafka节点的主机和端口,如果是集群模式可以指定任意一个,也可以逗号分隔指定多个(便于某一个节点宕机会自动切换到其它节点)
  • --topic:具体创建的Topic名,推荐小写并统一配合一种分隔符的命名方式,例如test.topictest-topictest_topic,不可使用其它符号或中文
  • --partitions:Topic分区数
  • --replication-factor:Topic副本数

创建完成后,我们可以用如下命令查看。

./bin/kafka-topics.sh --describe --bootstrap-server localhost:9092 --topic test.topic

收发消息

bin目录下也自带了收发消息的客户端,我们可以用其进行测试,下面命令可以启动消息消费客户端。

./bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic test.topic --group group0
  • --group:该参数指定消费者组ID,可以随意起名,如果不指定默认会生成一个新的消费者组。Kafka内部会记录某个组当前消费到了哪条消息,如果生成新的消费者组,那么就只能从订阅后新产生的消息进行消费,而不会读取到积压的消息。

下面命令用于启动发送消息的交互式客户端。

./bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic test.topic

这样我们发送消息后,在接收端就可以接收到了。

Java中使用Kafka客户端

Java中使用Kafka需要引入其客户端库,Maven依赖如下。

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-clients</artifactId>
    <version>3.9.2</version>
</dependency>

除此之外,该库依赖SLF4J输出日志,因此我们的工程中还需要具备SLF4J以及具体的日志实现,我们这里使用Logback。

<dependency>
    <groupId>org.slf4j</groupId>
    <artifactId>slf4j-api</artifactId>
    <version>2.0.17</version>
</dependency>
<dependency>
    <groupId>ch.qos.logback</groupId>
    <artifactId>logback-classic</artifactId>
    <version>1.5.32</version>
</dependency>

下面例子代码创建了Kafka消息消费者。

package com.gacfox.demo;

import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.serialization.StringDeserializer;

import java.time.Duration;
import java.util.Collections;
import java.util.Properties;

@Slf4j
public class Main {
    public static void main(String[] args) {
        // 创建KafkaConsumer
        Properties props = new Properties();
        // 集群节点,我这里是单节点,如果有多个需要逗号分隔
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        // 序列化器
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        // 消费者组ID
        props.put(ConsumerConfig.GROUP_ID_CONFIG, "group0");
        try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props)) {
            // 订阅消息
            String topic = "test.topic";
            consumer.subscribe(Collections.singletonList(topic));
            // 接收达到或超过10条消息后退出
            int received = 0;
            do {
                ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
                for (ConsumerRecord<String, String> record : records) {
                    log.info(record.value());
                    received++;
                }
            } while (received < 10);
        }
    }
}

下面例子代码创建了Kafka消息生产者。

package com.gacfox.demo;

import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.serialization.StringSerializer;

import java.util.Properties;

@Slf4j
public class Main {
    public static void main(String[] args) {
        // 创建KafkaProducer
        Properties props = new Properties();
        // 集群节点,我这里是单节点,如果有多个需要逗号分隔
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        // 序列化器
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        KafkaProducer<String, String> producer = new KafkaProducer<>(props);

        // 发送1条消息
        String topic = "test.topic";
        String msg = "Hello, world!";
        ProducerRecord<String, String> record = new ProducerRecord<>(topic, msg);
        // 异步发送不等待
        producer.send(record);

        producer.close();
    }
}

这里注意KafkaProducer是线程安全的,我们可以在多线程中共享同一个KafkaProducer;但KafkaConsumer则不是线程安全的,多线程情况下我们需要创建多个KafkaConsumer

在Java中配置使用Kafka客户端库非常简单,但这不意味着使用Kafka本身简单,实际开发中我们要考虑的更多,包括消息同步发送、消息异步发送、发送回调、操作幂等性与事务等,这意味着我们需要对Kafka整体架构设计有一定了解后才能正确使用它,有关更深入的内容我们将在后续章节介绍。

作者:Gacfox
版权声明:本网站为非盈利性质,文章如非特殊说明均为原创,版权遵循知识共享协议CC BY-NC-ND 4.0进行授权,转载必须署名,禁止用于商业目的或演绎修改后转载。