Kafka是目前最流行的分布式消息队列之一广泛应用于大数据流处理日志收集等场景。

LinkedIn在2010年开发了Kafka后来开源贡献给了Apache基金会现在是,Apache的顶级项目。很多大公司,比如LinkedInNetflixUberTwitter等等都在大规模使用Kafka。

很多人想学习Kafka但是不知道从哪里开始觉得Kafka很复杂很难学。

今天想写一篇Kafka入门指南从零开始讲解Kafka的基本概念架构原理安装使用帮大家快速入门Kafka。

一、什么是Kafka

在讲具体内容之前,先理解什么是Kafka。

Kafka是一个分布式的流处理平台也是一个消息队列系统。它能发布和订阅消息流能存储消息流能处理消息流。

简单来说Kafka就是一个高吞吐量的消息队列能处理大量的消息支持分布式部署高可用高可靠。

Kafka的核心能力有三个:

  1. 发布和订阅:能发布消息到Topic也能订阅Topic的消息。
  2. 存储:能持久化存储消息支持消息回溯和重放。
  3. 流处理:能实时处理消息流支持复杂的流处理操作。

Kafka和传统的消息队列,比如RabbitMQActiveMQ有一些不同:

  • Kafka更注重吞吐量能处理大量的消息适合大数据场景。
  • Kafka消息是持久化的能保留一段时间支持消息回溯和重放。
  • Kafka是分布式的支持水平扩展高可用。
  • Kafka不仅是消息队列还是流处理平台支持实时流处理。

二、Kafka的核心概念

学习Kafka首先,要理解它的核心概念。

1. Broker(代理)

Broker是Kafka的服务器节点一个Kafka集群由多个Broker组成。每个Broker负责存储一部分消息处理客户端的请求。

Broker是无状态的通过ZooKeeper来管理集群状态和协调。

2. Topic(主题)

Topic是消息的分类,或者说消息的频道。Producer把消息发布到某个TopicConsumer订阅某个Topic来接收消息。

Topic是逻辑上的概念一个Topic可以分布在多个Broker上。

3. Partition(分区)

Partition是Topic的分区一个Topic可以分成多个Partition每个Partition是一个有序的消息队列。

Partition是Kafka并行和扩展的基础。通过把Topic分成多个Partition能提高吞吐量支持并行读写。

每个Partition在物理上对应一个文件夹里面存储消息数据文件和索引文件。

Partition中的消息是有序的每个消息有一个唯一的偏移量(Offset)用来标识消息在Partition中的位置。

4. Producer(生产者)

Producer是消息的生产者负责把消息发布到Topic。

Producer可以指定消息发送到哪个Partition也可以用默认的分区策略,比如轮询,或者根据消息的Key哈希分配Partition。

Producer发送消息是异步的支持批量发送提高吞吐量。

5. Consumer(消费者)

Consumer是消息的消费者负责订阅Topic拉取消息。

Consumer从Broker拉取消息而不是Broker推送消息给Consumer。这样Consumer能自己控制消费速度根据自己的处理能力拉取消息。

6. Consumer Group(消费者组)

Consumer Group是消费者的组一个Consumer Group包含多个Consumer它们共同消费一个Topic的消息。

同一个Consumer Group中的Consumer是竞争关系一个Partition的消息只会被同一个Consumer Group中的一个Consumer消费。这样能实现消息的负载均衡提高消费能力。

不同的Consumer Group之间,是独立的每个Consumer Group都能消费完整的Topic的消息互不影响。这样能实现消息的广播不同的业务系统可以用不同的Consumer Group消费同一个Topic的消息。

7. Offset(偏移量)

Offset是消息在Partition中的位置是一个递增的整数唯一标识一个Partition中的消息。

Consumer需要记录自己消费到哪个Offset了下次从这个Offset继续消费。Kafka默认把Offset存储在一个特殊的Topic叫__consumer_offsets中。

8. ZooKeeper

ZooKeeper是Kafka的协调服务用来管理集群元数据Broker注册Topic配置Partition分配Consumer Group管理等等。

