Article

消息队列 Kafka

更新于:2026-07-14

第一章:Kafka 基础概念与架构

1.1 什么是消息队列与 Kafka 的定位

概念名称说明注意事项
消息队列(Message Queue)一种应用程序间通信的中间件,通过发送和接收消息实现异步通信和解耦。消息队列通常用于解耦系统、削峰填谷、异步处理等场景。
同步通信调用方等待被调用方返回结果后才继续执行,如 HTTP 请求。容易造成服务阻塞,系统耦合度高。
异步通信发送方发送消息后无需等待接收方处理结果,提高系统响应速度和吞吐量。需要处理消息丢失、顺序、幂等性等问题。
Kafka 定位分布式流处理平台,具备高吞吐、低延迟、可扩展、持久化存储等特点,适用于大数据场景。Kafka 不仅是消息队列,还可用于日志聚合、流处理、事件溯源等。
主要应用场景日志收集、监控数据聚合、业务解耦、流式处理、消息广播等。根据场景选择合适的消息中间件,Kafka 适合高吞吐、大数据量场景,而非低延迟微服务通信。

1.2 Kafka 核心组件介绍(Broker、Topic、Partition、Producer、Consumer、Consumer Group)

组件名称说明注意事项
BrokerKafka 集群中的一个服务器节点,负责存储和转发消息。每个 Broker 有唯一 ID,多个 Broker 组成集群,提升可用性和吞吐量。
Topic消息的主题分类,生产者将消息发送到特定 Topic,消费者订阅 Topic 接收消息。Topic 是逻辑概念,可划分为多个 Partition 以实现并行处理。
PartitionTopic 的分区,每个 Partition 是一个有序、不可变的消息序列,分布在不同 Broker 上。分区是 Kafka 实现水平扩展和并发处理的基础,每个 Partition 只能由一个 Consumer 消费。
Producer生产者,向 Kafka 的 Topic 发送消息的应用程序。生产者可指定消息发送到哪个 Partition,或由分区策略自动选择。
Consumer消费者,从 Kafka 订阅 Topic 并拉取消息的应用程序。消费者通过拉取(pull)方式获取消息,控制消费速度。
Consumer Group消费者组,一组消费者的逻辑集合,共同消费一个或多个 Topic,实现负载均衡。同一 Consumer Group 内的消费者共享位移(offset),每个 Partition 只能被组内一个消费者消费。

1.3 Kafka 架构原理与数据流向

概念名称说明注意事项
数据流向Producer → Topic (Partition) → Consumer (Consumer Group)消息按 Partition 有序存储,消费者组内消费者并行消费不同 Partition。
Leader Partition每个 Partition 的主副本,负责处理所有读写请求。Leader 由 Controller Broker 选举产生,客户端直接与 Leader 交互。
Follower PartitionPartition 的副本,从 Leader 同步数据,用于容灾。Follower 不处理客户端请求,仅被动复制数据。
ISR(In-Sync Replicas)与 Leader 保持同步的副本集合,包含 Leader 和同步良好的 Follower。ISR 列表由 Broker 维护,若 Follower 落后过多将被移出 ISR。
Controller Broker集群中的控制器,负责 Partition Leader 选举、Broker 上下线处理等元数据管理。每个集群有且仅有一个 Controller,由 ZooKeeper 或 KRaft 协议选举产生。
Replication Factor每个 Partition 的副本数量,决定数据冗余度和可用性。建议设置为 3 以保证高可用,但会增加存储开销。

1.4 ZooKeeper 在 Kafka 中的作用(旧版)与 KRaft 模式(新版)

概念名称说明注意事项
ZooKeeper(旧版)Kafka 依赖 ZooKeeper 管理集群元数据,如 Broker 注册、Topic 配置、Leader 选举等。Kafka 3.0 前必须依赖 ZooKeeper,增加了系统复杂性和运维成本。
KRaft 模式Kafka Raft Metadata Mode,Kafka 自研的共识协议,替代 ZooKeeper 进行元数据管理。从 Kafka 2.8 开始支持,3.0+ 推荐使用,简化架构,提升性能和可维护性。
Controller 选举在 KRaft 模式下,Controller 由集群内部通过 Raft 协议选举产生,无需外部 ZooKeeper。KRaft 模式下集群可独立运行,减少外部依赖。
元数据存储KRaft 将元数据存储在内部 Topic(如 __cluster_metadata)中,由 Controller 管理。数据一致性由 Raft 协议保证,避免 ZooKeeper 成为瓶颈。
混合模式过渡阶段支持 ZooKeeper 和 KRaft 共存,便于升级。建议新集群直接使用 KRaft,老集群逐步迁移。
KRaft 节点角色支持 Controller 和 Broker 角色,通过配置决定节点类型。可部署专用 Controller 节点或混合部署,根据集群规模选择。

第二章:开发环境搭建与快速入门

2.1 搭建本地 Kafka 开发环境(单机/集群)

操作步骤说明注意事项
下载 Kafka 发行版从 Apache 官网下载 Kafka 二进制包,如 kafka_2.13-3.8.0.tgz选择与 Scala 版本兼容的包,开发环境建议使用最新稳定版。
启动 ZooKeeper(旧版)执行 bin/zookeeper-server-start.sh config/zookeeper.properties仅在使用 ZooKeeper 模式时需要,KRaft 模式可跳过。
启动 Kafka Broker执行 bin/kafka-server-start.sh config/server.properties确保配置文件中 broker.idlistenerslog.dirs 等正确设置。
配置 KRaft 模式(新版)使用 bin/kafka-storage.sh 格式化存储目录,并启动 server.properties 启用 KRaft。需设置 process.roles=broker,controllercontroller.quorum.voters 等参数。
单机环境验证使用命令行创建 Topic、发送和接收消息测试。确保防火墙开放 9092 端口,主机名解析正常。
伪集群搭建修改多个 server.properties 文件,配置不同 broker.id、端口和日志目录。用于本地测试多 Broker 场景,模拟集群行为。
环境变量与脚本封装设置 KAFKA_HOME,编写启动/停止脚本简化操作。提高开发效率,避免重复命令输入。

2.2 创建第一个 Kafka 生产者程序

方法名称语法用途注意事项
KafkaProducer(Properties)new KafkaProducer<>(props)构造生产者实例,传入配置属性。必须配置 bootstrap.serverskey.serializervalue.serializer
send(ProducerRecord)producer.send(record)发送消息到指定 Topic。方法返回 Future 对象,可调用 get() 实现同步发送。
send(ProducerRecord, Callback)producer.send(record, callback)异步发送消息,并指定回调函数处理结果。回调在生产者线程中执行,避免阻塞。
flush()producer.flush()强制将所有待发送消息立即发送出去。通常在批量发送后调用,确保消息发出。
close()producer.close()关闭生产者,释放资源。必须在程序结束前调用,否则可能导致消息丢失或资源泄漏。
close(Duration)producer.close(timeout)带超时时间的关闭方法,等待在途请求完成。避免无限等待,合理设置超时时间。

代码示例:构造生产者实例

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");
KafkaProducer<String, String> producer = new KafkaProducer<>(props);

代码示例:异步发送(带回调)

ProducerRecord<String, String> record = new ProducerRecord<>("test-topic", "key1", "value1");
producer.send(record, (metadata, exception) -> {
    if (exception != null) {
        System.err.println("发送失败: " + exception);
    } else {
        System.out.println("发送成功到: " + metadata.topic() + "-" + metadata.partition());
    }
});

2.3 创建第一个 Kafka 消费者程序

方法名称语法用途注意事项
KafkaConsumer(Properties)new KafkaConsumer<>(props)构造消费者实例,传入配置属性。必须配置 bootstrap.serversgroup.idkey.deserializervalue.deserializer
subscribe(Collection)consumer.subscribe(Arrays.asList("topic1", "topic2"))订阅一个或多个 Topic。订阅后会自动触发再平衡(Rebalance)。
subscribe(Pattern)consumer.subscribe(Pattern.compile("test.*"))订阅符合正则表达式的所有 Topic。适用于动态创建的 Topic。
poll(Duration)consumer.poll(timeout)从服务器拉取消息,返回 ConsumerRecords必须在循环中持续调用,timeout 设置拉取阻塞时间。
commitSync()consumer.commitSync()同步提交当前消费位移(offset),阻塞直到提交成功。确保消息处理成功后再提交,防止重复消费。
commitAsync()consumer.commitAsync()异步提交位移,不阻塞主线程。提交失败不会重试,可在关闭前调用 commitSync() 补救。
close()consumer.close()关闭消费者,释放资源。必须调用,通常在 finally 块中执行。

代码示例:构造消费者实例

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");
props.put("auto.commit.interval.ms", "1000");
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);

代码示例:消费消息并同步提交位移

consumer.subscribe(Arrays.asList("test-topic"));
try {
    while (true) {
        ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
        for (ConsumerRecord<String, String> record : records) {
            System.out.printf("topic=%s, partition=%d, offset=%d, key=%s, value=%s%n",
                    record.topic(), record.partition(), record.offset(), record.key(), record.value());
        }
        consumer.commitSync();
    }
} finally {
    consumer.close();
}

2.4 使用命令行工具测试消息收发

命令名称语法用途注意事项
kafka-topics.sh --createbin/kafka-topics.sh --create --topic <name> --partitions <n> --replication-factor <n> --bootstrap-server <host>创建 Topic单机测试 replication-factor 可设为 1,生产环境建议为 3。
kafka-topics.sh --listbin/kafka-topics.sh --list --bootstrap-server <host>列出所有 Topic验证 Topic 是否创建成功。
kafka-topics.sh --describebin/kafka-topics.sh --describe --topic <name> --bootstrap-server <host>查看 Topic 详细信息(分区、副本、Leader 等)用于排查分区分布和副本状态。
kafka-console-producer.shbin/kafka-console-producer.sh --bootstrap-server <host> --topic <name>启动控制台生产者,手动输入消息发送每行输入即发送一条消息,Ctrl+C 退出。
kafka-console-consumer.shbin/kafka-console-consumer.sh --bootstrap-server <host> --topic <name> --from-beginning启动控制台消费者,消费指定 Topic 消息--from-beginning 从头开始消费,否则从最新消息开始。
kafka-consumer-groups.sh --listbin/kafka-consumer-groups.sh --list --bootstrap-server <host>列出所有消费者组查看当前活跃的消费者组。
kafka-consumer-groups.sh --describebin/kafka-consumer-groups.sh --describe --group <name> --bootstrap-server <host>查看消费者组的消费进度(位移、LAG 等)LAG 表示未消费消息数,用于监控消费延迟。

