全部笔记All notes

Apache Kafka 分布式消息队列完整指南

阅读 9m 09s9m 09s read

概述

Apache Kafka 是一个高性能、可扩展的分布式消息队列系统,最初由 LinkedIn 开发。它能够处理每秒数百万条消息,是大数据生态系统中重要的组件。

核心特性

  • 高吞吐量:支持每秒处理数百万条消息
  • 可扩展性:支持动态扩容,无需停机
  • 持久性:消息持久化存储,支持数据回放
  • 容错性:分布式架构,自动故障转移
  • 实时性:低延迟消息传递

应用场景

场景描述示例
日志聚合收集分布式系统的日志ELK Stack + Kafka
实时数据流实时数据处理和分析Kafka + Flink/Spark
系统解耦微服务间异步通信订单系统 → 库存系统
事件驱动基于事件的架构用户行为分析

💡 提示: Kafka 特别适合需要高吞吐量、低延迟的大规模数据处理场景。

环境准备

系统要求

  • Docker 20.10+
  • 内存:至少 2GB RAM
  • 存储:SSD 推荐
  • Java 8+ (如果不使用 Docker)

架构组件

┌─────────────┐    ┌─────────────┐    ┌─────────────┐
│   Producer  │    │   Broker    │    │  Consumer   │
│   (生产者)   │───▶│  (代理)     │───▶│  (消费者)    │
└─────────────┘    └─────────────┘    └─────────────┘
                          │
                   ┌─────────────┐
                   │  ZooKeeper  │
                   │  (协调服务)  │
                   └─────────────┘

Docker 安装部署

1. 安装 ZooKeeper

ZooKeeper 用于管理 Kafka 集群的元数据:

# 启动 ZooKeeper 容器
docker run -d \
  --restart=always \
  --name zookeeper \
  -p 2181:2181 \
  -e ALLOW_ANONYMOUS_LOGIN=yes \
  -v /etc/localtime:/etc/localtime \
  --log-driver json-file \
  --log-opt max-size=100m \
  --log-opt max-file=2 \
  bitnami/zookeeper:latest

# 验证 ZooKeeper 状态
docker logs zookeeper

2. 安装 Kafka

# 启动 Kafka 容器(替换 <IP> 为实际 IP 地址)
docker run -d \
  --restart=always \
  --name kafka \
  -p 9092:9092 \
  -e KAFKA_BROKER_ID=0 \
  -e KAFKA_ZOOKEEPER_CONNECT=<IP>:2181/kafka \
  -e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://<IP>:9092 \
  -e KAFKA_LISTENERS=PLAINTEXT://0.0.0.0:9092 \
  -e KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1 \
  -v /etc/localtime:/etc/localtime \
  --log-driver json-file \
  --log-opt max-size=100m \
  --log-opt max-file=2 \
  bitnami/kafka:latest

# 验证 Kafka 状态
docker logs kafka

3. 使用 Docker Compose 部署

创建 docker-compose.yml 文件:

version: '3.8'

services:
  zookeeper:
    image: bitnami/zookeeper:latest
    container_name: zookeeper
    restart: unless-stopped
    ports:
      - "2181:2181"
    environment:
      - ALLOW_ANONYMOUS_LOGIN=yes
    volumes:
      - zookeeper_data:/bitnami/zookeeper
    logging:
      driver: json-file
      options:
        max-size: "100m"
        max-file: "2"

  kafka:
    image: bitnami/kafka:latest
    container_name: kafka
    restart: unless-stopped
    ports:
      - "9092:9092"
    environment:
      - KAFKA_BROKER_ID=0
      - KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181/kafka
      - KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9092
      - KAFKA_LISTENERS=PLAINTEXT://0.0.0.0:9092
      - KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1
      - KAFKA_AUTO_CREATE_TOPICS_ENABLE=true
    volumes:
      - kafka_data:/bitnami/kafka
    depends_on:
      - zookeeper
    logging:
      driver: json-file
      options:
        max-size: "100m"
        max-file: "2"

  kafka-manager:
    image: dushixiang/kafka-map:latest
    container_name: kafka-manager
    restart: unless-stopped
    ports:
      - "9001:8080"
    environment:
      - DEFAULT_USERNAME=admin
      - DEFAULT_PASSWORD=admin
    volumes:
      - kafka_manager_data:/usr/local/kafka-map/data
    depends_on:
      - kafka