Kafka依赖ZooKeeper来协调集群,所以部署Kafka之前,需要先部署ZooKeeper。

不过最新的Kafka版本正在逐步去掉对ZooKeeper的依赖用Kafka自己的协调机制代替叫KRaft模式,但是目前还是需要ZooKeeper。

三、Kafka的架构

理解了核心概念再看看Kafka的整体架构。

Kafka集群由多个Broker组成每个Broker是一个Kafka服务器节点。

Topic被分成多个Partition每个Partition有多个副本(Replica)分布在不同的Broker上保证高可用。

每个Partition的副本中有一个是Leader(领导者)其他是Follower(跟随者)。所有的读写请求都由Leader处理Follower从Leader同步数据。如果Leader挂了会从Follower中选举一个新的Leader保证服务不中断。

Producer把消息发送到Partition的LeaderLeader把消息写入本地日志Follower从Leader拉取消息写入自己的日志。

Consumer从Partition的Leader拉取消息消费。

ZooKeeper负责管理集群元数据Broker注册Topic配置Partition分配Leader选举Consumer Group管理等等。

整体架构是分布式的支持水平扩展高可用高吞吐量。

四、Kafka的安装

下面讲讲Kafka的安装和基本使用。

Kafka需要Java环境,所以先安装JDK推荐JDK 8或者JDK 11。

然后下载Kafka从Apache官网下载二进制包解压就能用。

步骤1:启动ZooKeeper

Kafka依赖ZooKeeper所以先启动ZooKeeper。Kafka的二进制包里自带了一个单节点的ZooKeeper配置可以直接用。

bin/zookeeper-server-start.sh config/zookeeper.properties

步骤2:启动Kafka Broker

然后启动KafkaBroker。

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

这样一个单节点的Kafka就启动了。生产环境需要部署多个Broker组成集群保证高可用。

五、Kafka的基本使用

启动Kafka之后,就可以用命令行工具来操作Kafka了。

1. 创建Topic

用kafka-topics.sh工具创建Topic。

bin/kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 1 --partitions 3 --topic test

这个命令创建了一个叫test的Topic有3个Partition1个副本。

2. 查看Topic列表

bin/kafka-topics.sh --list --zookeeper localhost:2181

3. 查看Topic详情

bin/kafka-topics.sh --describe --zookeeper localhost:2181 --topic test

4. 生产消息

用kafka-console-producer.sh工具生产消息。

bin/kafka-console-producer.sh --broker-list localhost:9092 --topic test

然后输入消息按回车发送。

5. 消费消息

用kafka-console-consumer.sh工具消费消息。

bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic test --from-beginning

--from-beginning表示从最早的消息开始消费。不加这个参数只消费启动之后,的新消息。

六、Java客户端使用

除了命令行工具Kafka还提供了各种语言的客户端,比如JavaPythonGo等等。下面简单讲讲Java客户端的使用。

1. 添加依赖

用Maven添加Kafka客户端依赖。

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

2. 生产者示例

import org.apache.kafka.clients.producer.*;
import java.util.Properties;

public class KafkaProducerExample {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put("bootstrap.servers", "localhost:9092");
        props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
        props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
        
        Producer<String, String> producer = new KafkaProducer<>(props);
        
        for (int i = 0; i < 10; i++) {
            ProducerRecord<String, String> record = 
                new ProducerRecord<>("test", "key" + i, "value" + i);
            producer.send(record, (metadata, exception) -> {
                if (exception != null) {
                    exception.printStackTrace();
                } else {
                    System.out.println("Sent to " + metadata.topic() + 
                        " partition " + metadata.partition() + 
                        " offset " + metadata.offset());
                }
            });
        }
        
        producer.close();
    }
}

3. 消费者示例

import org.apache.kafka.clients.consumer.*;
import java.util.Collections;
import java.util.Properties;

public class KafkaConsumerExample {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put("bootstrap.servers", "localhost:9092");
        props.put("group.id", "test-group");
        props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        
        Consumer<String, String> consumer = new KafkaConsumer<>(props);
        consumer.subscribe(Collections.singletonList("test"));
        