命令行使用示例:

# 创建 Topic
bin/kafka-topics.sh --create --topic test-topic --partitions 3 --replication-factor 1 --bootstrap-server localhost:9092

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

# 查看 Topic 详情
bin/kafka-topics.sh --describe --topic test-topic --bootstrap-server localhost:9092

# 发送消息
bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic test-topic

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

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

# 查看消费进度
bin/kafka-consumer-groups.sh --describe --group test-group --bootstrap-server localhost:9092

第三章:Kafka 生产者(Producer)详解

3.1 生产者核心配置参数

配置参数语法示例用途注意事项
bootstrap.serversprops.put("bootstrap.servers", "host1:9092,host2:9092")指定 Kafka 集群的初始连接地址列表,生产者通过它发现集群元数据。至少配置两个 Broker 以防止单点故障,地址用逗号分隔。
key.serializerprops.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer")指定消息键(Key)的序列化类,必须实现 Serializer 接口。常见值:StringSerializerIntegerSerializerByteArraySerializer 等。
value.serializerprops.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer")指定消息值(Value)的序列化类。key.serializer 类似,必须正确配置,否则发送失败。
acksprops.put("acks", "all")控制消息写入副本的确认机制。0=不等待确认,1=Leader 确认,all=所有 ISR 副本确认。acks=all 提供最高持久性,但延迟较高;acks=1 平衡性能与可靠性。
retriesprops.put("retries", 3)设置生产者在遇到可重试异常(如网络抖动、Leader 选举)时的重试次数。设为 Integer.MAX_VALUE 可实现无限重试,需配合 retry.backoff.ms 使用。
batch.sizeprops.put("batch.size", 16384)每个批次(Batch)的字节数,生产者会累积消息直到达到此大小再发送。较大值可提高吞吐量,但增加延迟;默认 16KB。
linger.msprops.put("linger.ms", 5)批次等待更多消息的时间(毫秒),用于增加批次大小。设为 >0 可减少小消息发送次数,提升吞吐;但会增加延迟。
buffer.memoryprops.put("buffer.memory", 33554432)生产者缓冲区总大小,用于暂存待发送消息。默认 32MB,若消息产生速度过快可能导致缓冲区满,抛出 BufferExhaustedException
compression.typeprops.put("compression.type", "snappy")消息压缩类型,可选 nonegzipsnappylz4zstd压缩可减少网络传输和存储开销,但增加 CPU 消耗。
enable.idempotenceprops.put("enable.idempotence", true)启用幂等性生产者,确保消息在单分区不重复、不丢失、有序。需配合 acks=allretries=Integer.MAX_VALUEmax.in.flight.requests.per.connection=5(≤5)使用。
max.in.flight.requests.per.connectionprops.put("max.in.flight.requests.per.connection", 5)每个连接最多允许的未确认请求个数。幂等性模式下必须 ≤5,设为 1 可保证严格顺序,但降低吞吐。

3.2 发送消息的三种方式(同步、异步、带回调)

发送方式语法与代码示例用途注意事项
同步发送RecordMetadata metadata = producer.send(record).get();调用 send 后立即阻塞,等待 Broker 确认,返回消息元数据。简单直观,但吞吐低,适用于关键消息或调试。
异步发送(无回调)producer.send(record);发送后立即返回,不关心结果,吞吐最高。无法得知发送是否成功,可能丢失消息,仅用于测试或容忍丢失的场景。
异步发送(带回调)producer.send(record, new Callback() { ... });发送后立即返回,通过回调函数处理成功或失败。推荐方式,兼顾性能与错误处理,回调在生产者 I/O 线程中执行,避免阻塞。

代码示例:

// 同步发送
ProducerRecord<String, String> record = new ProducerRecord<>("topic", "key", "value");
RecordMetadata metadata = producer.send(record).get();
System.out.println("发送到: " + metadata.topic() + "-" + metadata.partition() + " offset=" + metadata.offset());

// 异步发送(无回调)
producer.send(record);

// 异步发送(带回调)
producer.send(record, new Callback() {
    @Override
    public void onCompletion(RecordMetadata metadata, Exception exception) {
        if (exception != null) {
            System.err.println("发送失败: " + exception.getMessage());
        } else {
            System.out.println("成功发送到 " + metadata.topic() + "-" + metadata.partition());
        }
    }
});

3.3 消息序列化与自定义 Serializer

方法/概念说明注意事项
Serializer 接口实现 org.apache.kafka.common.serialization.Serializer 接口,定义自定义序列化逻辑,将 Java 对象转为字节数组。必须实现 serialize 方法,注意处理 null 值。
配置自定义 Serializerprops.put("value.serializer", "com.example.CustomSerializer"),指定自定义序列化类的全限定名。类必须在 classpath 中,且无参构造函数可用。
close() 方法在生产者关闭时调用,用于清理资源(如关闭连接、释放缓冲区)。通常不需要显式资源管理,但复杂序列化器(如涉及连接池)应实现。
configure() 方法在序列化器初始化时调用,可获取生产者配置参数。isKey 为 true 表示用于序列化 Key,false 表示 Value。

代码示例:

public class CustomSerializer implements Serializer<MyObject> {
    @Override
    public byte[] serialize(String topic, MyObject data) {
        // 自定义序列化逻辑,如 JSON、Protobuf 等
        return JSON.toJSONString(data).getBytes(StandardCharsets.UTF_8);
    }

    @Override
    public void configure(Map<String, ?> configs, boolean isKey) {
        this.encoding = (String) configs.get("serializer.encoding");
    }

    @Override
    public void close() {
        // 如关闭 Protobuf 池
    }
}

3.4 分区策略与自定义 Partitioner

方法/概念说明用途注意事项
默认分区策略若消息有 Key,则对 Key 的 hash 取模;若无 Key,则轮询分配。实现负载均衡,相同 Key 的消息总在同一分区,保证顺序。Key 为 null 时,生产者使用粘性分区(sticky partitioning)优化批次发送。
自定义 Partitioner实现 org.apache.kafka.clients.producer.Partitioner 接口定义自己的分区分配逻辑,如按用户 ID、地理位置等路由。必须实现 partition 方法,返回 0 到 numPartitions-1 之间的整数。
配置自定义 Partitionerprops.put("partitioner.class", "com.example.CustomPartitioner")在生产者配置中指定自定义分区器类名。类必须在 classpath 中,且无参构造函数可用。
close()configure()分别用于初始化和关闭时执行逻辑。可在 configure 中读取配置参数,在 close 中释放资源。通常可留空,除非分区器持有状态或资源。

代码示例:

public class CustomPartitioner implements Partitioner {
    @Override
    public int partition(String topic, Object key, byte[] keyBytes,
                         Object value, byte[] valueBytes, Cluster cluster) {
        List<PartitionInfo> partitions = cluster.partitionsForTopic(topic);
        int numPartitions = partitions.size();
        // 自定义逻辑,如按用户 ID 哈希
        return (int) (Math.abs(((String) key).hashCode()) % numPartitions);
    }

    @Override
    public void configure(Map<String, ?> configs) { }

    @Override
    public void close() { }
}

3.5 消息确认机制(acks)与可靠性保证