volumes:
  zookeeper_data:
  kafka_data:
  kafka_manager_data:

启动服务:

docker-compose up -d

基本概念

核心术语

术语描述说明
Topic主题消息的逻辑分类,类似于数据库中的表
Partition分区Topic 的物理分片,提供并行处理能力
Producer生产者发送消息到 Kafka 的客户端
Consumer消费者从 Kafka 读取消息的客户端
Broker代理Kafka 集群中的服务器节点
Offset偏移量消息在分区中的唯一标识
Consumer Group消费者组多个消费者的逻辑组合

数据流架构

Producer → Topic(Partition 0) → Consumer Group A
              ├─ Partition 1  → Consumer Group B
              └─ Partition 2  → Consumer Group C

核心操作

1. 进入 Kafka 容器

# 进入 Kafka 容器
docker exec -it kafka /bin/bash

# 进入 Kafka 工具目录
cd /opt/bitnami/kafka/bin

# 查看 Kafka 版本
kafka-topics.sh --version

2. Topic 管理

创建 Topic

# 创建一个简单的 Topic
./kafka-topics.sh --create \
  --topic test-kafka \
  --bootstrap-server localhost:9092

# 创建带参数的 Topic
./kafka-topics.sh --create \
  --topic user-events \
  --bootstrap-server localhost:9092 \
  --partitions 3 \
  --replication-factor 1 \
  --config retention.ms=86400000

查看 Topic

# 列出所有 Topic
./kafka-topics.sh --list --bootstrap-server localhost:9092

# 查看 Topic 详细信息
./kafka-topics.sh --describe \
  --topic test-kafka \
  --bootstrap-server localhost:9092

# 查看所有 Topic 详细信息
./kafka-topics.sh --describe --bootstrap-server localhost:9092

修改 Topic

# 修改 Topic 分区数(只能增加)
./kafka-topics.sh --alter \
  --topic test-kafka \
  --partitions 5 \
  --bootstrap-server localhost:9092

# 修改 Topic 配置
./kafka-configs.sh --alter \
  --entity-type topics \
  --entity-name test-kafka \
  --add-config retention.ms=172800000 \
  --bootstrap-server localhost:9092

删除 Topic

# 删除 Topic
./kafka-topics.sh --delete \
  --topic test-kafka \
  --bootstrap-server localhost:9092

3. 消息生产和消费

生产者测试

# 启动控制台生产者
./kafka-console-producer.sh \
  --topic test-kafka \
  --bootstrap-server localhost:9092

# 带 key 的生产者
./kafka-console-producer.sh \
  --topic test-kafka \
  --bootstrap-server localhost:9092 \
  --property "parse.key=true" \
  --property "key.separator=:"

在生产者终端输入测试数据:

# 简单消息
Hello World!
{"id":1,"name":"张三","age":25}

# 带 key 的消息
user1:{"id":1,"name":"张三","age":25}
user2:{"id":2,"name":"李四","age":30}

消费者测试

# 启动控制台消费者(从最新消息开始)
./kafka-console-consumer.sh \
  --topic test-kafka \
  --bootstrap-server localhost:9092

# 从最早消息开始消费
./kafka-console-consumer.sh \
  --topic test-kafka \
  --from-beginning \
  --bootstrap-server localhost:9092

# 显示 key 和 value
./kafka-console-consumer.sh \
  --topic test-kafka \
  --from-beginning \
  --bootstrap-server localhost:9092 \
  --property print.key=true \
  --property print.value=true \
  --property key.separator=":"

4. 消费者组管理

# 查看消费者组列表
./kafka-consumer-groups.sh --list --bootstrap-server localhost:9092

# 查看消费者组详细信息
./kafka-consumer-groups.sh --describe \
  --group test-group \
  --bootstrap-server localhost:9092

# 重置消费者组偏移量
./kafka-consumer-groups.sh --reset-offsets \
  --group test-group \
  --topic test-kafka \
  --to-earliest \
  --bootstrap-server localhost:9092 \
  --execute

消费策略详解

1. 自动提交偏移量(Auto Offset Commit)

特点:

  • Kafka 默认的消费方式
  • 定期自动提交当前消费的偏移量
  • 简单易用,适合大多数场景

配置:

enable.auto.commit=true
auto.commit.interval.ms=5000

