Kafka
Kafka 不是把消息交出去便删除的队列。它把事件追加到分区日志,让生产者、Broker、消费者和外部数据库分别推进自己的位置。理解这些位置之间的空隙,比背诵参数更重要:生产者超时不等于写入失败,数据库提交不等于 offset 已提交,副本在线也不等于已经追到可选举的位置。
一、是什么
一条订单事件如何穿过 Kafka
订单服务把 order-1001 created 发给 bootstrap 地址。客户端取得元数据,依据 key 的序列化字节选择分区,把多条 record 聚成 record batch 后发给该分区 leader。leader 追加日志,ISR follower 拉取并确认;当 acks=all 与 min.insync.replicas 的条件满足,生产者才得到成功响应。
消费者不是从 Broker “领走”事件,而是读取分区 offset。Consumer Group 由协调器分配分区;应用把事件写入 PostgreSQL,再提交下一条待读 offset。若进程死在数据库提交之后、offset 提交之前,同一事件会再次到达,因此外部数据库必须用稳定 event id 去重。
生产者保存幂等序列与尚未确认的请求;分区日志保存 record batch 和 offset;组协调器保存成员、分配和 committed offset;PostgreSQL 保存业务结果与已处理 event id。任何“恰好一次”判断都必须先说明覆盖的是哪些边界。
Record batch、segment、HW 与 LSO
record 包含 key、value、headers 和时间戳。网络传输、压缩、磁盘追加与复制的基本单位通常是 record batch,因此增大 linger.ms 能改善吞吐,也会增加等待。分区文件由多个 segment 组成:active segment 持续追加,关闭的 segment 才适合按时间、大小或压缩策略清理。
offset 是分区内逻辑位置,不是全局序号;事务控制批次、回滚与保留清理都可能使业务可见 offset 不连续。log start offset 是仍可读取的最早位置,LEO 是副本日志末端,HW 是已安全复制的可读上界。事务消费者使用 read_committed 时还受 LSO 限制,不会越过最早未完成事务。磁盘里已有数据而消费者读不到时,需要分别检查 HW、LSO 和过滤条件。
Partition 决定顺序与并行度
Kafka 只保证单分区内的顺序。相同订单必须使用稳定 key,才能在分区器不变时进入同一分区。传统 Consumer Group 中,一个分区同一时刻最多交给一个成员,因此分区数也是并行上限。
增加分区会改变默认 key 映射,一部分既有 key 的新事件可能落入新分区,跨分区便不再有顺序关系;更多分区还会增加 leader、文件句柄、复制和恢复开销。要求长期顺序的主题应先估算峰值吞吐与未来并行度,再固定 key 契约。
ISR、ELR 与不干净选举
leader 接收读写,ISR 是当前与 leader 同步、具备正常选举资格的副本集合。三副本、min.insync.replicas=2 与 acks=all 联合表示一次写入至少要得到两个同步副本支持;副本不足时拒绝写入,是以可用性保护已确认数据。
Kafka 4.3 新建集群默认启用 Eligible Leader Replicas。ELR 记录最近仍足够同步、可在 ISR 为空时优先恢复领导权的副本。它不是额外副本,也不能代替跨故障域部署。开启不干净选举仍可能以丢失已确认 record 换取恢复服务,不能当作日常修复。
KRaft 元数据 quorum
Kafka 4.3.1 使用 KRaft,不再依赖 ZooKeeper。controller quorum 通过 Raft 保存 broker、topic、partition、ACL 和 feature level 等元数据,由 active controller 作出领导权决策。controller 元数据日志、快照与多数派是集群恢复的根。
开发机可以让单进程同时承担 broker 与 controller;生产环境应分离角色,通常用三个或五个 controller。动态 quorum 可先建立一个 controller,再把 observer 安全加入。无论采用静态还是动态 quorum,同一集群的 cluster id 必须一致,已有数据目录绝不能因 controller 暂时不可用而重新 format。
Consumer、Share 与 Streams Group
Consumer Group 把分区所有权分给成员,适合按分区保序并由应用提交 offset。Kafka 4.3 客户端默认仍使用 classic 协议;使用新 consumer group protocol 需要显式设置 group.protocol=consumer,心跳和会话参数由 broker 管理。
客户端配置只声明协议和业务参数,不再设置 classic 协议的 session timeout:
group.id=orders-projector-v1
group.protocol=consumer
enable.auto.commit=false
auto.offset.reset=none
isolation.level=read_committed新协议的会话边界由 broker 的 group.consumer.session.timeout.ms 和 group.consumer.heartbeat.interval.ms 管理。迁移前要删除客户端的 session.timeout.ms、heartbeat.interval.ms 与 partition.assignment.strategy;Kafka 4.3 并不会为旧客户端自动切换协议。
Share Group 面向按 record 并行的工作队列场景。成员可以共享同一分区,Broker 保存获取锁、确认结果和投递次数;成员可 accept、release 或 reject。确认并不立刻从日志删除 record,保留仍由 topic 策略决定。需要严格 key 顺序时应使用 Consumer Group。Streams Group 则服务于 Kafka Streams 的任务分配与状态管理,三者是不同处理模型而非性能档位。
Kafka 4.3 已把 Share Group 作为生产可用能力,但 4.3 broker 的默认 group.coordinator.rebalance.protocols=classic,consumer,streams 仍不包含 share。要使用它,须先在所有 broker 增加 share 并完成滚动重启,再由管理员检查和提升 share.version;客户端使用 KafkaShareConsumer,不能把普通 KafkaConsumer 的 group.protocol 改成 share 冒充。启用命令放在生产变更窗口执行:
kafka-features.sh --bootstrap-server kafka1.example.net:9094 \
--command-config /run/secrets/kafka-admin.properties describe
kafka-features.sh --bootstrap-server kafka1.example.net:9094 \
--command-config /run/secrets/kafka-admin.properties \
upgrade --feature share.version=1二、为什么
可回放日志把生产与消费解耦
同一份订单事件可以同时驱动库存、搜索、风控和审计,每个订阅者维护自己的位置。新消费者可从历史重建状态,旧消费者可修复后回放。这个能力来自事件在保留期内仍在日志里,而非来自消费者确认。
Kafka 适合持续事件流、多个独立订阅者和可容忍异步的业务。低流量命令、单条复杂路由或要求 Broker 直接管理每条任务重试时,传统队列可能更自然。Share Group 缩小了差距,但不会把 Kafka 的存储模型变成删除式队列。
可靠性是一组联合约束
replication factor 只有在副本跨主机或故障域时才有意义。acks=all 依赖 ISR,min.insync.replicas 决定最少同步副本数,生产者重试又依赖幂等性。只改单个参数,常得到与名称相反的结果。
幂等 producer 避免同一生产会话因重试产生重复序列,但网络超时仍可能让调用方不知道写入结果。业务若每次重试生成新 event id,就会造成业务重复,因此 event id 必须稳定。
Kafka transaction 能原子地写多个 Kafka 分区并提交消费 offset,却不能把 PostgreSQL commit 自动纳入同一事务。跨数据库边界通常使用 Transactional Outbox 或 Inbox。Outbox 在业务数据库事务中写业务行和待发事件;Inbox 在消费数据库事务中先插入唯一 event id,再更新投影。它们接受至少一次传递,把重复变成可判定、可重放的正常路径。
Retention、compaction 与 tiered storage
cleanup.policy=delete 按时间或大小删除旧 segment,适合审计流水与可回放事件。compact 为每个 key 最终保留最近值,用 tombstone 表示删除;压缩是后台渐进过程,读取时仍可能暂时看到旧值。
Tiered Storage 把关闭的 segment 复制到远端,降低本地磁盘压力并延长历史。Kafka 只定义 RemoteStorageManager 接口,不附带可直接用于生产的实现,而且 compacted topic 不支持 tiering。选型必须验证供应商插件、对象存储一致性、删除语义、恢复时间和成本。
自建、托管与跨集群复制
托管 Kafka 减少 controller、磁盘和升级工作,但 topic 设计、key 契约、消费者幂等与外部事务仍属于应用。自建集群还要负责证书、SASL、ACL、容量、机架分布与故障演练。
MirrorMaker 2 基于 Kafka Connect 复制 topic、heartbeat 与 checkpoint。它是异步复制,不是同步双活共识;两端同时写相同业务 key 会产生冲突。灾备切换必须先围栏旧生产者,再确认复制差距和目标端消费位置。
三、怎么做
用 Compose 跑通开发链
开发环境固定使用官方 apache/kafka:4.3.1 和 PostgreSQL 18.6。单节点 KRaft 只用于理解调用链,不代表生产拓扑。
KAFKA_CLUSTER_ID=MkU3OEVBNTcwNTJENDM2Qk
POSTGRES_PASSWORD=dev-only-change-meservices:
kafka:
image: apache/kafka:4.3.1
hostname: kafka
ports: ["127.0.0.1:19092:19092"]
environment:
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:29093
KAFKA_LISTENERS: CONTROLLER://:29093,INTERNAL://:29092,HOST://:19092
KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka:29092,HOST://localhost:19092
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,INTERNAL:PLAINTEXT,HOST:PLAINTEXT
KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
CLUSTER_ID: ${KAFKA_CLUSTER_ID}
postgres:
image: postgres:18.6
ports: ["127.0.0.1:15432:5432"]
environment:
POSTGRES_DB: orders
POSTGRES_USER: orders_app
POSTGRES_PASSWORD: ${POSTGRES_PASSWORD}
volumes: ["./sql:/docker-entrypoint-initdb.d:ro"]
healthcheck:
test: ["CMD-SHELL", "pg_isready -U orders_app -d orders"]
interval: 5s
timeout: 3s
retries: 20Compose 会在 PostgreSQL 首次初始化时执行 ./sql 目录中的脚本。先创建 sql/001-inbox.sql,再启动容器;若已有旧数据卷,新增初始化脚本不会自动补跑,应显式迁移或先确认可以重建开发数据。
CREATE TABLE consumed_event (
event_id text PRIMARY KEY,
topic text NOT NULL,
partition_id integer NOT NULL,
offset_id bigint NOT NULL,
consumed_at timestamptz NOT NULL DEFAULT now()
);
CREATE TABLE order_projection (
order_id text PRIMARY KEY,
last_event_id text NOT NULL UNIQUE,
last_event_type text NOT NULL,
updated_at timestamptz NOT NULL DEFAULT now()
);docker compose up -d
until docker compose exec -T postgres pg_isready -U orders_app -d orders; do sleep 2; done
until docker compose exec -T kafka /opt/kafka/bin/kafka-metadata-quorum.sh \
--bootstrap-server kafka:29092 describe --status; do sleep 2; done
docker compose exec -T kafka /opt/kafka/bin/kafka-metadata-quorum.sh \
--bootstrap-server kafka:29092 describe --status
docker compose exec -T kafka /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server kafka:29092 --create \
--topic orders.events.v1 --partitions 3 --replication-factor 1
printf 'order-1001\tevt-001|order-1001|CREATED\n' |
docker compose exec -T kafka /opt/kafka/bin/kafka-console-producer.sh \
--bootstrap-server kafka:29092 --topic orders.events.v1 \
--reader-property parse.key=true --reader-property key.separator=$'\t' \
--command-property acks=all --command-property enable.idempotence=true
docker compose exec kafka /opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server kafka:29092 --topic orders.events.v1 \
--from-beginning --max-messages 1 \
--formatter-property print.key=true \
--formatter-property print.partition=true \
--formatter-property print.offset=true这些是 Kafka 4.2 以后替代旧控制台参数的写法。容器内客户端使用 kafka:29092,宿主机应用使用 localhost:19092;advertised listener 返回不可达地址时,bootstrap 可以成功,后续元数据连接仍会失败。
用 Java 与 PostgreSQL Inbox 处理重复
项目使用 JDK 17、kafka-clients:4.3.1 和 postgresql:42.7.13。生产项目应使用 schema 管理;这里用 eventId|orderId|type 聚焦提交顺序。
<project xmlns="http://maven.apache.org/POM/4.0.0">
<modelVersion>4.0.0</modelVersion>
<groupId>demo</groupId><artifactId>kafka-inbox-demo</artifactId><version>1.0.0</version>
<properties><maven.compiler.release>17</maven.compiler.release></properties>
<dependencies>
<dependency><groupId>org.apache.kafka</groupId><artifactId>kafka-clients</artifactId><version>4.3.1</version></dependency>
<dependency><groupId>org.postgresql</groupId><artifactId>postgresql</artifactId><version>42.7.13</version></dependency>
</dependencies>
<build><plugins>
<plugin><groupId>org.apache.maven.plugins</groupId><artifactId>maven-compiler-plugin</artifactId><version>3.14.1</version><configuration><release>17</release></configuration></plugin>
<plugin><groupId>org.codehaus.mojo</groupId><artifactId>exec-maven-plugin</artifactId><version>3.5.1</version></plugin>
</plugins></build>
</project>生产者复用稳定 event id,并等待可靠确认。
package demo;
import java.util.Properties;
import org.apache.kafka.clients.producer.*;
import org.apache.kafka.common.serialization.StringSerializer;
public final class OrderProducer {
public static void main(String[] args) throws Exception {
if (args.length != 2) throw new IllegalArgumentException("eventId orderId");
Properties p = new Properties();
p.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,
System.getenv().getOrDefault("KAFKA_BOOTSTRAP", "localhost:19092"));
p.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
p.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
p.put(ProducerConfig.ACKS_CONFIG, "all");
p.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "true");
try (KafkaProducer<String,String> producer = new KafkaProducer<>(p)) {
String value = args[0] + "|" + args[1] + "|CREATED";
RecordMetadata m = producer.send(
new ProducerRecord<>("orders.events.v1", args[1], value)).get();
System.out.printf("partition=%d offset=%d%n", m.partition(), m.offset());
}
}
}消费者的核心必须是一个数据库事务:先执行 INSERT INTO consumed_event ... ON CONFLICT DO NOTHING;只有插入成功才 upsert order_projection,然后 commit。事务成功后,Kafka 客户端再提交该分区的 record.offset() + 1。
package demo;
import java.sql.DriverManager;
import java.time.Duration;
import java.util.List;
import java.util.Map;
import java.util.Properties;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.serialization.StringDeserializer;
public final class OrderProjector {
public static void main(String[] args) throws Exception {
Properties p = new Properties();
p.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,
System.getenv().getOrDefault("KAFKA_BOOTSTRAP", "localhost:19092"));
p.put(ConsumerConfig.GROUP_ID_CONFIG, "orders-projector-v1");
p.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
p.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
p.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
p.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
p.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed");
p.put("group.protocol", "consumer");
try (var consumer = new KafkaConsumer<String, String>(p);
var db = DriverManager.getConnection(
System.getenv("JDBC_URL"),
System.getenv("JDBC_USER"),
System.getenv("JDBC_PASSWORD"))) {
consumer.subscribe(List.of("orders.events.v1"));
while (true) {
for (var record : consumer.poll(Duration.ofSeconds(1))) {
String[] event = record.value().split("\\|", -1);
if (event.length != 3) {
throw new IllegalArgumentException("invalid event: " + record.value());
}
db.setAutoCommit(false);
boolean first;
try {
try (var inbox = db.prepareStatement(
"INSERT INTO consumed_event(event_id,topic,partition_id,offset_id) " +
"VALUES (?,?,?,?) ON CONFLICT DO NOTHING")) {
inbox.setString(1, event[0]);
inbox.setString(2, record.topic());
inbox.setInt(3, record.partition());
inbox.setLong(4, record.offset());
first = inbox.executeUpdate() == 1;
}
if (first) {
try (var projection = db.prepareStatement(
"INSERT INTO order_projection(order_id,last_event_id,last_event_type) " +
"VALUES (?,?,?) ON CONFLICT(order_id) DO UPDATE SET " +
"last_event_id=excluded.last_event_id," +
"last_event_type=excluded.last_event_type,updated_at=now()")) {
projection.setString(1, event[1]);
projection.setString(2, event[0]);
projection.setString(3, event[2]);
projection.executeUpdate();
}
}
db.commit();
if ("1".equals(System.getenv("CRASH_AFTER_DB_COMMIT"))) {
Runtime.getRuntime().halt(17);
}
} catch (Exception failure) {
db.rollback();
throw failure;
}
var partition = new TopicPartition(record.topic(), record.partition());
consumer.commitSync(Map.of(
partition, new OffsetAndMetadata(record.offset() + 1)));
System.out.println("eventId=" + event[0] + " businessChanged=" + first);
}
}
}
}
}把两个 Java 文件放入 src/main/java/demo。数据库事务提交后才推进 offset;若 offset commit 超时,第二次 Inbox 插入命中唯一键,业务投影不再变化,随后同一 group 才把位置推进。先编译并发送事件,再在一个终端模拟“数据库已提交、offset 未提交”的崩溃;随后在第二个终端去掉故障开关重启同一 group。消费者是常驻进程,看到第二次输出后用 Ctrl+C 停止。
export KAFKA_BOOTSTRAP=localhost:19092
export JDBC_URL='jdbc:postgresql://localhost:15432/orders'
export JDBC_USER=orders_app
export JDBC_PASSWORD='dev-only-change-me'
mvn -q compile exec:java -Dexec.mainClass=demo.OrderProducer \
-Dexec.args='evt-002 order-1002'
CRASH_AFTER_DB_COMMIT=1 mvn -q compile exec:java \
-Dexec.mainClass=demo.OrderProjector
# 上一条命令以 17 退出后,在第二个终端沿用相同 JDBC/Kafka 环境变量:
mvn -q compile exec:java -Dexec.mainClass=demo.OrderProjector
docker compose exec postgres psql -U orders_app -d orders -c \
"SELECT c.event_id,p.order_id,p.last_event_type FROM consumed_event c JOIN order_projection p ON p.last_event_id=c.event_id;"第一轮应在数据库中留下 evt-002 后退出,第二轮日志应显示 businessChanged=false,查询仍只有一份事件和一份投影;这才证明 Inbox 吸收了真实重复窗口。生产环境通常把 auto.offset.reset 设为 none,避免位置丢失时静默跳转。
在 Linux 上拆分 KRaft 角色
从 Apache 下载 Kafka 4.3.1 二进制包和同目录 SHA-512 文件,校验后以固定目录安装。生产 controller 与 broker 使用不同节点、目录和证书。
curl -fLO https://downloads.apache.org/kafka/4.3.1/kafka_2.13-4.3.1.tgz
curl -fLO https://downloads.apache.org/kafka/4.3.1/kafka_2.13-4.3.1.tgz.sha512
expected=$(sed '1s/^[^:]*:[[:space:]]*//;2,$s/^[[:space:]]*//' \
kafka_2.13-4.3.1.tgz.sha512 | tr -d '[:space:]' | tr '[:upper:]' '[:lower:]')
actual=$(sha512sum kafka_2.13-4.3.1.tgz | awk '{print $1}')
test "$actual" = "$expected" || { echo 'SHA-512 mismatch' >&2; exit 1; }
sudo useradd --system --home /var/lib/kafka --shell /usr/sbin/nologin kafka
sudo install -d -o kafka -g kafka -m 0750 /var/lib/kafka/data /var/lib/kafka/meta /var/log/kafka
sudo tar -xzf kafka_2.13-4.3.1.tgz -C /opt
sudo ln -sfn /opt/kafka_2.13-4.3.1 /opt/kafkacontroller 节点保存独立配置;三台机器分别使用 node.id=1、2、3,并替换监听 FQDN 与各自证书。min.insync.replicas 在 ELR 启用时是 controller 维护的集群级默认值,因此必须在首个 controller 建群前统一写入 controller 配置,不能只放进某台 broker 的本地文件。
process.roles=controller
node.id=1
controller.listener.names=CONTROLLER
controller.quorum.bootstrap.servers=kctrl1.example.net:9093,kctrl2.example.net:9093,kctrl3.example.net:9093
listeners=CONTROLLER://kctrl1.example.net:9093
listener.security.protocol.map=CONTROLLER:SSL
security.protocol=SSL
listener.name.controller.ssl.client.auth=required
ssl.keystore.type=PKCS12
ssl.keystore.location=/etc/kafka/tls/controller.p12
ssl.keystore.password=RENDER_FROM_SECRET
ssl.truststore.type=PKCS12
ssl.truststore.location=/etc/kafka/tls/ca.p12
ssl.truststore.password=RENDER_FROM_SECRET
authorizer.class.name=org.apache.kafka.metadata.authorizer.StandardAuthorizer
allow.everyone.if.no.acl.found=false
super.users=User:CN=kafka-admin;User:CN=kafka-controller
min.insync.replicas=2
metadata.log.dir=/var/lib/kafka/metabroker 节点只承担数据面;下面还把 CLIENT 的 SCRAM 与内部 mTLS 分开,避免全局强制客户端证书后,只有 truststore 的 SCRAM 应用永远握手失败。
process.roles=broker
node.id=101
broker.rack=az-a
controller.listener.names=CONTROLLER
controller.quorum.bootstrap.servers=kctrl1.example.net:9093,kctrl2.example.net:9093,kctrl3.example.net:9093
listeners=INTERNAL://:9092,CLIENT://:9094,ADMIN://127.0.0.1:9095
advertised.listeners=INTERNAL://kafka1.example.net:9092,CLIENT://kafka1.example.net:9094
listener.security.protocol.map=CONTROLLER:SSL,INTERNAL:SSL,CLIENT:SASL_SSL,ADMIN:SSL
inter.broker.listener.name=INTERNAL
ssl.keystore.type=PKCS12
ssl.keystore.location=/etc/kafka/tls/broker.p12
ssl.keystore.password=RENDER_FROM_SECRET
ssl.truststore.type=PKCS12
ssl.truststore.location=/etc/kafka/tls/ca.p12
ssl.truststore.password=RENDER_FROM_SECRET
listener.name.internal.ssl.client.auth=required
listener.name.admin.ssl.client.auth=required
listener.name.client.ssl.client.auth=none
listener.name.client.sasl.enabled.mechanisms=SCRAM-SHA-512
listener.name.client.scram-sha-512.sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required;
authorizer.class.name=org.apache.kafka.metadata.authorizer.StandardAuthorizer
allow.everyone.if.no.acl.found=false
super.users=User:CN=kafka-admin;User:CN=kafka-controller
default.replication.factor=3
offsets.topic.replication.factor=3
transaction.state.log.replication.factor=3
transaction.state.log.min.isr=2
log.dirs=/var/lib/kafka/dataCLIENT 的 JAAS 配置让 Broker 加载 SCRAM 服务端模块,应用密码仍通过管理命令写入 Kafka 元数据;此处内部通信使用 mTLS,不把应用用户名和密码写进服务端模块。认证通过与 ACL 授权要分别验证。
每个节点必须替换唯一 node id、FQDN、机架和 PKCS12 密钥库,密码由主机密钥系统渲染进仅 kafka 可读的配置。controller 文件保存为 /etc/kafka/controller.properties;其中 security.protocol=SSL 供后续 Admin CLI 使用。kctrl2、kctrl3 的 listeners 必须分别填自己的可达 FQDN,不能保留空 host:成员添加工具可能将空值登记为 localhost。动态 quorum 与前面的单节点静态 Compose 是两套互斥方案,生产配置不得再加入 controller.quorum.voters。先在 kctrl1 生成一次 cluster id 并安全分发到其余新节点;格式化前必须确认元数据目录为空:
先由管理员在各角色主机保存 unit。broker 使用 /etc/systemd/system/kafka-broker.service:
[Unit]
Description=Apache Kafka broker
After=network-online.target
Wants=network-online.target
[Service]
User=kafka
Group=kafka
ExecStart=/opt/kafka/bin/kafka-server-start.sh /etc/kafka/broker.properties
ExecStop=/opt/kafka/bin/kafka-server-stop.sh
Restart=on-failure
RestartSec=10
LimitNOFILE=100000
TimeoutStopSec=180
[Install]
WantedBy=multi-user.targetcontroller 主机把 Description 改为 Apache Kafka controller、ExecStart 的配置改为 /etc/kafka/controller.properties,保存为 /etc/systemd/system/kafka-controller.service。网络只放行 broker 间、controller quorum 与业务客户端所需端口;证书 SAN 必须覆盖 FQDN,数据目录和私钥不得给应用身份读取。
保存后在每台主机执行 sudo systemctl daemon-reload,再做下面的格式化与启动。服务文件不创建集群元数据,已有数据的节点不重复 format。
# kctrl1:建立只有一个 voter 的初始动态 quorum
CLUSTER_ID=$(/opt/kafka/bin/kafka-storage.sh random-uuid)
sudo -u kafka /opt/kafka/bin/kafka-storage.sh format \
--cluster-id "$CLUSTER_ID" --config /etc/kafka/controller.properties --standalone
sudo systemctl enable --now kafka-controller
# kctrl2、kctrl3:使用同一 CLUSTER_ID,各自先作为 observer 启动
sudo -u kafka /opt/kafka/bin/kafka-storage.sh format \
--cluster-id "$CLUSTER_ID" --config /etc/kafka/controller.properties \
--no-initial-controllers
sudo systemctl enable --now kafka-controller
sudo -u kafka /opt/kafka/bin/kafka-metadata-quorum.sh \
--bootstrap-controller kctrl1.example.net:9093 \
--command-config /etc/kafka/controller.properties describe --replication
sudo -u kafka /opt/kafka/bin/kafka-metadata-quorum.sh \
--bootstrap-controller kctrl1.example.net:9093 \
--command-config /etc/kafka/controller.properties add-controller --dry-run在待加入的 kctrl2 本机执行以上检查:复制输出中该节点应为 Observer,持续取得元数据并追平;dry-run 输出的 node ID、directory ID、FQDN 和端口应与本机一致。CLI 从配置读取 node.id、controller listener 与元数据目录,再读取本机 meta.properties 的 directory ID,普通管理客户端配置不能替代这些输入。dry-run 只检查本地参数,不验证网络、追平或真实成员变更;确认上述条件后才执行:
sudo -u kafka /opt/kafka/bin/kafka-metadata-quorum.sh \
--bootstrap-controller kctrl1.example.net:9093 \
--command-config /etc/kafka/controller.properties add-controller
sudo -u kafka /opt/kafka/bin/kafka-metadata-quorum.sh \
--bootstrap-controller kctrl1.example.net:9093 \
--command-config /etc/kafka/controller.properties describe --status
# 每台全新 broker:确认 log.dirs 为空后格式化并启动
sudo -u kafka /opt/kafka/bin/kafka-storage.sh format \
--cluster-id "$CLUSTER_ID" --config /etc/kafka/broker.properties \
--no-initial-controllers
sudo systemctl enable --now kafka-broker确认 kctrl2 进入 CurrentVoters 后,再在 kctrl3 本机重复 observer、追平、dry-run 和 add-controller。不得为已存在数据的节点重新格式化。最终用 describe --status 确认 leader 和三名 voter,并用集群动态配置检查 min.insync.replicas=2。成员操作见 KRaft 运维。在已经写入 ELR 状态的集群上修改该值会清空现有 ELR,因此后续变更必须独立评估并观察 ISR/ELR,而不是临时救火。
TLS、SASL 与 ACL
应用配置放入 0600 secret 文件,管理命令统一使用 --command-config。不要在进程参数中传 SCRAM 密码;凭据创建与轮换应通过 Admin API 和密钥工作流完成。
bootstrap.servers=kafka1.example.net:9094,kafka2.example.net:9094,kafka3.example.net:9094
security.protocol=SASL_SSL
sasl.mechanism=SCRAM-SHA-512
sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required username="orders-consumer" password="RENDER_FROM_SECRET";
ssl.truststore.type=PKCS12
ssl.truststore.location=/run/secrets/kafka-client-ca.p12
ssl.truststore.password=RENDER_FROM_SECRET
ssl.endpoint.identification.algorithm=httpskafka-acls.sh --bootstrap-server kafka1.example.net:9094 \
--command-config /run/secrets/kafka-admin.properties \
--add --allow-principal 'User:orders-producer' \
--operation Write --operation Describe --topic orders.events.v1
kafka-acls.sh --bootstrap-server kafka1.example.net:9094 \
--command-config /run/secrets/kafka-admin.properties \
--add --allow-principal 'User:orders-consumer' \
--operation Read --operation Describe --topic orders.events.v1
kafka-acls.sh --bootstrap-server kafka1.example.net:9094 \
--command-config /run/secrets/kafka-admin.properties \
--add --allow-principal 'User:orders-consumer' \
--operation Read --group orders-projector-v1生产者和消费者使用不同的 SCRAM 用户、secret 文件和最小 ACL;上面的客户端文件属于消费者,生产者文件只把用户名换成 orders-producer。内部 broker、controller 与 ADMIN listener 继续要求 mTLS,公网业务 listener 则以 SCRAM 认证,不混用两套身份边界。
身份认证成功而 ACL 失败时会出现 TopicAuthorizationException、GroupAuthorizationException、ClusterAuthorizationException 或其他授权拒绝,不应把账号加入 super users;应核对服务端看到的 principal、resource pattern 和操作。
灾备、offset 与升级
MirrorMaker 2 的最小方向是 primary 到 dr。两端还需要各自的 TLS/SASL 客户端配置,复制身份不应是超级用户。
clusters=primary,dr
primary.bootstrap.servers=kafka-a1:9094,kafka-a2:9094,kafka-a3:9094
dr.bootstrap.servers=kafka-b1:9094,kafka-b2:9094,kafka-b3:9094
primary->dr.enabled=true
primary->dr.topics=orders\.events\..*
primary->dr.groups=orders-projector-.*
replication.policy.class=org.apache.kafka.connect.mirror.DefaultReplicationPolicy
emit.heartbeats.enabled=true
emit.checkpoints.enabled=true
sync.group.offsets.enabled=true
replication.factor=3
heartbeats.topic.replication.factor=3
checkpoints.topic.replication.factor=3
offset-syncs.topic.replication.factor=3默认复制策略会把目标 topic 命名为 primary.orders.events.v1,灾备消费者必须订阅这个名字;若改用 IdentityReplicationPolicy,则要先证明目标端不存在同名双写和环形复制。checkpoint 只应更新目标端未活跃的 group,而且翻译位置依赖最新 offset-sync,不能把“有 checkpoint”误判成零 RPO。MM2 的 topic 配置和 ACL 同步也不等于复制了 SCRAM 用户、group ACL、目标集群 min.insync.replicas、Schema Registry 或数据库,这些必须在 DR 端独立配置和验证。切换前先围栏源生产者并停止源消费者,确认 heartbeat、offset-sync/checkpoint 新鲜度、目标 end offset 与可接受 RPO,再让拥有目标 topic/group ACL 的目标消费者以恢复 group 启动并核对业务计数,最后切入口。回切同样要求单写围栏。
升级先逐节点替换二进制,持续观察 controller quorum、ISR、请求错误和消费延迟。4.3.0 到 4.3.1 的补丁升级不需要提升 feature level;从较早 KRaft 版本到 4.3,所有节点稳定运行新版本后才执行:
kafka-features.sh --bootstrap-controller kctrl1.example.net:9093 \
--command-config /run/secrets/kafka-admin.properties describe
kafka-features.sh --bootstrap-controller kctrl1.example.net:9093 \
--command-config /run/secrets/kafka-admin.properties \
upgrade --release-version 4.34.3 元数据 feature 完成提升后不支持降级。旧二进制只能覆盖 finalization 之前的回退窗口;变更前必须确认 controller 快照、客户端兼容与回滚边界。
四、问题处理
Bootstrap 成功,随后连接超时
常见原因是 bootstrap 返回了客户端不可达的 advertised listener,或证书 SAN 与 FQDN 不匹配。先用 getent hosts 检查每个 advertised FQDN,再用 openssl s_client -connect host:9094 -servername host -verify_hostname host -verify_return_error 检查 TLS,最后从应用所在网络重新取元数据。修正 listener、DNS 或证书后,必须再次验证全部 broker 端点,而非只测 bootstrap。
KRaft 没有 active controller
metadata quorum 无 leader,通常意味着多数派不可达、node id 重复、cluster id 不一致或 TLS 身份失败。结合三个节点的 systemd 日志、端口与证书 principal 判断,先恢复原 quorum 多数派。不要对已有元数据目录再次 format,也不要用新 cluster id 启动空集群接管旧 broker。
NotEnoughReplicas 或分区没有 leader
先描述 topic 的 leader、ISR 与 ELR,再检查磁盘空间、I/O、长暂停和网络。恢复故障副本并等待追平后再开放写入;临时降低 min ISR、改成 acks=1 或开启不干净选举,会把暂时不可写转化为数据丢失风险。
Consumer lag 持续增加
同时比较各分区 end offset、committed offset、最老待处理事件年龄和处理耗时。单分区异常通常是 hot key;所有分区一起增长更可能是数据库、线程池或配额瓶颈。若处理超过 max poll interval,成员会被移出组并反复 rebalance。先缩小批次或把耗时工作移出 poll 线程,再按分区吞吐扩容;实例数超过分区数不会提升传统 group 并行度。
数据库有结果,但消息再次到达
数据库 commit 成功而 offset commit 超时,重复投递是预期恢复路径。检查 event id 是否稳定、Inbox 与业务更新是否同一事务、提交位置是否为当前 offset 加一。保留同一 group 重放,第二次 Inbox 插入应命中唯一键且不重复改变业务;删除 group 或改为 latest 只会掩盖缺口。
Share Group 反复投递
查看成员、record lock、delivery count 与 ack 类型。处理超过锁期限、成员宕机、显式 release 或可重试失败都会再次投递。任务可并行时缩短处理或续锁并设置失败去向;业务依赖同 key 严格顺序时,应改用 Consumer Group。
认证成功但 ACL 拒绝
服务端日志里的 principal 才是授权输入。核对 mTLS 映射或 SCRAM 用户名、topic/group/transactional-id、literal 或 prefixed pattern 和具体操作,再增加最小权限。加入 super users 或开启无 ACL 默认放行,会把局部错误变成全局越权。
记录过大或 producer 缓冲耗尽
RecordTooLarge 要同时比较 producer 最大请求、broker/topic 最大消息与 follower 最大 fetch;只放大一端会让复制继续失败。缓冲耗尽还应查 Broker 延迟、配额、网络和 delivery timeout。大对象通常放对象存储,Kafka 只传引用、摘要和校验信息。
Offset 已被保留策略删除
OFFSET_OUT_OF_RANGE 表明 committed offset 早于 log start offset。先用 kafka-get-offsets.sh --time earliest、--time latest 与 group describe 量化缺口,再从源系统、归档或 DR 恢复。确认业务可重建且具备幂等性后,才预览并执行 reset;直接设 latest 会永久跳过缺失事件。
Compaction 后仍读到旧值
压缩按 segment 后台运行,不保证新值写入后旧值立即消失。检查 key 序列化字节、segment 是否关闭、cleaner backlog 与 tombstone 保留时间。需要当前值的消费者应按 offset 折叠同 key;需要完整历史则不能只依赖 compacted topic。
DR 有数据却不能切换
目标 topic 有数据不代表 checkpoint、事务标记、ACL、schema 与外部数据库都满足恢复点。先围栏源写入,确认 MM2 heartbeat 和复制 lag,确保目标 group 不活跃时再同步位置,并以 event id 对账。无法保证单写与 RPO 时应保持只读,不能用 DNS 切换制造双写。
升级后协议或事务异常
区分二进制版本、metadata feature level、group protocol 和客户端版本。查看 feature、controller、transaction coordinator 与 group 状态;finalization 前可滚回旧二进制,之后不能把 4.3 元数据降级。此时应修复兼容客户端或前滚,而不是在生产元数据目录上尝试降级格式。
权威资料与规范地址
元数据与成员维护
| 查阅内容 | 官方地址 |
|---|---|
| KRaft 运维 | https://kafka.apache.org/43/operations/kraft/ |