配置参数说明用途注意事项
acks=0生产者不等待任何确认。消息可能丢失,吞吐最高,可靠性最低,适用于可容忍丢失的场景(如日志)。不推荐用于关键业务。
acks=1只需 Leader 副本写入成功即确认,Follower 可能未同步。平衡性能与可靠性,但若 Leader 故障且 Follower 未同步,消息会丢失。常用折中方案。
acks=all(或 -1需所有 ISR(In-Sync Replicas)副本都写入成功才确认。提供最高持久性,即使 Leader 故障,消息也不会丢失。需配合 min.insync.replicas 使用。
min.insync.replicasBroker 端参数,定义 ISR 最小副本数。当 ISR 数量 < 此值时,Producer 写入会失败。acks=all 配合使用,确保数据冗余。例如 replication.factor=3, min.insync.replicas=2,允许一个副本故障。Broker 配置:min.insync.replicas=2

3.6 消息重试机制与幂等性生产者

配置/概念语法示例用途注意事项
retriesprops.put("retries", Integer.MAX_VALUE)设置重试次数,应对临时性故障(如网络抖动、Leader 选举)。默认 0,不重试。设为 Integer.MAX_VALUE 可无限重试。
retry.backoff.msprops.put("retry.backoff.ms", 100)重试前等待的时间(毫秒),避免密集重试。默认 100ms,可根据网络状况调整。
enable.idempotenceprops.put("enable.idempotence", true)启用幂等性,确保每条消息在单个分区只被写入一次(不重复、不丢失、有序)。自动设置 retries=Integer.MAX_VALUE, acks=all, max.in.flight.requests.per.connection=5
幂等性原理生产者为每个消息分配 PID(Producer ID)和 Sequence NumberBroker 根据 PID 和 Sequence Number 去重,防止重复写入。仅保证单分区内的幂等性,跨分区或生产者重启后不保证。
事务性生产者(Transactional)props.put("transactional.id", "txn-01"); producer.initTransactions();提供跨分区、跨会话的原子性写入,确保一组消息全部成功或全部失败。需启用幂等性,使用 producer.beginTransaction(), send(), commitTransaction() 等方法。

3.7 生产者拦截器(Interceptor)

方法名称说明用途注意事项
onSend(ProducerRecord)ProducerRecord<K, V> onSend(ProducerRecord<K, V> record),在消息发送前拦截。可修改消息内容或记录日志。返回值为实际发送的消息,可包装或修改。返回值为实际发送的消息。
onAcknowledgement(RecordMetadata, Exception)void onAcknowledgement(RecordMetadata metadata, Exception exception),在收到 Broker 确认后拦截。可记录发送结果或监控。无论成功或失败都会调用。可用于监控和告警。
close()拦截器关闭时调用。释放资源,通常用于清理工作。如关闭监控客户端等。
configure(Map)拦截器初始化时调用。获取配置参数,可读取生产者配置中的自定义参数。this.metricPrefix = (String) configs.get("metric.prefix")
配置拦截器props.put("interceptor.classes", "com.example.ProducerInterceptor")指定拦截器类名,可配置多个,按顺序执行。拦截器链中任一环节抛出异常会中断流程,后续拦截器不执行。

代码示例:

public class LoggingInterceptor implements ProducerInterceptor<String, String> {
    @Override
    public ProducerRecord<String, String> onSend(ProducerRecord<String, String> record) {
        System.out.println("发送消息: " + record.value());
        return record; // 可返回修改后的 record
    }

    @Override
    public void onAcknowledgement(RecordMetadata metadata, Exception exception) {
        if (exception != null) {
            System.err.println("发送失败: " + exception);
        } else {
            System.out.println("发送成功: " + metadata.offset());
        }
    }

    @Override
    public void configure(Map<String, ?> configs) {
        this.metricPrefix = (String) configs.get("metric.prefix");
    }

    @Override
    public void close() {
        // 关闭监控客户端等
    }
}

// 配置多个拦截器
props.put("interceptor.classes",
        "com.example.LoggingInterceptor,com.example.MonitoringInterceptor");

第四章:Kafka 主题(Topic)与分区管理

4.1 Topic 的创建、查看与删除

操作方式语法/命令用途注意事项
命令行创建bin/kafka-topics.sh --create --topic <name> --partitions <n> --replication-factor <n> --bootstrap-server <host>创建指定分区数和副本因子的 Topic副本因子不能大于 Broker 数量,否则创建失败。
命令行查看列表bin/kafka-topics.sh --list --bootstrap-server <host>列出集群中所有 Topic可用于验证 Topic 是否创建成功。
命令行查看详情bin/kafka-topics.sh --describe --topic <name> --bootstrap-server <host>查看 Topic 的分区分布、Leader、副本、ISR 等信息输出包含 Partition、Leader、Replicas、Isr 字段,用于排查问题。
命令行删除bin/kafka-topics.sh --delete --topic <name> --bootstrap-server <host>删除指定 Topic(需 server.propertiesdelete.topic.enable=true删除后数据不会立即清除,需等待日志清理线程执行。
Java API 创建Admin.createTopics(new NewTopic(topic, partitions, replication))使用 Admin 客户端在代码中创建 Topic需引入 kafka-clients 依赖,适合自动化部署场景。
Java API 删除Admin.deleteTopics(Collection)使用 Admin 客户端删除 Topic异步操作,返回 DeleteTopicsResult 可监听结果。
Java API 查看Admin.describeTopics(Collection)获取 Topic 配置和元数据信息返回 TopicDescription 对象,包含分区和副本信息。

命令行示例:

# 创建 Topic
bin/kafka-topics.sh --create --topic user-logs --partitions 3 --replication-factor 2 --bootstrap-server localhost:9092

# 查看详情
bin/kafka-topics.sh --describe --topic user-logs --bootstrap-server localhost:9092

# 删除 Topic
bin/kafka-topics.sh --delete --topic user-logs --bootstrap-server localhost:9092

Java API 示例:

Admin admin = Admin.create(config);

// 创建 Topic
NewTopic newTopic = new NewTopic("orders", 6, (short) 3);
CreateTopicsResult result = admin.createTopics(Collections.singleton(newTopic));

// 查看 Topic
DescribeTopicsResult describeResult = admin.describeTopics(Arrays.asList("user-logs"));
KafkaFuture<Map<String, TopicDescription>> future = describeResult.values();

// 删除 Topic
DeleteTopicsResult deleteResult = admin.deleteTopics(Arrays.asList("temp-data"));

4.2 分区与副本机制(Leader、Follower)

概念名称说明注意事项
PartitionTopic 的分区,每个 Partition 是一个有序、不可变的消息序列,独立存储和处理。分区数决定并行度,增加分区可提升吞吐,但过多分区会增加管理开销。
Replication Factor每个 Partition 的副本数量,包括 1 个 Leader 和多个 Follower。副本用于容灾,建议生产环境设置为 3,确保高可用。
Leader每个 Partition 的主副本,负责处理所有读写请求,客户端直接与 Leader 交互。Leader 选举由 Controller 负责,若 Leader 宕机,Controller 会从 ISR 中选举新 Leader。
Follower副本的从节点,被动从 Leader 拉取消息,保持数据同步。Follower 不对外提供服务,仅用于数据冗余和故障转移。
数据分布分区均匀分布在多个 Broker 上,实现负载均衡。避免所有 Leader 集中在少数 Broker 上,影响性能。
写流程生产者发送消息到 Leader → Leader 写入本地日志 → Follower 从 Leader 拉取数据同步。消息写入 Leader 成功即返回(根据 acks 配置),Follower 异步复制。
读流程消费者从 Leader 拉取消息,Follower 不参与读取。保证数据一致性,避免从不同副本读取导致顺序错乱。

4.3 副本同步与 ISR(In-Sync Replicas)机制

概念名称说明注意事项
ISR(In-Sync Replicas)与 Leader 保持同步的副本集合,包括 Leader 自身和同步进度未落后的 Follower。ISR 是动态变化的,Broker 会定期检查 Follower 的同步延迟。
同步条件Follower 在 replica.lag.time.max.ms 时间内从 Leader 获取最新消息且未断开连接。默认值为 30 秒,超过此时间 Follower 将被移出 ISR。
OSR(Out-of-Sync Replicas)落后过多或断开连接的 Follower 副本,不在 ISR 中。OSR 中的副本无法参与 Leader 选举,存在数据丢失风险。
Leader 选举当 Leader 宕机时,Controller 从 ISR 中选举新 Leader,确保数据不丢失。仅从 ISR 中选举,避免选择落后副本导致数据不一致。
Unclean Leader Election若 ISR 为空,是否允许从 OSR 中选举 Leader(由 unclean.leader.election.enable 控制)。设为 true 可能导致数据丢失,生产环境建议设为 false。
高水位(HW)表示已提交(committed)的消息的偏移量,所有 ISR 副本都已写入该位置之前的消息。消费者只能读取到 HW 之前的消息,保证一致性。
LEO(Log End Offset)日志末端偏移量,表示下一条待写入消息的位置。Leader 和 Follower 都有自己的 LEO,用于计算同步进度。

4.4 分区再平衡与 Reassignment

操作方式语法/命令用途注意事项
自动再平衡消费者组内成员变化时自动触发消费者加入或退出时重新分配 Partition再平衡期间消费暂停,应尽量减少触发频率。
手动 Reassignmentbin/kafka-reassign-partitions.sh --generate --topics-to-move-json-file <file> --broker-list生成分区迁移计划用于 Broker 扩容或缩容时重新分布分区。
执行 Reassignmentbin/kafka-reassign-partitions.sh --execute --reassignment-json-file <file> --bootstrap-server执行生成的分区迁移计划迁移过程不影响服务,但会增加网络和磁盘 I/O。
验证 Reassignmentbin/kafka-reassign-partitions.sh --verify --reassignment-json-file <file> --bootstrap-server验证分区迁移是否完成迁移未完成前不要删除临时文件。
取消 Reassignmentbin/kafka-reassign-partitions.sh --cancel --bootstrap-server取消正在进行的分区迁移取消后可能留下部分迁移数据,需手动清理。
增加分区数bin/kafka-topics.sh --alter --topic <name> --partitions <n>增加 Topic 的分区数量(只能增加不能减少)分区数增加后,消费者组会触发再平衡,注意 key 到分区的映射可能变化。

命令行示例:

# 生成分区迁移计划
echo '{"topics":[{"topic":"user-logs"}],"version":1}' > topics.json
bin/kafka-reassign-partitions.sh --generate --topics-to-move-json-file topics.json --broker-list "0,1,2"

# 执行迁移
bin/kafka-reassign-partitions.sh --execute --reassignment-json-file reassignment.json --bootstrap-server localhost:9092

# 验证迁移
bin/kafka-reassign-partitions.sh --verify --reassignment-json-file reassignment.json --bootstrap-server localhost:9092

# 取消迁移
bin/kafka-reassign-partitions.sh --cancel --bootstrap-server localhost:9092

# 增加分区数
bin/kafka-topics.sh --alter --topic user-logs --partitions 6 --bootstrap-server localhost:9092

4.5 Topic 配置参数调优

参数名称默认值用途推荐设置/调优建议注意事项
retention.ms604800000(7 天)控制消息保留时间,超过此时间的消息将被清理。根据业务需求设置,如 1 天、3 天或更长。设置过短可能导致消费者来不及消费,过长则占用存储。
retention.bytes-1(不限制)控制单个 Partition 的最大日志大小。根据磁盘容量和数据量设置,如 1073741824(1GB)。retention.ms 共同作用,任一条件满足即触发清理。
segment.bytes1073741824(1GB)日志段文件大小,达到该值后滚动到新文件。可根据写入频率调整,频繁写入可设小些(如 512MB)。较小的段便于清理,但增加文件句柄开销。
segment.ms604800000(7 天)日志段最大存活时间,超过后强制滚动。可设置为 24 小时,确保每天一个段文件便于管理。segment.bytes 共同控制日志滚动。
cleanup.policydelete清理策略:delete(按时间/大小删除)或 compact(压缩,保留 key 最新值)。普通日志用 delete,需要保留状态的用 compactcompact 模式适用于 key-value 状态存储场景。
max.message.bytes1048588(1MB)单条消息最大大小。若需发送大消息,可调大(如 10MB),需同步调整生产者和消费者配置。过大消息影响性能和延迟,建议拆分或使用外部存储。
min.insync.replicas1ISR 中最小副本数,配合 acks=all 使用,确保消息写入多数副本。建议设为 2(副本数 ≥3 时),提高数据可靠性。若 ISR 数不足,生产者将收到 NotEnoughReplicasException
unclean.leader.election.enablefalse是否允许从非 ISR 副本中选举 Leader。生产环境建议 false,避免数据丢失;对可用性要求极高可设 true。设为 true 可能导致数据不一致。
compression.typeproducer消息压缩类型:nonegzipsnappylz4zstd建议 lz4zstd,压缩比高且性能好。压缩减少网络和磁盘开销,但增加 CPU 使用。

第五章:Kafka 消费者(Consumer)详解

5.1 消费者核心配置参数

配置参数名称语法示例用途注意事项
bootstrap.serversprops.put("bootstrap.servers", "host1:9092,host2:9092")指定 Kafka 集群的初始连接地址列表。客户端会自动发现所有 Broker,建议配置 2-3 个以提高可用性。
group.idprops.put("group.id", "consumer-group-1")指定消费者所属的消费者组 ID。同一组内的消费者共享 Topic 的消费负载。
key.deserializerprops.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer")指定 key 的反序列化类。必须实现 Deserializer 接口,常用有 String、Integer、ByteArray 等。
value.deserializerprops.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer")指定 value 的反序列化类。同上,需与生产者序列化方式匹配。
enable.auto.commitprops.put("enable.auto.commit", "true")是否开启自动提交位移。设为 true 时由 auto.commit.interval.ms 控制提交频率。
auto.commit.interval.msprops.put("auto.commit.interval.ms", "5000")自动提交位移的间隔时间(毫秒)。仅在 enable.auto.commit=true 时生效,过短增加 Broker 压力,过长增加重复消费风险。
auto.offset.resetprops.put("auto.offset.reset", "earliest")当位移不存在或无效时的重置策略。可选值:earliest(从头开始)、latest(从最新消息开始)、none(抛异常)。
session.timeout.msprops.put("session.timeout.ms", "45000")消费者与 Broker 的心跳超时时间。超时后 Broker 认为消费者失效并触发 Rebalance,建议设置为心跳间隔的 3 倍以上。
heartbeat.interval.msprops.put("heartbeat.interval.ms", "3000")消费者向 Broker 发送心跳的间隔。必须小于 session.timeout.ms,通常为其 1/3。
max.poll.recordsprops.put("max.poll.records", "500")每次 poll() 调用返回的最大记录数。控制单次处理消息量,避免处理超时导致 Rebalance。
max.poll.interval.msprops.put("max.poll.interval.ms", "300000")两次 poll() 调用之间的最大间隔。若处理消息时间过长超过此值,消费者会被踢出组。
fetch.min.bytesprops.put("fetch.min.bytes", "1")每次拉取请求返回的最小数据量。设为较大值可减少网络请求,但增加延迟。
fetch.max.wait.msprops.put("fetch.max.wait.ms", "500")拉取请求在 Broker 端等待数据积累的最大时间。配合 fetch.min.bytes 使用,实现批量拉取。

5.2 订阅主题与消息拉取机制

方法名称语法用途注意事项
subscribe(Collection)consumer.subscribe(Arrays.asList("t1"))订阅指定的一个或多个 Topic。订阅后自动参与再平衡,不可与 assign 混用。
subscribe(Pattern)consumer.subscribe(Pattern.compile("t.*"))订阅符合正则表达式的 Topic。动态匹配新创建的 Topic,适合监控类应用。
assign(Collection)consumer.assign(Arrays.asList(tp))手动为消费者分配特定的分区(TopicPartition)。跳过消费者组管理,适用于特定分区处理或外部位移管理。
poll(Duration)consumer.poll(Duration.ofMillis(100))拉取消息,阻塞指定时间或直到有数据可返回。必须在循环中调用,是消费者主循环的核心。
wakeup()consumer.wakeup()唤醒阻塞的 poll()commit 操作,用于优雅关闭。唯一可安全跨线程调用的方法,用于中断阻塞操作。
position(TopicPartition)consumer.position(tp)获取指定分区下一次拉取的起始位移。仅在分配了分区后调用有效。
endOffsets(Collection)consumer.endOffsets(partitions)获取指定分区当前的最后一条消息的位移(即日志末尾)。用于计算 LAG 或初始化位移。

代码示例:

// 订阅 Topic
consumer.subscribe(Arrays.asList("order-topic", "log-topic"));

// 正则订阅
consumer.subscribe(Pattern.compile("metrics-.*"));

// 手动分配分区
TopicPartition tp = new TopicPartition("test-topic", 0);
consumer.assign(Arrays.asList(tp));

// 获取位移
long nextOffset = consumer.position(new TopicPartition("t", 0));
Set<TopicPartition> partitions = consumer.assignment();
Map<TopicPartition, Long> ends = consumer.endOffsets(partitions);

5.3 反序列化与自定义 Deserializer

方法名称语法用途注意事项
configure(Map, boolean)public void configure(Map<String, ?> configs, boolean isKey)初始化 Deserializer,获取配置参数。isKey 为 true 表示反序列化 key,false 表示 value。
deserialize(String, byte[])public T deserialize(String topic, byte[] data)执行反序列化逻辑,将字节数组转换为 Java 对象。可根据 topic 或 data 内容进行差异化处理。
close()public void close()关闭 Deserializer,释放资源。通常用于关闭缓存或连接,但大多数 Deserializer 无状态,无需实现。

注: 自定义 Deserializer 需实现 org.apache.kafka.common.serialization.Deserializer<T> 接口。

代码示例:

public class CustomDeserializer implements Deserializer<String> {
    @Override
    public void configure(Map<String, ?> configs, boolean isKey) {
        this.isKey = isKey;
    }

    @Override
    public String deserialize(String topic, byte[] data) {
        return data == null ? null : new String(data, StandardCharsets.UTF_8);
    }

    @Override
    public void close() {
        // 清理资源
    }
}

5.4 消费者组(Consumer Group)与再平衡(Rebalance)

概念名称说明注意事项
消费者组(Consumer Group)一组消费者的逻辑集合,共同消费一个或多个 Topic,实现消息的广播或负载均衡。同一 Group 内每个 Partition 只能被一个 Consumer 消费,实现队列模式。
再平衡(Rebalance)当消费者组成员变化(加入/退出)或订阅 Topic 分区数变化时,Broker 重新分配 Partition 的过程。过程中所有消费者暂停消费,影响吞吐,应尽量减少不必要的 Rebalance。
Group Coordinator负责管理消费者组的 Broker 节点,处理成员加入、同步、位移提交等。每个 Group 有唯一的 Coordinator,由 Group ID 哈希决定。
JoinGroup Request消费者请求加入组的消息。首次订阅或 Rebalance 时发送。
SyncGroup Request在 Leader Consumer 选出后,向成员分发分区分配方案。分配方案由 Group Leader 制定(默认 StickyAssignor)。
Heartbeat Request消费者定期发送心跳以保持活跃状态。heartbeat.interval.ms 控制频率,超时则被踢出组。
Rebalance 避免策略减少处理时间、合理设置 max.poll.interval.ms、避免长时间 GC。处理消息逻辑不应阻塞 poll() 循环。

5.5 位移(Offset)管理:自动提交与手动提交

方法名称语法用途注意事项
commitSync()consumer.commitSync()同步提交当前所有分区的位移,阻塞直到提交成功。确保消息处理成功后再提交,防止数据丢失。
commitSync(Map)consumer.commitSync(offsets)同步提交指定分区的位移。精确控制每个分区的提交位移,常用于精确一次语义。
commitAsync()consumer.commitAsync()异步提交位移,不阻塞线程。性能更高,但失败不会重试。
commitAsync(Callback)consumer.commitAsync(callback)异步提交并指定回调处理结果。回调执行在消费者线程,避免复杂操作。
committed(TopicPartition)consumer.committed(tp)查询指定分区已提交的位移。返回 null 表示该分区无提交记录。
seek(TopicPartition, long)consumer.seek(tp, offset)手动将消费者位移定位到指定位置。用于重新处理消息或跳过数据,需在 poll() 前调用。
seekToBeginning(Collection)consumer.seekToBeginning(partitions)将指定分区的位移重置到起始位置。常用于测试或重新消费。
seekToEnd(Collection)consumer.seekToEnd(partitions)将指定分区的位移定位到末尾(最新消息之后)。实现从最新消息开始消费,等效于 auto.offset.reset=latest

代码示例:

// 同步提交
try {
    while (true) {
        ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
        for (ConsumerRecord<String, String> record : records) {
            // 处理消息
        }
        consumer.commitSync();
    }
} finally {
    consumer.close();
}

// 指定分区提交
Map<TopicPartition, OffsetAndMetadata> offsets = new HashMap<>();
offsets.put(tp, new OffsetAndMetadata(offset + 1));
consumer.commitSync(offsets);

// 异步提交
consumer.commitAsync((offsets, exception) -> {
    if (exception != null) {
        log.error("Commit failed", exception);
    }
});

// 查询已提交位移
OffsetAndMetadata committedOffset = consumer.committed(tp);

// 定位位移
consumer.seek(tp, 100L);
consumer.seekToBeginning(consumer.assignment());
consumer.seekToEnd(consumer.assignment());

5.6 再平衡监听器(RebalanceListener)

方法名称语法用途注意事项
onPartitionsRevoked(Collection)public void onPartitionsRevoked(Collection<TopicPartition> partitions)在分区被撤销前调用,用于提交位移或清理资源。必须在此方法中完成位移提交,否则可能丢失。
onPartitionsAssigned(Collection)public void onPartitionsAssigned(Collection<TopicPartition> partitions)在分区被分配后调用,可用于重置本地状态或 seek 位移。可在此方法中进行初始化操作。

使用方式: 在调用 subscribe 时传入监听器实例:

consumer.subscribe(topics, new ConsumerRebalanceListener() {
    @Override
    public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
        consumer.commitSync();
    }

    @Override
    public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
        for (TopicPartition tp : partitions) {
            consumer.seek(tp, 0);
        }
    }
});