优势:

  • 实现简单,无需手动管理偏移量
  • 适合容错性要求不高的场景

劣势:

  • 可能导致消息重复消费或丢失
  • 消费者崩溃时的数据一致性问题

2. 手动提交偏移量(Manual Offset Commit)

特点:

  • 消费者手动控制偏移量提交时机
  • 确保数据完全处理后再提交
  • 提供更好的数据一致性保证

配置:

enable.auto.commit=false

实现方式:

// 同步提交
consumer.commitSync();

// 异步提交
consumer.commitAsync();

// 指定偏移量提交
Map<TopicPartition, OffsetAndMetadata> offsets = new HashMap<>();
offsets.put(new TopicPartition("topic", 0), new OffsetAndMetadata(100));
consumer.commitSync(offsets);

3. 批量消费(Batch Consumption)

特点:

  • 一次性处理多条消息
  • 提高处理效率和吞吐量
  • 减少网络开销

配置:

max.poll.records=500
fetch.min.bytes=1024
fetch.max.wait.ms=500

4. 消费指定偏移量(Consume from Specific Offset)

使用场景:

  • 数据回放和重新处理
  • 跳过特定的错误消息
  • 从特定位置开始消费

实现方式:

// 移动到指定偏移量
TopicPartition partition = new TopicPartition("topic", 0);
consumer.seek(partition, 100);

// 移动到分区开始位置
consumer.seekToBeginning(Arrays.asList(partition));

// 移动到分区结束位置
consumer.seekToEnd(Arrays.asList(partition));

5. 消费最新/最早消息

配置选项:

# 从最新消息开始消费
auto.offset.reset=latest

# 从最早消息开始消费
auto.offset.reset=earliest

# 如果没有偏移量信息则报错
auto.offset.reset=none

6. 基于时间戳的消费

实现方式:

// 根据时间戳查找偏移量
long timestamp = System.currentTimeMillis() - 24 * 60 * 60 * 1000; // 24小时前
Map<TopicPartition, Long> timestampsToSearch = new HashMap<>();
timestampsToSearch.put(new TopicPartition("topic", 0), timestamp);

Map<TopicPartition, OffsetAndMetadata> offsets = consumer.offsetsForTimes(timestampsToSearch);

编程接口

1. Python 客户端

生产者实现

import json
import time
from datetime import datetime
from kafka import KafkaProducer
from kafka.errors import KafkaError

class KafkaProducerClient:
    def __init__(self, bootstrap_servers, topic):
        self.topic = topic
        self.producer = KafkaProducer(
            bootstrap_servers=bootstrap_servers,
            key_serializer=lambda k: json.dumps(k).encode('utf-8'),
            value_serializer=lambda v: json.dumps(v).encode('utf-8'),
            # 性能配置
            batch_size=16384,
            linger_ms=10,
            buffer_memory=33554432,
            # 可靠性配置
            acks='all',
            retries=3,
            retry_backoff_ms=100
        )
    
    def send_message(self, key, value, partition=None):
        """发送消息"""
        try:
            future = self.producer.send(
                self.topic,
                key=key,
                value=value,
                partition=partition
            )
            
            # 等待发送结果
            record_metadata = future.get(timeout=10)
            print(f"Message sent successfully: {record_metadata}")
            return True
            
        except KafkaError as e:
            print(f"Failed to send message: {e}")
            return False
    
    def send_batch_messages(self, messages):
        """批量发送消息"""
        for message in messages:
            self.send_message(
                key=message.get('key'),
                value=message.get('value'),
                partition=message.get('partition')
            )
        
        # 确保所有消息都发送完成
        self.producer.flush()
    
    def close(self):
        """关闭生产者"""
        self.producer.close()

# 使用示例
if __name__ == "__main__":
    producer = KafkaProducerClient(
        bootstrap_servers=['localhost:9092'],
        topic='user-events'
    )
    
    # 发送单条消息
    producer.send_message(
        key='user1',
        value={
            'user_id': 1,
            'action': 'login',
            'timestamp': datetime.now().isoformat()
        }
    )
    
    # 发送批量消息
    messages = [
        {'key': 'user2', 'value': {'user_id': 2, 'action': 'view_page'}},
        {'key': 'user3', 'value': {'user_id': 3, 'action': 'purchase'}},
    ]
    producer.send_batch_messages(messages)
    
    producer.close()

消费者实现