        while (true) {
            ConsumerRecords<String, String> records = consumer.poll(100);
            for (ConsumerRecord<String, String> record : records) {
                System.out.printf("offset = %d, key = %s, value = %s%n",
                    record.offset(), record.key(), record.value());
            }
        }
    }
}

七、Kafka的应用场景

Kafka的应用场景很广泛下面介绍几个常见的场景。

1. 消息队列

Kafka最基本的用途就是作为消息队列实现系统之间,的解耦异步削峰填谷。

比如订单系统下单后发送消息到Kafka库存系统物流系统通知系统等等订阅消息各自处理实现系统解耦。

2. 日志收集

Kafka很适合收集各种日志,比如应用日志访问日志审计日志等等。

各个应用服务器把日志发送到Kafka然后日志处理系统,比如ELK订阅Kafka的日志进行处理和分析。

Kafka的高吞吐量能处理大量的日志数据,而且支持多个消费者,同时消费同一份日志数据。

3. 流处理

Kafka不仅是消息队列还是流处理平台。Kafka Streams是Kafka自带的流处理库能实时处理Kafka中的消息流。

比如实时计算网站的访问统计实时监控系统指标实时风控等等都可以用Kafka Streams来实现。

Kafka也能和其他流处理框架配合,比如Spark StreamingFlinkStorm等等Kafka作为数据源流处理框架处理数据。

4. 事件溯源

事件溯源是一种架构模式把系统的所有状态变化都记录为事件存储起来需要的时候,通过重放事件来恢复状态。

Kafka的持久化存储和消息重放能力很适合事件溯源。把所有的事件发送到Kafka存储起来需要的时候,重放事件恢复状态。

5. 大数据管道

在大数据场景中Kafka经常作为数据管道连接各个大数据组件。比如数据从各个数据源采集到Kafka然后导入到HadoopSparkFlink等进行处理和分析。

Kafka的高吞吐量和分布式架构能处理海量的数据是大数据管道的理想选择。

八、学习Kafka的建议

最后给一些学习Kafka的建议。

1. 先理解核心概念

学习Kafka首先,要理解核心概念,比如TopicPartitionProducerConsumerConsumer GroupOffsetBrokerZooKeeper等等。这些概念是基础理解了这些才能进一步学习。

2. 动手实践

光看理论不够要动手实践。安装Kafka用命令行工具生产消费消息用Java客户端写生产者和,消费者体验一下Kafka的基本使用。

3. 理解架构和原理

掌握了基本使用之后,要深入理解Kafka的架构和原理,比如Partition的副本机制Leader选举消息存储机制消息可靠性保证高可用机制等等。

4. 学习高级特性

然后学习Kafka的高级特性,比如Kafka ConnectKafka StreamsExactly Once语义事务等等。

5. 了解生态系统

Kafka有丰富的生态系统,比如Kafka ConnectKafka StreamsSchema RegistryKSQL等等了解这些组件能更好地使用Kafka。

6. 看官方文档

Kafka的官方文档很详细也很权威学习Kafka要多看官方文档。

九、写在最后

以上就是Kafka消息队列的入门指南从零开始讲解了Kafka的基本概念架构原理安装使用以及,常见的应用场景。

Kafka是一个很强大的分布式消息队列和流处理平台,虽然看起来复杂,但是,只要理解了核心概念和架构原理就不难掌握。

现在Kafka越来越流行很多公司都在使用Kafka掌握Kafka对后端开发和,大数据开发来说都是很有价值的技能。

希望这篇入门指南能帮大家快速入门Kafka。如果想深入学习建议多看官方文档多动手实践在实际项目中使用Kafka。

如果有什么问题,或者不同的看法欢迎在评论区留言我们一起交流。

最后用一句话结束这篇文章:"Kafka是大数据时代的数据管道掌握它能帮你处理海量的数据。"

愿大家都能学好Kafka用好Kafka。