5.7 消费者拦截器(Interceptor)

方法名称语法用途注意事项
onConsume(ConsumerRecord)public ConsumerRecord onConsume(ConsumerRecord record)在消息返回给用户前拦截并处理。可修改消息内容或头信息,或记录日志。
onCommit(Map)public void onCommit(Map<TopicPartition, OffsetAndMetadata> offsets)在位移提交完成后调用。仅通知作用,不能修改提交行为。
configure(Map)public void configure(Map<String, ?> configs)初始化拦截器。获取消费者配置参数。
close()public void close()关闭拦截器,释放资源。在消费者关闭时调用。

配置方式:

props.put("interceptor.classes", "com.example.MyConsumerInterceptor");

代码示例:

public class LoggingConsumerInterceptor implements ConsumerInterceptor<String, String> {
    @Override
    public ConsumerRecords<String, String> onConsume(ConsumerRecords<String, String> records) {
        for (ConsumerRecord<String, String> record : records) {
            System.out.println("消费: " + record.topic() + " " + record.value());
        }
        return records;
    }

    @Override
    public void onCommit(Map<TopicPartition, OffsetAndMetadata> offsets) {
        System.out.println("提交位移: " + offsets);
    }

    @Override
    public void configure(Map<String, ?> configs) {
        this.configs = configs;
    }