import json
from kafka import KafkaConsumer
from kafka.errors import KafkaError

class KafkaConsumerClient:
    def __init__(self, bootstrap_servers, topics, group_id):
        self.consumer = KafkaConsumer(
            *topics,
            bootstrap_servers=bootstrap_servers,
            group_id=group_id,
            # 序列化配置
            key_deserializer=lambda k: json.loads(k.decode('utf-8')) if k else None,
            value_deserializer=lambda v: json.loads(v.decode('utf-8')),
            # 消费配置
            auto_offset_reset='earliest',
            enable_auto_commit=False,  # 手动提交
            max_poll_records=100,
            # 性能配置
            fetch_min_bytes=1024,
            fetch_max_wait_ms=500,
            # 心跳配置
            heartbeat_interval_ms=3000,
            session_timeout_ms=30000
        )
    
    def consume_messages(self, callback=None):
        """消费消息"""
        try:
            for message in self.consumer:
                try:
                    # 处理消息
                    result = self.process_message(message, callback)
                    
                    if result:
                        # 手动提交偏移量
                        self.consumer.commit()
                        print(f"Message processed successfully: {message.value}")
                    else:
                        print(f"Failed to process message: {message.value}")
                        
                except Exception as e:
                    print(f"Error processing message: {e}")
                    
        except KafkaError as e:
            print(f"Kafka error: {e}")
    
    def process_message(self, message, callback=None):
        """处理单条消息"""
        try:
            data = {
                'topic': message.topic,
                'partition': message.partition,
                'offset': message.offset,
                'key': message.key,
                'value': message.value,
                'timestamp': message.timestamp
            }
            
            if callback:
                return callback(data)
            else:
                print(f"Received message: {data}")
                return True
                
        except Exception as e:
            print(f"Error in process_message: {e}")
            return False
    
    def close(self):
        """关闭消费者"""
        self.consumer.close()

# 使用示例
def message_handler(data):
    """自定义消息处理函数"""
    print(f"Processing message from {data['topic']}: {data['value']}")
    
    # 根据消息类型进行不同处理
    if data['value'].get('action') == 'login':
        print("Processing login event")
    elif data['value'].get('action') == 'purchase':
        print("Processing purchase event")
    
    return True

if __name__ == "__main__":
    consumer = KafkaConsumerClient(
        bootstrap_servers=['localhost:9092'],
        topics=['user-events'],
        group_id='event-processor'
    )
    
    try:
        consumer.consume_messages(callback=message_handler)
    except KeyboardInterrupt:
        print("Stopping consumer...")
    finally:
        consumer.close()

2. Java 客户端

生产者实现

import org.apache.kafka.clients.producer.*;
import org.apache.kafka.common.serialization.StringSerializer;
import java.util.Properties;
import java.util.concurrent.Future;

public class KafkaProducerClient {
    private KafkaProducer<String, String> producer;
    private String topic;
    
    public KafkaProducerClient(String bootstrapServers, String topic) {
        this.topic = topic;
        
        Properties props = new Properties();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        
        // 性能配置
        props.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384);
        props.put(ProducerConfig.LINGER_MS_CONFIG, 10);
        props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 33554432);
        
        // 可靠性配置
        props.put(ProducerConfig.ACKS_CONFIG, "all");
        props.put(ProducerConfig.RETRIES_CONFIG, 3);
        props.put(ProducerConfig.RETRY_BACKOFF_MS_CONFIG, 100);
        
        this.producer = new KafkaProducer<>(props);
    }
    
    public void sendMessage(String key, String value) {
        ProducerRecord<String, String> record = new ProducerRecord<>(topic, key, value);
        
        Future<RecordMetadata> future = producer.send(record, new Callback() {
            @Override
            public void onCompletion(RecordMetadata metadata, Exception exception) {
                if (exception != null) {
                    System.err.println("Error sending message: " + exception.getMessage());
                } else {
                    System.out.println("Message sent successfully: " + metadata.toString());
                }
            }
        });
    }
    
    public void close() {
        producer.close();
    }
}

高级配置

1. 性能调优配置

生产者性能配置

# 批量发送配置
batch.size=16384
linger.ms=10
buffer.memory=33554432

# 压缩配置
compression.type=snappy

# 网络配置
send.buffer.bytes=131072
receive.buffer.bytes=32768

消费者性能配置

