概述
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


常见问题
问题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