    @Override
    public void close() {
        // 释放资源
    }
}

第六章:Kafka 高级特性

6.1 Kafka Connect:数据导入导出

概念/组件名称说明注意事项
Kafka Connect分布式、可扩展的数据集成工具,用于在 Kafka 与其他系统间可靠地导入导出数据。支持批量和流式数据同步,无需编写代码。
Source Connector源连接器,从外部系统(如数据库、文件)读取数据并写入 Kafka Topic。常见实现:JDBC Source、File Source、Debezium(CDC)。
Sink Connector目标连接器,从 Kafka 读取数据并写入外部系统(如数据库、HDFS、Elasticsearch)。常见实现:JDBC Sink、HDFS Sink、Elasticsearch Sink。
Connect Worker运行 Connector 和 Task 的进程,支持独立模式(单节点)和分布式模式(多节点)。分布式模式支持自动故障转移和水平扩展。
TaskConnector 的实际执行单元,一个 Connector 可拆分为多个 Task 并行运行。Task 数量通常与 Topic 分区数匹配,提升吞吐量。
Converter负责在 Connect 内部数据格式与外部格式(如 JSON、Avro)之间转换。常用:JsonConverterAvroConverter(需 Schema Registry)。
Transformations在数据流动过程中进行轻量级转换,如字段过滤、重命名、路由等。使用 SMT(Single Message Transform)实现,避免额外处理服务。

配置示例(JDBC Source):

{
  "name": "jdbc-source",
  "config": {
    "connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
    "tasks.max": "1",
    "connection.url": "jdbc:mysql://localhost:3306/db",
    "table.whitelist": "users",
    "mode": "incrementing",
    "incrementing.column.name": "id",
    "topic.prefix": "db-"
  }
}

注意: 需确保数据库驱动在 CLASSPATH 中。

6.2 Kafka Streams:流处理入门

方法/组件名称语法/用途注意事项
StreamsBuilder构建流处理拓扑的构建器。用于定义 KStream 和 KTable。
KStream表示一个持续的数据流,支持 map、filter、flatMap 等操作。适用于事件流处理。
KTable表示一个 changelog 流,代表某个数据的最新状态。适用于聚合和状态查询。
stream(String topic)从指定 Topic 创建 KStream。消费 Topic 所有消息。
process(Processor)对每条记录执行自定义处理逻辑。可访问记录元数据和状态存储。
map(KeyValueMapper)转换每条记录的 key 和 value。无状态转换。
filter(Predicate)过滤满足条件的记录。保留返回 true 的记录。
groupByKey()按 key 分组,用于后续聚合。要求 key 已存在或使用 map 操作设置。
count()对分组数据进行计数。结果写入状态存储,支持容错。
to(String topic)将流写入指定 Topic。输出 Topic 需预先创建或自动创建。
KafkaStreams.start()启动流处理应用。应用为长期运行服务。
KafkaStreams.close()关闭流处理应用。释放资源,持久化状态。

代码示例:

StreamsBuilder builder = new StreamsBuilder();

// 从 Topic 创建流
KStream<String, String> source = builder.stream("orders");
KStream<String, String> stream = builder.stream("input-topic");

// 转换
stream.map((key, value) -> new KeyValue<>(key.toUpperCase(), value.length()));

// 过滤
stream.filter((key, value) -> value.contains("error"));

// 分组与聚合
KTable<String, Long> counts = stream.groupByKey().count(Materialized.as("counts-store"));

// 输出到 Topic
stream.to("output-topic");

// 启动
KafkaStreams streams = new KafkaStreams(builder.build(), config);
streams.start();

// 关闭
streams.close();

6.3 KSQL:流式 SQL 查询

命令/语句类型语法用途注意事项
CREATE STREAMCREATE STREAM stream_name (column_defs) WITH (kafka_topic='topic', value_format='FORMAT');创建流,映射到 Kafka Topic。value_format 支持 JSON、AVRO、DELIMITED 等。
CREATE TABLECREATE TABLE table_name (column_defs) WITH (kafka_topic='topic', value_format='FORMAT', key='key_column');创建表,表示实体的当前状态。需指定主键(key)。
SELECTSELECT * FROM stream_name WHERE condition;查询流或表中的数据。支持 WHERE、GROUP BY、JOIN 等标准 SQL 操作。
INSERT INTOINSERT INTO stream_name SELECT ...;将查询结果插入到另一个流。用于数据路由和转换。
CREATE STREAM AS SELECTCREATE STREAM output_stream AS SELECT ... FROM input_stream ...;创建持续查询,将结果写入新 Topic。查询持续运行,输出流自动创建。
JOINSELECT s.id, t.name FROM stream s LEFT JOIN table t ON s.id = t.id;流与表或流与流之间的连接操作。支持 INNER、LEFT JOIN,流-流 JOIN 需窗口定义。
WINDOWWINDOW TUMBLING (SIZE 5 MINUTES)定义时间窗口,用于聚合操作。支持 Tumbling、Hopping、Session 窗口。
TERMINATETERMINATE query_id;停止持续查询。停止后不再产生数据,但 Topic 保留。

示例:

-- 创建流
CREATE STREAM page_views (view_time BIGINT, user_id VARCHAR)
  WITH (kafka_topic='page_views', value_format='JSON');

-- 创建表
CREATE TABLE users (user_id VARCHAR PRIMARY KEY, email VARCHAR)
  WITH (kafka_topic='users', value_format='JSON', key='user_id');

-- 查询
SELECT user_id, view_time FROM page_views WHERE view_time > 1600000000;

-- 数据路由
INSERT INTO high_value_users SELECT * FROM users WHERE premium = true;

-- 持续查询
CREATE STREAM errors AS SELECT * FROM logs WHERE level = 'ERROR';

-- 窗口聚合
SELECT user_id, COUNT(*) FROM page_views
  WINDOW TUMBLING (SIZE 1 MINUTE) GROUP BY user_id;

-- 停止查询
TERMINATE CSAS_ERRORS_1;

6.4 消息压缩(Compression)机制

压缩类型特点适用场景注意事项
none无压缩,消息原文传输。网络和磁盘资源充足,或消息本身已压缩。默认值,性能最好,但占用资源最多。
gzip高压缩比,CPU 消耗较高。对存储和带宽要求高,可接受较高 CPU 开销。适合日志归档等场景。
snappy压缩比适中,速度快,由 Google 开发。平衡压缩比和性能,通用推荐。Kafka 内置支持,广泛使用。
lz4压缩速度极快,压缩比接近 snappy。高吞吐场景,追求低延迟。Kafka 默认压缩类型(从 0.10+)。
zstd高压缩比,可调压缩级别,性能优秀。需要高压缩比且 CPU 资源充足,现代推荐。Kafka 2.1+ 支持,压缩比优于 gzip,速度优于 snappy。
批处理压缩压缩发生在消息批次(batch)级别,而非单条消息。提升整体压缩效率。batch.sizelinger.ms 影响压缩效果。
消费者解压消费者自动解压,无需配置。透明处理,开发者无感知。Broker 存储压缩后的消息,消费者拉取后解压。

生产者配置:

props.put("compression.type", "lz4"); // 推荐
props.put("compression.type", "zstd"); // 高压缩比

6.5 消息头(Headers)的使用

方法名称语法用途注意事项
Headers.add(String, byte[])record.headers().add("trace-id", bytes)添加消息头键值对。值为 byte[],需自行序列化。
Headers.lastHeader(String)Header header = record.headers().lastHeader("trace-id");获取指定名称的最后一个头(允许重复键)。返回 Header 对象,value() 返回 byte[]
Headers.toArray()Header[] headers = record.headers().toArray();获取所有头的数组。可用于遍历所有头信息。
Headers.remove(String)record.headers().remove("temp-header");移除指定名称的头。移除后不再传递。
生产者设置头ProducerRecord 创建后添加。传递上下文信息,如认证、追踪、路由指令。头信息不参与分区计算。
消费者读取头ConsumerRecord 中通过 headers() 方法获取。在处理消息前读取元数据。头在消息传递中保持不变。
序列化注意事项头值需手动序列化为 byte[],如 String 需 getBytes,对象可 JSON 或二进制序列化。确保生产者和消费者使用相同编码。使用 StandardCharsets.UTF_8 确保一致性。

代码示例:

// 生产者添加头
ProducerRecord<String, String> record = new ProducerRecord<>("topic", "key", "value");
record.headers().add("trace-id", "abc123".getBytes(StandardCharsets.UTF_8));
record.headers().add("auth-token", tokenBytes);