# 拉取配置
fetch.min.bytes=1024
fetch.max.wait.ms=500
max.poll.records=500

# 网络配置
receive.buffer.bytes=32768
send.buffer.bytes=131072

2. 可靠性配置

# 生产者可靠性
acks=all
retries=3
retry.backoff.ms=100
max.in.flight.requests.per.connection=1

# 消费者可靠性
enable.auto.commit=false
isolation.level=read_committed

3. 集群配置

# Broker 配置
num.network.threads=8
num.io.threads=8
socket.send.buffer.bytes=102400
socket.receive.buffer.bytes=102400
socket.request.max.bytes=104857600

# 日志配置
log.retention.hours=168
log.segment.bytes=1073741824
log.retention.check.interval.ms=300000

监控与管理

1. Kafka Manager 安装

# 启动 Kafka Manager
docker run -d \
  --name kafka-manager \
  -p 9001:8080 \
  -e DEFAULT_USERNAME=admin \
  -e DEFAULT_PASSWORD=admin \
  --restart always \
  dushixiang/kafka-map:latest

访问管理界面:http://localhost:9001

2. 监控指标

关键指标

指标类型指标名称描述
吞吐量MessagesInPerSec每秒接收消息数
吞吐量BytesInPerSec每秒接收字节数
延迟ProduceRequestLatency生产请求延迟
延迟FetchRequestLatency拉取请求延迟
错误率ErrorRate错误率

JMX 监控

# 启动 JMX 监控
export KAFKA_JMX_OPTS="-Dcom.sun.management.jmxremote=true -Dcom.sun.management.jmxremote.authenticate=false -Dcom.sun.management.jmxremote.ssl=false -Dcom.sun.management.jmxremote.port=9999"

# 查看 JMX 指标
kafka-run-class.sh kafka.tools.JmxTool \
  --jmx-url service:jmx:rmi:///jndi/rmi://localhost:9999/jmxrmi \
  --object-name kafka.server:type=BrokerTopicMetrics,name=MessagesInPerSec

Kafka生产者测试

Kafka消费者测试

常见问题

问题1:消息丢失

现象: 生产者发送的消息消费者接收不到 原因:

  • 生产者 acks 配置不当
  • 消费者偏移量提交问题
  • 网络故障导致消息丢失

解决方案:

# 生产者配置
acks=all
retries=3
enable.idempotence=true

# 消费者配置
enable.auto.commit=false
# 手动提交偏移量

问题2:消息重复消费

现象: 同一条消息被多次消费 原因:

  • 消费者崩溃后重新消费
  • 自动提交偏移量机制问题
  • 消费者组管理问题

解决方案:

// 实现幂等性处理
public void processMessage(String messageId, String content) {
    if (isAlreadyProcessed(messageId)) {
        return; // 已处理过,跳过
    }
    
    // 处理消息
    doProcess(content);
    
    // 记录处理状态
    markAsProcessed(messageId);
}

问题3:消费延迟

现象: 消费者处理消息速度跟不上生产速度 原因:

  • 消费者处理逻辑耗时
  • 分区数不够
  • 消费者数量不足

解决方案:

# 增加分区数
kafka-topics.sh --alter --topic my-topic --partitions 10 --bootstrap-server localhost:9092

# 增加消费者实例
# 启动多个消费者进程,使用相同的 group.id

问题4:磁盘空间不足

现象: Kafka 服务器磁盘空间耗尽 原因:

  • 日志保留时间过长
  • 日志文件过大
  • 监控不及时

解决方案:

# 调整日志保留策略
log.retention.hours=24
log.retention.bytes=1073741824
log.segment.bytes=536870912

# 启用日志压缩
log.cleanup.policy=compact

⚠️ 注意: 在生产环境中,建议使用专业的监控工具如 Prometheus + Grafana 来监控 Kafka 集群状态。

相关文章

大数据生态

数据库相关

其他中间件

运维部署

监控运维

相关阅读

  • Flink DataStream 高级特性:../../../大数据/Flink/DataStream API/Flink DataStream API 高级特性.md
  • Kubernetes 本地部署:../../../运维/Linux/容器编排/Kubernetes/本地部署K8s集群(一键部署).md
  • MySQL 基础:../../../数据库/MySQL/MySQL.md
  • Redis 入门:./Redis.md
  • RabbitMQ 入门:./RabbitMQ.md