// 消费者读取头
ConsumerRecord<String, String> cr = ...;
Header header = cr.headers().lastHeader("trace-id");
if (header != null) {
    String traceId = new String(header.value(), StandardCharsets.UTF_8);
}

// 遍历所有头
for (Header h : cr.headers()) {
    System.out.println(h.key() + ": " + Arrays.toString(h.value()));
}

// 移除头
record.headers().remove("internal-flag");

6.6 事务消息(Transactions)支持

方法名称语法用途注意事项
initTransactions()producer.initTransactions();初始化事务,必须在启用事务的生产者上调用。每个 transactional.id 在集群内唯一,用于恢复事务状态。
beginTransaction()producer.beginTransaction();开始新事务。必须在 initTransactions() 之后调用。
send(ProducerRecord)producer.send(record);在事务中发送消息,消息对消费者不可见。所有发送在事务提交前对消费者不可见。
commitTransaction()producer.commitTransaction();提交事务,消息变为可见。仅当所有消息成功写入且日志已同步。
abortTransaction()producer.abortTransaction();中止事务,丢弃所有已发送消息。发生异常时调用,确保原子性。
sendOffsetsToTransactionproducer.sendOffsetsToTransaction(offsets, consumerGroupId);在事务中提交消费者位移,实现精确一次(exactly-once)处理。需配合 Consumer 使用,确保消费-处理-生产的原子性。
EOS(Exactly-Once Semantics)Kafka 0.11+ 支持保证消息在流处理中处理一次,不丢失不重复。结合事务和幂等性实现,需设置 processing.guarantee="exactly_once_v2"(Streams)。

代码示例:

// 事务性生产者配置
Properties props = new Properties();
props.put("enable.idempotence", "true");
props.put("transactional.id", "tx-producer-1"); // 全局唯一
// ... 其他配置
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.initTransactions();

// 事务发送
try {
    producer.beginTransaction();
    producer.send(new ProducerRecord<>("t1", "k1", "v1"));
    producer.send(new ProducerRecord<>("t2", "k2", "v2"));
    producer.commitTransaction();
} catch (Exception e) {
    producer.abortTransaction();
}

// 精确一次处理(EOS)
Map<TopicPartition, OffsetAndMetadata> offsets = ...;
producer.sendOffsetsToTransaction(offsets, "group-1");

注意: 同一 transactional.id 只能有一个活跃生产者,否则会踢出旧实例。

第七章:性能调优与监控

7.1 生产者性能优化策略

优化策略配置参数/方法说明推荐值/建议注意事项
批量发送batch.size控制批次大小,达到后发送。16KB ~ 128KB,根据消息大小调整。过大增加延迟,过小降低吞吐。
延迟等待linger.ms批次未满时等待更多消息的时间。5 ~ 100ms,平衡延迟与吞吐。batch.size 协同作用,提升批次效率。
压缩compression.type启用压缩减少网络和磁盘开销。lz4zstd,兼顾压缩比与性能。压缩在 Producer 端进行,增加 CPU 使用。
异步发送enable.idempotence=true启用幂等性,确保单分区不重复。必须设置,配合 max.in.flight.requests.per.connection=15(Kafka 2.1+)。transactional.id 用于事务,保证原子性。
并行度多线程生产或增加分区提升整体吞吐量。分区数应大于等于生产者线程数。分区是并行度的基础。
序列化器优化自定义高效序列化(如 Avro、Protobuf)减少消息体积和序列化开销。避免使用 Java 原生序列化。结合 Schema Registry 管理模式。
ACK 机制acks=1acks=allacks=1:Leader 写入即返回,低延迟;acks=all:ISR 全同步,高可靠。高可靠性场景用 all,配合 min.insync.replicas>=2acks=0 不可靠,仅用于可丢失场景。
连接池与缓冲buffer.memory生产者缓存消息的总内存大小。根据吞吐需求调整,如 32MB ~ 128MB。溢出时 send() 会阻塞或抛出异常。

7.2 消费者性能优化策略

优化策略配置参数/方法说明推荐值/建议注意事项
批量拉取fetch.min.bytes每次拉取请求返回的最小数据量。1KB ~ 10KB,减少小批次拉取。过大增加延迟,需与 fetch.max.wait.ms 平衡。
拉取等待fetch.max.wait.ms若数据不足 fetch.min.bytes,Broker 等待时间。100 ~ 500ms。避免频繁空拉取。
拉取大小fetch.max.bytes单次拉取请求的最大数据量。50MB ~ 100MB(需 Broker message.max.bytes 支持)。防止单次拉取过大阻塞。
会话超时session.timeout.msConsumer 与 Group Coordinator 的心跳超时。10s ~ 30s。过短导致误判宕机,过长故障恢复慢。
心跳间隔heartbeat.interval.msConsumer 向 Coordinator 发送心跳的间隔。应小于 session.timeout.ms 的 1/3,如 3s。确保心跳及时,避免不必要的再平衡。
最大轮询间隔max.poll.interval.ms两次 poll() 调用的最大间隔,超过则触发再平衡。根据处理逻辑耗时设置,如 30s ~ 5min。长时间处理需拆分或启用手动提交。
位移提交enable.auto.commit=false禁用自动提交,手动控制位移提交时机。配合 commitSync()commitAsync() 实现精确控制。避免重复消费或数据丢失。
反序列化器优化高效反序列化(如 Avro、Protobuf)减少 CPU 开销。与生产者匹配。避免反序列化成为瓶颈。
多线程消费单 Consumer 多线程处理消息解耦拉取与处理,提升 CPU 利用率。注意位移提交的线程安全。避免在多线程中直接提交位移。

7.3 Broker 配置调优

参数名称默认值说明调优建议注意事项
num.network.threads3处理网络请求(接收)的线程数。6 ~ 12,根据网络流量调整。接收客户端请求。
num.io.threads8执行磁盘 I/O(读写日志)的线程数。16 ~ 32,通常为磁盘数的倍数。Kafka 依赖顺序 I/O,线程数影响吞吐。
socket.send.buffer.bytes100KBServer 端 SO_SNDBUF 缓冲区大小。1MB ~ 4MB,提升网络吞吐。与客户端 socket.request.max.bytes 匹配。
socket.receive.buffer.bytes100KBServer 端 SO_RCVBUF 缓冲区大小。1MB ~ 4MB。接收客户端数据。
queued.max.requests500网络线程可发送给 I/O 线程的最大请求数。1000 ~ 2000,缓冲更多请求。防止请求堆积。
log.flush.interval.messages-强制刷盘的消息条数间隔(不推荐使用)。通常依赖 log.flush.interval.ms 或操作系统刷盘。启用后影响性能,一般关闭。
log.flush.interval.ms-强制刷盘的时间间隔(不推荐)。通常不设置,依赖 log.flush.scheduler.interval.ms影响持久性和性能。
log.retention.hours168(7 天)消息保留时间。根据业务需求设置,如 24、72 小时。log.retention.bytes 共同作用。
log.segment.bytes1GB日志段文件大小。512MB ~ 2GB,根据清理频率调整。较小文件便于删除。
log.retention.check.interval.ms300000(5 分钟)检查是否需要清理日志的间隔。300000 ~ 3600000(1 小时)。频繁检查增加开销。
zookeeper.connection.timeout.ms6000Broker 连接 ZooKeeper 超时时间。60000(1 分钟),网络不稳定时增大。连接 ZooKeeper 失败会导致 Broker 不可用。
delete.topic.enabletrue是否允许通过 API 删除 Topic。生产环境建议 true,便于管理。删除后数据不会立即清除。

7.4 Kafka 监控指标与常用工具(JMX、Prometheus、Grafana)

监控维度关键指标(JMX MBean)说明工具集成
Brokerkafka.server:type=BrokerTopicMetrics,name=MessagesInPerSec每秒入站消息数,衡量生产负载。Prometheus 通过 JMX Exporter 抓取,Grafana 展示。
kafka.server:type=BrokerTopicMetrics,name=BytesInPerSec每秒入站字节数。
kafka.server:type=BrokerTopicMetrics,name=BytesOutPerSec每秒出站字节数,衡量消费负载。
kafka.server:type=ReplicaManager,name=UnderReplicatedPartitions非同步副本分区数,应为 0。告警关键指标。
kafka.server:type=KafkaRequestHandlerPool,name=RequestHandlerAvgIdlePercent请求处理线程空闲率,过低表示处理不过来。
Producerkafka.producer:type=producer-metrics,client-id="..."包括 record-send-raterequest-latency-avgoutgoing-byte-rate 等。客户端监控,定位生产瓶颈。
Consumerkafka.consumer:type=consumer-fetch-manager-metrics,client-id="..."包括 records-lag-max(最大滞后)、fetch-ratefetch-latency-avg 等。records-lag-max 是核心指标,监控消费延迟。
Topic/Partitionkafka.log:type=Log,name=LogEndOffset,topic=...,partition=...分区 LEO,结合消费者位移计算滞后量。需结合消费者组位移计算 consumer_lag
ZooKeeperkafka.server:type=SessionExpireListener,name=ZooKeeperAuthFailuresZooKeeper 认证失败次数。监控 ZK 健康状态。
kafka.server:type=SessionExpireListener,name=ZooKeeperDisconnectsZooKeeper 断开连接次数。

工具链:

工具说明
JMX Exporter + Prometheus + Grafana标准监控栈。部署 JMX Exporter 暴露指标,Prometheus 抓取,Grafana 可视化。
Kafka Manager / Cruise Control / Confluent Control Center专用管理工具,提供集群视图、再平衡、监控等。Cruise Control 支持自动负载均衡。

告警设置:

  • UnderReplicatedPartitions > 0:及时发现集群异常
  • ISR 变化频繁
  • Consumer Lag > 阈值
  • Request Latency 骤升

使用 Prometheus Alertmanager 或工具内置告警。

7.5 日志清理策略(Log Retention)

策略类型相关配置参数说明触发条件注意事项
基于时间删除log.retention.hours(或 minutes/days)、log.roll.hours按消息时间戳删除过期数据。消息时间戳超过 retention.ms时间基于消息的 timestamp,需确保生产者时间准确。
基于大小删除log.retention.byteslog.segment.bytes按分区日志总大小删除最旧段。分区日志总大小超过 retention.byteslog.retention.hours 任一满足即触发清理。
清理检查周期log.retention.check.interval.ms(默认 300000ms=5 分钟)清理服务检查日志是否可删除的频率。check.interval.ms 执行一次清理检查。频繁检查增加 CPU 开销。
日志段滚动log.segment.byteslog.segment.ms控制何时创建新日志段文件。当前段大小超过 segment.bytes 或存活时间超过 segment.ms滚动是删除的前提,小段文件便于精确清理。
删除流程-异步删除,不影响服务。确定可删除段 → 标记为 .deleted → 后台线程删除文件。删除不可逆。
压缩(Compact)cleanup.policy=compactmin.compaction.lag.msmax.compaction.lag.ms保留每个 key 的最新值,适用于状态存储。满足 min.compaction.lag.ms 后,后台线程合并日志,保留最新 key。不删除消息,仅合并;deletecompact 可同时启用。
混合策略cleanup.policy=delete,compact先按时间/大小删除过期数据,再对剩余数据压缩。结合两种策略优点。适用于需要长期保留但又需节省空间的场景。
特殊保留retention.ms=-1永久保留,永不删除。用于审计日志等。需手动管理磁盘空间。

第八章:Java Kafka 实战案例

8.1 日志收集系统设计与实现

场景描述: 构建一个分布式日志收集系统,将多台服务器的应用日志(如 Nginx、应用日志)统一收集到 Kafka,供后续分析、监控或存储。

架构设计:

[应用服务器] → (Filebeat/Fluentd) → [Kafka Cluster] → (Kafka Consumer) → [Elasticsearch/Logstash/S3]

日志采集端:

  • 使用 Filebeat 或 Fluentd 作为轻量级日志采集器,监控日志文件变化。
  • 配置 Filebeat 输出到 Kafka:
output.kafka:
  hosts: ["kafka-broker1:9092", "kafka-broker2:9092"]
  topic: app-logs-topic
  partition.round_robin:
    reachable_only: false
  required_acks: 1
  compression: gzip
  max_message_bytes: 1000000

Kafka Producer(Java,可选): 若需自定义处理(如过滤、格式化),可用 Java 编写 Producer。

Properties props = new Properties();
props.put("bootstrap.servers", "kafka1:9092,kafka2:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("acks", "1"); // 平衡吞吐与可靠性
props.put("retries", 3);
props.put("batch.size", 16384); // 16KB
props.put("linger.ms", 10); // 等待 10ms 凑批
props.put("compression.type", "gzip"); // 减少网络传输
props.put("buffer.memory", 33554432); // 32MB

Topic 设计:

  • 分区数:根据日志源数量和吞吐量预估,如 10-50 个。
  • 副本数:replication.factor=3,保证高可用。
  • 保留策略:retention.ms=604800000(7 天),或按大小 retention.bytes

Consumer(Java):

  • 消费日志写入 Elasticsearch 或 HDFS。
  • 幂等写入:确保日志不重复(如使用 Logstash 或自定义去重逻辑)。
  • 批量处理:提升写入下游存储的效率。

8.2 订单系统解耦:异步处理订单

场景描述: 电商系统中,用户下单后需执行库存扣减、积分计算、短信通知、物流准备等多个操作。使用 Kafka 解耦,提升系统响应速度和可靠性。

架构设计:

[订单服务] → (Kafka Producer) → [order-topic] → [库存服务] (Consumer)
                                    ├──→ [积分服务] (Consumer)
                                    ├──→ [通知服务] (Consumer)
                                    └──→ [物流服务] (Consumer)

订单服务(Producer):

@Service
public class OrderService {
    @Autowired
    private KafkaTemplate<String, String> kafkaTemplate;

    public void createOrder(Order order) {
        // 1. 保存订单到数据库(本地事务)
        orderRepository.save(order);

        // 2. 发送订单消息到 Kafka
        try {
            String orderJson = objectMapper.writeValueAsString(order);
            kafkaTemplate.send("order-topic", order.getOrderId().toString(), orderJson);
            log.info("Order message sent: {}", order.getOrderId());
        } catch (Exception e) {
            log.error("Failed to send order message", e);
            // 处理发送失败:重试、记录日志、或降级
        }
    }
}

下游服务(Consumer - 以库存服务为例):

@Component
public class InventoryConsumer {
    @KafkaListener(topics = "order-topic", groupId = "inventory-group")
    public void consumeOrder(String key, String orderJson,
                             ConsumerRecord<String, String> record) {
        try {
            Order order = objectMapper.readValue(orderJson, Order.class);
            // 执行扣减库存逻辑
            inventoryService.deductStock(order.getItems());
            log.info("Inventory deducted for order: {}", order.getOrderId());
        } catch (Exception e) {
            log.error("Failed to process order for inventory", e);
            // 关键:不要提交位移!让 Kafka 重试
            // 可以记录错误日志,或发送到死信队列(DLQ)
            throw e;
        }
    }
}

配置要点:

  • Producer:acks=all, retries=Integer.MAX_VALUE, enable.idempotence=true,确保消息不丢失。
  • Consumer:enable.auto.commit=false,手动提交位移,确保处理成功后再提交。
  • 死信队列(DLQ):处理失败的消息可发送到 order-failed-topic,供人工或后台程序处理。

8.3 实时数据统计:使用 Kafka Streams

场景描述: 实时统计每分钟的订单金额总额、订单数,或用户点击流的 PV/UV。

public class RealTimeAnalytics {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, "realtime-analytics");
        props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka1:9092");
        props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
        props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());

        StreamsBuilder builder = new StreamsBuilder();

        // 1. 从订单 Topic 读取数据
        KStream<String, String> orderStream = builder.stream("order-topic");

        // 2. 转换为 <key, Order> 流
        KStream<String, Order> orderObjectStream = orderStream.mapValues(json -> {
            try {
                return objectMapper.readValue(json, Order.class);
            } catch (Exception e) {
                throw new RuntimeException("Deserialization failed", e);
            }
        });

        // 3. 按时间窗口聚合(例如:每分钟)
        KTable<Windowed<String>, Long> orderCountTable = orderObjectStream
            .groupByKey()
            .windowedBy(TimeWindows.of(Duration.ofMinutes(1)))
            .count();

        KTable<Windowed<String>, Double> orderAmountTable = orderObjectStream
            .map((key, order) -> new KeyValue<>(key, order.getAmount()))
            .groupByKey()
            .windowedBy(TimeWindows.of(Duration.ofMinutes(1)))
            .reduce(Double::sum);

        // 4. 将结果写回 Kafka Topic
        orderCountTable.toStream()
            .to("order-count-stats",
                Produced.with(WindowedSerdes.timeWindowedSerdeFrom(String.class),
                              Serdes.Long()));
        orderAmountTable.toStream()
            .to("order-amount-stats",
                Produced.with(WindowedSerdes.timeWindowedSerdeFrom(String.class),
                              Serdes.Double()));

        KafkaStreams streams = new KafkaStreams(builder.build(), props);
        streams.start();

        // 优雅关闭
        Runtime.getRuntime().addShutdownHook(new Thread(streams::close));
    }
}

优势:

  • 轻量级:直接嵌入 Java 应用。
  • Exactly-Once:支持精确一次处理语义。
  • 状态存储:内置状态存储(RocksDB),支持复杂计算。
  • 容错:自动故障恢复。

8.4 消息幂等性与去重处理

问题: 网络重试可能导致 Producer 重复发送,Consumer 处理失败重试导致重复消费。

Producer 端幂等性:

  • 启用 enable.idempotence=true(Kafka 0.11+)。
  • 配合 max.in.flight.requests.per.connection=1(旧版本)或 5(Kafka 2.1+,支持乱序)。
  • 原理:Broker 通过 <ProducerID, SequenceNumber> 识别并去重。

Consumer 端去重:

  • 业务层幂等: 设计幂等的业务操作。
    • 订单:order_id 作为数据库唯一键。
    • 扣库存:基于 order_iditem_id 判断是否已扣减。
    • 发短信:记录 user_id+event 是否已发送。
  • 外部存储去重: 使用 Redis 记录已处理的 message_id 或业务 ID。
public void consume(ConsumerRecord<String, String> record) {
    String messageId = record.key(); // 或从 value 中提取业务 ID
    Boolean added = redisTemplate.opsForSet().add("processed_messages", messageId);
    if (!added) {
        log.info("Message already processed, skipping: {}", messageId);
        return; // 跳过
    }
    // 处理业务逻辑
    processBusiness(record.value());
    // 提交位移
    consumer.commitSync();
}

事务性: Consumer 与下游数据库在同一个事务中提交(较少用)。

8.5 容错与高可用架构设计

目标: 确保 Kafka 集群和应用在节点故障时仍能正常服务。

Broker 高可用:

  • 多副本:replication.factor >= 3,Leader 故障时 ISR 中 Follower 自动选举为新 Leader。
  • 跨机架/可用区部署:Broker 分布在不同物理机或云可用区,避免单点故障。
  • Controller 高可用:Kafka Controller 本身也是多副本,ZooKeeper 选举。

Producer 容错:

  • acks=all:确保消息被 ISR 所有副本确认。
  • retries=Integer.MAX_VALUE:无限重试(结合 retry.backoff.ms)。
  • enable.idempotence=true:防止重试导致重复。
  • 降级策略:发送失败时,可写入本地文件或内存队列,待恢复后重放。

Consumer 容错:

  • 消费者组:多个 Consumer 实例组成 Group,自动负载均衡和故障转移。
  • 位移管理:手动提交(enable.auto.commit=false),确保处理成功再提交;异常时不提交位移,由 Kafka 重试。
  • 死信队列(DLQ):处理失败的消息发送到专门的 Topic,避免阻塞主流程。

网络与 ZooKeeper:

  • ZooKeeper 集群:至少 3 个节点,跨机架部署。
  • 网络分区处理:合理设置 session.timeout.msheartbeat.interval.ms,避免误判。

监控与告警:

  • 监控 UnderReplicatedPartitionsIsrShrinksPerSecRequestHandlerAvgIdlePercent
  • 设置告警,及时发现并处理故障。

总结: 通过多副本、消费者组、幂等生产者、手动位移提交、DLQ 和全面监控,构建一个健壮的 Kafka 应用系统。

第九章:常见问题与最佳实践

9.1 常见错误与排查方法

问题现象可能原因排查方法解决方案
Producer 发送失败(Timeout、NotLeaderForPartition、NetworkException)Broker 宕机或网络不通telnet <broker> 9092 检查网络连通性。确认 Broker 正常运行,网络畅通。
Topic 不存在或分区数变化kafka-topics.sh --describe 查看 Topic 状态和 Leader 分布。创建 Topic 或增加分区。
网络延迟/超时设置过短检查 Broker 日志(server.log)是否有错误。调大 request.timeout.msdelivery.timeout.ms
acks=all 但 ISR 副本不足监控 RequestLatencyMsNetworkProcessorAvgIdlePercent确保 min.insync.replicas <= replication.factor,且 ISR 数量足够。
Consumer 消费延迟(Lag 增长)Consumer 处理速度慢监控 records-lag-maxfetch-rate优化 Consumer 业务逻辑,减少单次处理耗时。
Consumer 实例数不足使用 jstat 或 APM 工具分析 Consumer 应用 GC 情况。增加 Consumer 实例(需增加分区)。
GC 停顿导致心跳超时触发再平衡检查 max.poll.interval.ms 是否小于实际处理时间。调大 max.poll.interval.ms 或拆分长任务。
网络瓶颈监控 Broker 的 BytesOutPerSec启用批量拉取(fetch.min.bytes)。
集群负载不均(某些 Broker/CPU/磁盘高)分区 Leader 分布不均kafka-topics.sh --describe 检查各 Broker 上的 Leader 分区数。执行 kafka-replica-assignment.sh 或使用 Cruise Control 进行 Leader 重平衡。
某些 Topic 流量过大且分区少监控各 Broker 的 BytesInPerSecRequestHandlerAvgIdlePercent为热点 Topic 增加分区数。
数据倾斜(Key 分布不均)分析消息 Key 的分布。优化 Key 的设计,使其更均匀分布。
ZooKeeper 连接频繁中断ZooKeeper 集群性能瓶颈或网络抖动检查 ZooKeeper 服务器的 CPU、内存、网络和磁盘 I/O。优化 ZooKeeper 配置和部署(独立机器,SSD)。
zookeeper.session.timeout.ms 设置过短查看 ZooKeeper 日志是否有 Session expired 错误。适当增大 zookeeper.connection.timeout.mszookeeper.session.timeout.ms
监控 ZooKeeperDisconnects 指标。
磁盘空间不足日志保留策略过于宽松(retention.ms 太大)du -sh /kafka/data 检查磁盘使用。调整 log.retention.hours/bytes
消费者组停滞导致数据无法清理kafka-log-dirs.sh --describe 查看各分区日志大小。清理无用的消费者组或重启停滞的 Consumer。
生产速度远大于消费速度监控 LogEndOffset 和消费者位移,计算滞后量。提升消费能力。

9.2 生产环境部署建议

维度最佳实践
硬件选择
  • CPU:Kafka 依赖 I/O 而非 CPU,中等配置即可。但 Streams 应用可能需要更强 CPU。
  • 内存:OS 缓存至关重要!确保有充足空闲内存用于 Page Cache(至少 32GB+)。JVM 堆通常不宜过大(<8GB),避免长时间 GC。
  • 磁盘:必须使用 SSD。RAID 10 或 JBOD。避免 NAS/SAN。挂载选项:noatime,nodiratime
  • 网络:万兆网卡(10GbE),低延迟交换机。
集群规划
  • Broker 数量:至少 3 个,建议奇数个(便于 ZooKeeper 选举)。根据吞吐和存储需求扩展。
  • Topic 设计:预估峰值吞吐,合理设置初始分区数(可后续增加)。避免创建过多小 Topic。
  • 副本数:replication.factor=3 是生产标准。
操作系统调优
  • 关闭透明大页(echo never > /sys/kernel/mm/transparent_hugepage/enabled)。
  • 调整文件句柄数(ulimit -n 100000)。
  • 调整虚拟内存 vm.swappiness=1
  • 使用 deadline 或 noop 磁盘调度器。
Kafka 配置
  • num.io.threads = 磁盘数 × 2 ~ 4。
  • log.flush.interval.ms 通常不设,依赖 OS 刷盘。
  • auto.create.topics.enable=false,防止意外创建 Topic。
  • 启用 delete.topic.enable=true
监控告警必须集成 Prometheus + Grafana + Alertmanager。
核心告警指标:
  • UnderReplicatedPartitions > 0
  • ActiveControllerCount != 1
  • OfflinePartitionsCount > 0
  • ConsumerLag > 10000(阈值自定)
  • RequestHandlerAvgIdlePercent < 30%
备份与恢复
  • 定期使用 kafka-mirror-maker.sh 或 MirrorMaker 2.0 进行跨集群复制(灾备)。
  • 重要 Topic 启用压缩(cleanup.policy=compact)结合时间删除。

9.3 安全配置(SSL/SASL)

Kafka 提供多层安全机制:

加密(SSL/TLS):

  • 作用:加密客户端与 Broker、Broker 间、Broker 与 ZooKeeper 间的通信。
  • 配置:
    • 生成 CA 证书和密钥。
    • 为每个 Broker 签发证书。
# Broker 配置(server.properties)
listeners=SSL://:9093
security.inter.broker.protocol=SSL
ssl.keystore.location=/path/to/server.keystore.jks
ssl.keystore.password=keystore_password
ssl.key.password=key_password
ssl.truststore.location=/path/to/server.truststore.jks
ssl.truststore.password=truststore_password
ssl.client.auth=required  # 开启客户端认证

Client 配置需提供信任库。

认证(SASL):

  • SASL/PLAIN:简单用户名/密码认证。必须配合 SSL 使用,否则密码明文传输。
# server.properties
sasl.enabled.mechanisms=PLAIN
sasl.mechanism.inter.broker.protocol=PLAIN
// JAAS 配置文件
KafkaServer {
  org.apache.kafka.common.security.plain.PlainLoginModule required
  username="admin"
  password="admin-secret"
  user_admin="admin-secret"
  user_alice="alice-secret";
};
  • SASL/SCRAM:比 PLAIN 更安全,支持凭证动态增删。
  • SASL/GSSAPI(Kerberos):企业级认证,集成 AD/KDC。

授权(ACL):

  • 控制用户对 Topic、Consumer Group 等资源的操作权限(Read、Write、Describe、Create…)。
  • 启用:authorizer.class.name=kafka.security.authorizer.AclAuthorizer
# 示例命令
kafka-acls.sh --add --allow-principal User:alice --operation Read --topic orders-topic
kafka-acls.sh --add --allow-principal User:bob --operation Write --topic logs-topic

部署建议: 生产环境应同时启用 SSL(加密)+ SASL/SCRAM(认证)+ ACL(授权)。

9.4 版本升级与兼容性注意事项

升级路径:

  • Kafka 遵循语义化版本控制,主版本号变更可能引入不兼容改动。
  • 推荐:逐个 Broker 滚动升级。先升级 Follower,最后升级 Controller。
  • 使用 kafka-broker-api-versions.sh 检查 API 兼容性。

兼容性:

  • Producer/Consumer 与 Broker:通常向后兼容。新版本 Client 可连接旧 Broker,但可能无法使用新特性。旧 Client 连接新 Broker 一般没问题。
  • 消息格式:message.format.version 需与 inter.broker.protocol.version 协调。升级前确保所有 Broker 版本一致后再提升协议版本。
  • ZooKeeper:Kafka 3.0+ 开始逐步用 KRaft 替代 ZooKeeper。未来升级需考虑元数据存储的迁移。

关键步骤:

  • 备份: 备份 ZooKeeper 数据和 Kafka 日志。
  • 测试: 在预发环境充分测试新版本。
  • 文档: 仔细阅读官方发布说明(Release Notes),了解 Breaking Changes。
  • Rollback Plan: 准备好回滚到旧版本的方案。

9.5 Kafka 与其他中间件对比(RabbitMQ、RocketMQ)

特性Apache KafkaRabbitMQApache RocketMQ
架构模型发布/订阅(基于日志)AMQP 0.9.1,支持多种模式(Pub/Sub、Point-to-Point、Routing、RPC)发布/订阅,点对点
吞吐量极高(10w+ msg/s)中等(万级别)高(接近 Kafka)
延迟毫秒级(批处理时稍高)微秒级(低延迟)毫秒级
持久化基于磁盘日志,顺序 I/O,极高效内存+磁盘,可持久化基于磁盘日志
伸缩性极强,水平扩展容易,分区是并行单元较弱,队列不易扩展,镜像队列非线性扩展强,类似 Kafka
消息顺序分区内严格有序单队列有序Topic 内有序
消息路由简单(基于 Topic/Partition)极其灵活(Exchange + Binding + Routing Key)较灵活(Tag、SQL 过滤)
功能丰富度核心功能稳定,生态强大(Streams、Connect)功能最丰富(插件、管理界面、多种协议)功能全面(事务消息、定时消息、死信队列)
运维复杂度中等偏高(依赖 ZooKeeper/KRaft)低(单节点简单,集群需镜像)中等
适用场景大数据管道、日志收集、流处理、高吞吐解耦传统企业应用、复杂路由、RPC、低延迟任务金融、电商(阿里系),需要事务/定时消息的场景

选型建议:

  • 选 Kafka: 追求极致吞吐、大数据集成、流式处理、作为”数据中心”。
  • 选 RabbitMQ: 需要复杂路由规则、低延迟、协议多样(MQTT、STOMP)、易上手。
  • 选 RocketMQ: 需要强事务消息、定时/延时消息、在中国云环境下深度集成。