Kafka 运行链:分区日志、副本、位点与消费组
订单事件被追加到 Kafka 后,库存服务可以读取一次,审计服务可以从头读取,新的统计任务也可以从保留的历史位置开始。三个订阅者共用日志,各自保存进度。日志保留多久、谁负责某个分区、哪个位置已经处理完成,决定了这些读取能否持续下去。
从一条记录找到分区日志
Topic、Partition、Record 与 Key
Topic:order-events
├── Partition 0:offset 0 → 1 → 2 → ...
│ ├── leader:broker 1
│ └── follower replicas:broker 2、broker 3
├── Partition 1:自己的 offset 序列与副本
└── Partition 2:自己的 offset 序列与副本
Record:key + value + headers + timestamp
位置:TopicPartition + offset
订阅进度:group.id + TopicPartition → committed offsetTopic 是业务记录的分类名称,Partition 是追加日志和并行读取的基本单位。记录进入一个分区后得到 offset。不同分区的 offset 独立增长,partition 0 / offset 8 与 partition 1 / offset 8 是两个位置。
Key 通常选择订单号、账户号或设备号,使同一业务对象的记录落到同一分区。使用默认分区逻辑时,相同序列化 key 在分区数量不变的条件下得到稳定映射;显式指定 partition、修改分区器或增加分区,都可能改变这一结果。需要按订单依次执行“创建、支付、退款”时,还要保持分区内顺序处理,避免线程池把先读取的记录后完成。Kafka 日志与分区设计
增加分区后,旧记录留在原位置,新记录可能使用新映射。对顺序敏感的 Topic,可以采用新 Topic 切换、停写迁移或业务版本检查,明确新旧流量的衔接方式。单纯调大分区数无法重排历史日志。
客户端如何找到正确的 Broker
bootstrap.servers 提供初始联系地址。生产者取得元数据后,按目标分区找到 leader;消费者也根据元数据连接持有目标分区的 Broker。因此,能连 bootstrap 地址只完成了第一步,元数据中的 advertised.listeners 地址也必须能从客户端网络访问。
实验将所有客户端放进同一个 Docker 网络,Broker 对外通告 kafka1:9092、kafka2:9092、kafka3:9092。这些服务名由 Docker 内部 DNS 解析;宿主直接运行客户端时,需要另外设计可由宿主访问的 listener,而不能把容器服务名照搬到公网客户端。Listener 与 Broker 配置
KRaft controller 管理集群元数据和分区领导权,Broker 保存并提供业务日志。生产部署通常使用三个独立 controller 形成元数据仲裁。下面保留一个 controller、三个 Broker,专门观察数据副本变化;controller 停机将使这个实验失去控制面可用性,不能据此称它为完整高可用集群。
建立可执行的 Java 与 Docker 环境
下载完整实验工程,解压进入 message-event-lab。需要 Linux shell、Docker Engine、Compose v2+、至少约 3 GiB 可用容器内存。由有 Docker 权限的普通宿主用户操作,镜像使用 Apache 官方 apache/kafka:4.3.1,客户端固定 kafka-clients:4.3.1;不使用实验性的 native 镜像。官方 Docker 镜像说明
message-event-lab/
├── pom.xml 共享构建配置,Java 17 编译目标
├── compose.kafka.yaml 1 controller + 3 brokers
└── kafka/
├── pom.xml kafka-clients 4.3.1、SLF4J 2.0.17
└── src/main/java/example/
├── KafkaLabSupport.java 客户端、隔离 Topic 与有限时读取
├── KafkaRuntimeLab.java 位点、事务与 ISR 实验
└── KafkaGroupLab.java 真实订阅与分区重新分配共享工程还包含其他消息实验。这里使用 -pl kafka -am 只构建 Kafka 模块。Compose 的数据目录放在命名卷中,由镜像内的 appuser 使用;Java 客户端和 Maven 使用当前宿主 UID/GID。示例仅在专用网络开放明文端口,不配置生产账号、TLS 或 ACL,不应直接发布到外网。
test "$(id -u)" -ne 0 || { echo '请使用普通宿主用户'; exit 1; }
mkdir -p .m2
docker compose -f compose.kafka.yaml up -d
ready=0
for attempt in $(seq 1 30); do
if docker compose -f compose.kafka.yaml exec -T kafka1 \
/opt/kafka/bin/kafka-topics.sh --bootstrap-server kafka1:9092 --list; then
ready=1
break
fi
sleep 2
done
test "$ready" -eq 1 || { docker compose -f compose.kafka.yaml logs --tail=60; exit 1; }
docker run --rm --user "$(id -u):$(id -g)" \
-e MAVEN_CONFIG=/maven-cache \
-v "$PWD:/workspace" -v "$PWD/.m2:/maven-cache" -w /workspace \
maven:3.9.12-eclipse-temurin-25 \
mvn -B -ntp '-Dmaven.repo.local=/maven-cache' '-Duser.home=/tmp' \
-pl kafka -am clean package dependency:copy-dependencies预期 Maven 输出 BUILD SUCCESS。若网络与元数据已连通但后续建 Topic 提示副本不足,检查三个 Broker 是否全部注册:docker compose -f compose.kafka.yaml ps 与各 Broker 启动日志。一次 --list 成功不会检查所有分区副本健康。
定义后续命令使用的函数:
run_kafka() {
docker run --rm --user "$(id -u):$(id -g)" --network me12-kafka_default \
-v "$PWD:/workspace:ro" -w /workspace eclipse-temurin:25.0.4_7-jdk \
java '-Dorg.slf4j.simpleLogger.defaultLogLevel=warn' \
-cp 'kafka/target/classes:kafka/target/dependency/*' example.KafkaRuntimeLab "$@"
}如果更改了 Compose project name,函数中的网络名也要随之修改。各模式只重建名称以 me12-kafka- 开头的指定实验 Topic;它们的历史消息会被删除。不要将同名 Topic 用于业务数据,也不要并发运行同一个实验入口。
区分拉取位置、已提交位置与业务处理
先跑一次关闭后重放
run_kafka offsets程序发送三条记录,第一次读取后关闭消费者且不提交;第二次重开读取同三条记录,然后提交;第三次重开确认无重复记录:
record: partition=0 offset=0
record: partition=0 offset=1
record: partition=0 offset=2
beforeCommit: records=3 position=3 committed=null
restarted: replayed=3 committedNext=3
afterCommitRestart: records=0 position=3position=3 是当前客户端接下来读取的位置。committed=3 才是消费组保存的恢复起点。第一次关闭之后,客户端内存里的 position 消失,新的客户端因没有提交记录而按 auto.offset.reset=earliest 从可用起点读取。
这里的三个记录恰好连续。实际日志还可能有事务控制记录、压缩删除留下的位置空洞,不能把任意 offset 差值解释成业务消息条数。代码使用 ConsumerRecords.nextOffsets() 收集各分区的下次位置及 leader epoch,再提交已经全部处理完成的批次。KafkaConsumer API 与位点提交
这个实验使用 assign 固定分区,便于单独观察位点。assign 不加入消费者组的自动分配过程;后面的 KafkaGroupLab 使用 subscribe 和持续 poll,实际触发组协调与重平衡。
完整入口如何建立这些判断
公共辅助类负责建立客户端、重置专属 Topic、限时收集记录。这里的 read 只有在拿到准确条数后才返回,不会把超时后空集合误报成成功。
package example;
import java.time.Duration;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Properties;
import java.util.concurrent.TimeUnit;
import org.apache.kafka.clients.admin.Admin;
import org.apache.kafka.clients.admin.NewTopic;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.common.TopicPartition;
final class KafkaLabSupport {
static final String BOOTSTRAP = System.getenv().getOrDefault("KAFKA_BOOTSTRAP", "kafka1:9092,kafka2:9092,kafka3:9092");
static void check(boolean ok, String detail) { if (!ok) throw new IllegalStateException(detail); }
static Properties producerProperties() {
var p = new Properties(); p.put("bootstrap.servers", BOOTSTRAP);
p.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
p.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
p.put("acks", "all"); p.put("enable.idempotence", "true"); p.put("linger.ms", "0");
p.put("request.timeout.ms", "3000"); p.put("delivery.timeout.ms", "8000"); p.put("max.block.ms", "15000");
return p;
}
static KafkaProducer<String,String> producer() { return new KafkaProducer<>(producerProperties()); }
static Properties consumerProperties(String group, String isolation) {
var p = new Properties(); p.put("bootstrap.servers", BOOTSTRAP);
p.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
p.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
p.put("group.id", group); p.put("group.protocol", "classic");
p.put("enable.auto.commit", "false"); p.put("auto.offset.reset", "earliest");
p.put("isolation.level", isolation); p.put("default.api.timeout.ms", "15000");
return p;
}
static KafkaConsumer<String,String> consumer(String group, String isolation) {
return new KafkaConsumer<>(consumerProperties(group, isolation));
}
static Admin admin() { return Admin.create(Map.of("bootstrap.servers",BOOTSTRAP,"default.api.timeout.ms","15000","request.timeout.ms","5000")); }
static void resetGroup(String group) throws Exception {
try (var a = admin()) {
try { a.deleteConsumerGroups(List.of(group)).all().get(15, TimeUnit.SECONDS); }
catch (java.util.concurrent.ExecutionException e) {
if (!(e.getCause() instanceof org.apache.kafka.common.errors.GroupIdNotFoundException)) throw e;
}
}
}
static void resetTopic(String name, int partitions) throws Exception {
try (var a = admin()) {
if (a.listTopics().names().get(15,TimeUnit.SECONDS).contains(name)) {
a.deleteTopics(List.of(name)).all().get(15,TimeUnit.SECONDS);
long deadline=System.nanoTime()+TimeUnit.SECONDS.toNanos(15);
while(a.listTopics().names().get(15,TimeUnit.SECONDS).contains(name)) {
check(System.nanoTime()<deadline,"topic deletion timed out"); Thread.sleep(100);
}
}
a.createTopics(List.of(new NewTopic(name,partitions,(short)3).configs(Map.of("min.insync.replicas","2")))).all().get(15,TimeUnit.SECONDS);
}
}
record Batch(List<ConsumerRecord<String,String>> records, Map<TopicPartition,OffsetAndMetadata> nextOffsets) {}
static Batch read(KafkaConsumer<String,String> c, int count) {
var out = new ArrayList<ConsumerRecord<String,String>>();
var next = new HashMap<TopicPartition,OffsetAndMetadata>();
long deadline=System.nanoTime()+TimeUnit.SECONDS.toNanos(20);
while(out.size()<count && System.nanoTime()<deadline) {
var records=c.poll(Duration.ofMillis(250));
records.forEach(out::add); next.putAll(records.nextOffsets());
}
check(out.size()==count,"expected "+count+" records, received "+out.size());
return new Batch(out,next);
}
static List<String> values(KafkaConsumer<String,String> c, Duration window) {
var result=new ArrayList<String>(); long deadline=System.nanoTime()+window.toNanos();
while(System.nanoTime()<deadline) c.poll(Duration.ofMillis(200)).forEach(r->result.add(r.value()));
return result;
}
}主入口通过模式选择实验,所有成功、负例类型和恢复结果都有断言:
package example;
import java.time.Duration;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import org.apache.kafka.clients.admin.NewTopic;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.errors.NotEnoughReplicasException;
import org.apache.kafka.common.errors.NotEnoughReplicasAfterAppendException;
public final class KafkaRuntimeLab {
public static void main(String[] args) throws Exception {
KafkaLabSupport.check(args.length==1,"mode: offsets|transactions|transaction-pipe|isr-prepare|isr-ok|isr-fail");
switch(args[0]) {
case "offsets" -> offsets();
case "transactions" -> transactions();
case "transaction-pipe" -> transactionPipe();
case "isr-prepare" -> {
try(var a=KafkaLabSupport.admin()) {
String topic="me12-kafka-isr";
if(!a.listTopics().names().get(15,TimeUnit.SECONDS).contains(topic))
a.createTopics(List.of(new NewTopic(topic,Map.of(0,List.of(1,2,3)))
.configs(Map.of("min.insync.replicas","2")))).all().get(15,TimeUnit.SECONDS);
var partition=a.describeTopics(List.of(topic)).allTopicNames().get(15,TimeUnit.SECONDS).get(topic).partitions().get(0);
KafkaLabSupport.check(partition.isr().size()==3,"expected three ISR members");
KafkaLabSupport.check(partition.leader().id()==1,"keep current leader online; expected broker1");
System.out.println("isr: leader=1 replicas=3 minISR=2 isr="+partition.isr().stream().map(n->n.id()).toList());
}
sendIsr(false);
}
case "isr-ok" -> sendIsr(false);
case "isr-fail" -> sendIsr(true);
default -> throw new IllegalArgumentException("unknown mode: "+args[0]);
}
}
static void offsets() throws Exception {
String topic="me12-kafka-offsets",group="me12-offset-group";
KafkaLabSupport.resetTopic(topic,1); var tp=new TopicPartition(topic,0);
KafkaLabSupport.resetGroup(group);
try(var producer=KafkaLabSupport.producer()) {
for(int i=0;i<3;i++) {
var result=producer.send(new ProducerRecord<>(topic,0,"order-42","event-"+i)).get(15,TimeUnit.SECONDS);
System.out.println("record: partition="+result.partition()+" offset="+result.offset());
}
}
try(var c=KafkaLabSupport.consumer(group,"read_uncommitted")) {
c.assign(List.of(tp)); var batch=KafkaLabSupport.read(c,3);
KafkaLabSupport.check(c.committed(Set.of(tp)).get(tp)==null,"unexpected prior commit");
System.out.println("beforeCommit: records="+batch.records().size()+" position="+c.position(tp)+" committed=null");
}
try(var c=KafkaLabSupport.consumer(group,"read_uncommitted")) {
c.assign(List.of(tp)); var batch=KafkaLabSupport.read(c,3);
KafkaLabSupport.check(batch.records().get(0).offset()==0,"uncommitted records did not replay");
c.commitSync(batch.nextOffsets());
KafkaLabSupport.check(c.committed(Set.of(tp)).get(tp).offset()==3,"next offset not committed");
System.out.println("restarted: replayed=3 committedNext=3");
}
try(var c=KafkaLabSupport.consumer(group,"read_uncommitted")) {
c.assign(List.of(tp)); KafkaLabSupport.check(KafkaLabSupport.values(c,Duration.ofSeconds(2)).isEmpty(),"committed records replayed");
System.out.println("afterCommitRestart: records=0 position="+c.position(tp));
}
}
static void transactions() throws Exception {
String topic="me12-kafka-transactions"; KafkaLabSupport.resetTopic(topic,1);
var props=KafkaLabSupport.producerProperties(); props.put("transactional.id","me12-transaction-producer");
try(var p=new KafkaProducer<String,String>(props)) {
p.initTransactions(); p.beginTransaction();
p.send(new ProducerRecord<>(topic,0,"order-42","committed")).get(15,TimeUnit.SECONDS); p.commitTransaction();
p.beginTransaction(); p.send(new ProducerRecord<>(topic,0,"order-42","aborted")).get(15,TimeUnit.SECONDS); p.abortTransaction();
}
var tp=new TopicPartition(topic,0);
for(String isolation:List.of("read_committed","read_uncommitted")) {
try(var c=KafkaLabSupport.consumer("me12-"+isolation,isolation)) {
c.assign(List.of(tp)); c.seekToBeginning(List.of(tp));
var values=KafkaLabSupport.values(c,Duration.ofSeconds(4));
KafkaLabSupport.check(values.equals(isolation.equals("read_committed")?List.of("committed"):List.of("committed","aborted")),"unexpected transaction visibility: "+values);
System.out.println(isolation+": "+values);
}
}
}
static void transactionPipe() throws Exception {
String input="me12-kafka-input",output="me12-kafka-output",group="me12-pipe-group";
KafkaLabSupport.resetTopic(input,1); KafkaLabSupport.resetTopic(output,1);
KafkaLabSupport.resetGroup(group);
try(var p=KafkaLabSupport.producer()) { p.send(new ProducerRecord<>(input,"order-42","input")).get(15,TimeUnit.SECONDS); }
var props=KafkaLabSupport.producerProperties(); props.put("transactional.id","me12-pipe-producer");
var tp=new TopicPartition(input,0);
try(var c=KafkaLabSupport.consumer(group,"read_committed"); var p=new KafkaProducer<String,String>(props)) {
c.subscribe(List.of(input)); var batch=KafkaLabSupport.read(c,1); p.initTransactions();
p.beginTransaction(); p.send(new ProducerRecord<>(output,"order-42","aborted-output")).get(15,TimeUnit.SECONDS);
p.sendOffsetsToTransaction(batch.nextOffsets(),c.groupMetadata()); p.abortTransaction();
KafkaLabSupport.check(c.committed(Set.of(tp)).get(tp)==null,"aborted input offset committed");
System.out.println("aborted: inputCommitted=null currentPosition="+c.position(tp));
c.seek(tp,0); // Aborting the producer does not rewind this consumer.
batch=KafkaLabSupport.read(c,1);
p.beginTransaction(); p.send(new ProducerRecord<>(output,"order-42","committed-output")).get(15,TimeUnit.SECONDS);
p.sendOffsetsToTransaction(batch.nextOffsets(),c.groupMetadata()); p.commitTransaction();
KafkaLabSupport.check(c.committed(Set.of(tp)).get(tp).offset()==1,"input offset missing");
}
try(var c=KafkaLabSupport.consumer("me12-pipe-result","read_committed")) {
var out=new TopicPartition(output,0);c.assign(List.of(out));c.seekToBeginning(List.of(out));
var values=KafkaLabSupport.values(c,Duration.ofSeconds(4));
KafkaLabSupport.check(values.equals(List.of("committed-output")),"transaction output wrong: "+values);
System.out.println("committed: inputNext=1 visibleOutput="+values);
}
}
static void sendIsr(boolean expectFailure) throws Exception {
var props=KafkaLabSupport.producerProperties(); props.put("retries","1");
try(var p=new KafkaProducer<String,String>(props)) {
try {
var result=p.send(new ProducerRecord<>("me12-kafka-isr",0,"probe","probe")).get(15,TimeUnit.SECONDS);
KafkaLabSupport.check(!expectFailure,"write unexpectedly succeeded below minISR");
System.out.println("isrWrite: acknowledged=true offset="+result.offset());
} catch(ExecutionException e) {
Throwable cause=e.getCause();
KafkaLabSupport.check(expectFailure&&(cause instanceof NotEnoughReplicasException||cause instanceof NotEnoughReplicasAfterAppendException),"unexpected send failure: "+cause);
System.out.println("isrWrite: acknowledged=false error="+cause.getClass().getSimpleName());
}
}
}
}offsets() 每次重开都会创建新的消费者,恢复位置取决于 Broker 保存的 group offset 与 auto.offset.reset。前两次读取都得到三条记录,第二次从 offset 0 开始;提交 next offset 后,第三次重开从位置 3 继续。
enable.auto.commit=false 让提交位置显式可见。处理一个 poll 批次时,如果只有前半段成功,应为每个分区计算连续完成的位置,不能直接提交整个批次的 nextOffsets。异步处理尤其需要维护“连续完成前沿”:offset 12 完成而 11 未完成时,恢复起点仍应覆盖 11。
从一次业务写入选择提交时机
| 执行顺序 | 中断窗口 | 恢复办法 |
|---|---|---|
| 先提交 Kafka offset,再写数据库 | offset 已前进,数据库操作尚未完成 | 需要业务对账或重新指定位置;自动重启通常跳过该消息 |
| 数据库事务成功,再提交 offset | 数据已更新,offset 提交未成功 | 重放同一业务事件,由数据库唯一约束或版本检查抑制重复 |
| Kafka 输出与输入 offset 同一 Kafka 事务 | Kafka 事务整体提交或中止 | 下游使用 read_committed;外部数据库仍需另外协调 |
数据库更新和 Kafka 提交涉及两个系统。订单消费可以在一个数据库事务内插入 Inbox 事件 ID 并更新订单,再提交 Kafka 位点;实际 SQL 与重复交付恢复见 投递、幂等与顺序。给所有消息配置同一个 transactional.id 无法把外部数据库加入 Kafka 事务。
副本确认与 Kafka 事务怎样推进
acks、ISR 与 min.insync.replicas
一个分区的 leader 接收追加,follower 从 leader 拉取日志。ISR 记录当前符合同步条件的副本。三副本配置描述有三份分区副本,并不要求任意时刻三份都已追上。
acks=all 让生产者等待同步副本满足确认要求,min.insync.replicas=2 限制允许继续写入的最小 ISR 数量。两者一起使用时,ISR 下降到 1 将拒写。若只设置 acks=1,leader 本地追加后即可响应,客户端不会因此自动获得最小双副本确认。
生产者中的幂等功能通过 producer ID、epoch 与分区序号识别协议重试。应用重新调用一次 send 表达的新记录仍可能被接受;跨进程重建同一业务事件也要保留业务事件 ID。启用幂等还要求 acks=all、允许重试及受约束的在途请求数,冲突配置应在启动时处理。生产者配置
批次影响实际性能。batch.size 限制分区批次大小,linger.ms 给小记录留下聚合窗口,压缩在批次层面节省网络与存储;低流量下过大的 linger 会直接增加等待。调参时同时观察批次填充、发送延迟、缓冲池等待和请求错误,先确认瓶颈在客户端、网络还是 Broker。
实际让 ISR 从三份降到一份
先在三个 Broker 正常时准备专用分区。代码明确将副本放在 1、2、3,首次实验保持 broker 1 为 leader:
run_kafka isr-prepare
docker compose -f compose.kafka.yaml exec -T kafka1 \
/opt/kafka/bin/kafka-topics.sh --bootstrap-server kafka1:9092 \
--describe --topic me12-kafka-isr预期 leader=1、Replicas: 1,2,3、ISR 包含三个成员。若之前演练改变了 leader,先恢复所有 Broker 并确认副本,再根据真实 leader 重新规划停止对象;不要继续执行会停掉 leader 的故障序列。
每次停止后,用下面的函数等待实际 ISR 数量稳定,再运行写入:
wait_isr() {
wanted="$1"
for attempt in $(seq 1 30); do
report=$(docker compose -f compose.kafka.yaml exec -T kafka1 \
/opt/kafka/bin/kafka-topics.sh --bootstrap-server kafka1:9092 \
--describe --topic me12-kafka-isr) || return 1
printf '%s\n' "$report"
count=$(printf '%s\n' "$report" | awk '{for(i=1;i<=NF;i++) if($i=="Isr:"){n=split($(i+1),a,","); print n}}')
test "$count" = "$wanted" && return 0
sleep 1
done
return 1
}
docker compose -f compose.kafka.yaml stop kafka3
wait_isr 2 || exit 1
run_kafka isr-ok
docker compose -f compose.kafka.yaml stop kafka2
wait_isr 1 || exit 1
run_kafka isr-fail
docker compose -f compose.kafka.yaml up -d kafka2 kafka3
wait_isr 3 || exit 1
run_kafka isr-ok中间两次写入的关键结果分别是:
isrWrite: acknowledged=true offset=1
isrWrite: acknowledged=false error=NotEnoughReplicasExceptionoffset 随实验重复次数变化,不应固定比对;必须比对的是确认结果与异常类型。程序只接受副本不足相关异常作为负例成功,连接失败或任意超时都会使程序失败。恢复到三个 ISR 后再次得到确认,才完成这轮演练。实验中途退出也应先恢复 kafka2、kafka3。
Kafka 4.3 的新集群还启用 ELR(Eligible Leader Replicas)。本次 ISR 只剩 1 时,描述输出可出现 Elr: 2:controller 保留了额外的安全选主信息。讨论选主时需同时看 ISR、ELR 和 feature 状态,不能沿用“所有非 ISR 副本都只能不安全选主”的旧分类。ELR 选主规则
事务提交、中止与读隔离
run_kafka transactions
run_kafka transaction-pipe第一个入口分别提交与中止一条记录;第二个入口将输出记录和输入组位点纳入同一事务:
read_committed: [committed]
read_uncommitted: [committed, aborted]
aborted: inputCommitted=null currentPosition=1
committed: inputNext=1 visibleOutput=[committed-output]事务生产者先 initTransactions,随后在 beginTransaction 与提交/中止之间发送记录。Kafka 通过事务协调器维护事务状态,在参与分区写入相应控制标记;read_committed 的读取还受最后稳定位置 LSO 限制,避免越过未完成事务暴露尚未确定的结果。事务协议
输入—处理—输出拓扑中,sendOffsetsToTransaction(batch.nextOffsets(), consumer.groupMetadata()) 将输入进度绑定到同一个 Kafka 事务。第二个实验先中止,观察提交位置仍为空,再主动 seek 回输入起点重新处理。生产者的 abortTransaction 不会修改消费者对象已经前进的 position。KafkaProducer 事务 API
transactional.id 应能识别同一逻辑生产任务的重启,同时避免两个合法活动实例争用同一个 ID。旧实例遭到 fencing 后必须退出或重建相应任务;把所有异常统一重试,会把所有权冲突延长成持续故障。
消费组如何接管分区并完成恢复
运行真实的加入、撤销与重新分配
普通消费者组内,一个分区在稳定分配时交给一个成员。多个组则可以独立订阅同一日志,各自保存进度。Kafka 4.3 同时提供 Share Consumer API;下面使用普通 KafkaConsumer,并显式设置 group.protocol=classic,不把普通组的分区分配规则扩展到 Share Groups。消费者协议与配置
docker run --rm --user "$(id -u):$(id -g)" --network me12-kafka_default \
-v "$PWD:/workspace:ro" -w /workspace eclipse-temurin:25.0.4_7-jdk \
java '-Dorg.slf4j.simpleLogger.defaultLogLevel=warn' \
-cp 'kafka/target/classes:kafka/target/dependency/*' example.KafkaGroupLab 2入口先向一个三分区 Topic 发布六条 key 相同的记录,断言实际元数据只出现一个分区;随后逐步改变活跃消费者数量:
key: order-42 records=6 partitions=[0]
group: members=1 assigned=[3] union=3 disjoint=true
group: members=2 assigned=[1, 2] union=3 disjoint=true
group: members=1 assigned=[3] union=3 disjoint=true具体分区 ID 随分区算法和 key 改动而变化。相同 key 应落入一个目标分区,稳定分配后的所有权集合应无交叉且覆盖三个分区。程序通过 ConsumerRebalanceListener 更新集合;关闭第二个消费者后,第一个消费者会在重新分配完成后接管全部分区。
每个消费者始终由自己的线程创建、poll 和关闭。KafkaConsumer 通常不支持多线程并发调用;业务线程可以处理已交出的数据,但不能顺便跨线程调用它的 poll/commit/seek。安全停机先停止接收新工作,再收束在途处理,按分区提交连续完成的位置,最后关闭消费者。
长任务、再平衡与旧成员
session.timeout.ms 判断成员是否仍保持会话,max.poll.interval.ms 约束应用两次 poll 之间允许的最大间隔。心跳仍在并不代表应用可以无限时间不 poll。一次批次包含大量慢 SQL 时,应减小 max.poll.records,限制任务队列,或暂停对应分区而继续维持 poll 循环。
onPartitionsRevoked 是撤销之前整理进度的机会,适合提交已连续完成的工作;onPartitionsLost 表示分区可能已经被他人接管,处理必须更保守。协调器会检查提交请求中的组成员代际,但数据库并不认识 Kafka 的组代际;旧线程已经发出的 SQL 仍要通过唯一约束、版本条件或业务锁限制影响。
classic eager/cooperative assignor 与新 consumer group 协议的再平衡流程、配置归属不同。切换协议前应检查客户端版本、Broker 配置、回调行为与滚动升级办法,不把一次“减少停顿”参数变更当作完整迁移。消费者重平衡协议
日志保留决定还能从哪里恢复
delete 策略按时间、空间及日志段条件回收旧数据。压缩策略按照 key 保留较新的值,tombstone 表达删除;压缩后 offset 可以有空洞,且尚未被清理的旧值仍可能存在。事件审计需要保留每次变化,状态投影则可能适合 compaction,应按记录含义选择 Topic 策略。Topic 保留与清理配置
消费者停机超过保留窗口后,已提交位置可能早于日志起点。auto.offset.reset 此时决定自动从 earliest/latest 开始还是报错;对于金融流水、订单状态等不可静默跳过的输入,可以选择显式失败并从快照与可用日志恢复。
同样重要的是处理速度。若持续写入速度大于消费净处理速度,日志保留会一边推进一边淘汰尚未处理的数据。先测出每个分区的最老积压年龄和完成速率,再评估扩分区、扩消费者或降低下游开销。
故障信息对应到下一项检查
| 现象 | 优先检查 | 下一步 |
|---|---|---|
| bootstrap 可达,发送仍持续找不到 leader | 元数据中的 advertised listener、DNS、端口与分区 leader | 从客户端所在网络验证每个通告地址 |
| NotEnoughReplicas | Topic 实际 min ISR、ISR 数量与落后副本日志 | 恢复副本,不先降低门槛掩盖故障 |
| 消费位置增长但提交位置不动 | 手工提交分支、异常处理和异步完成前沿 | 对照业务事务与提交调用结果 |
| 重平衡频繁 | poll 间隔、进程退出、组协议与订阅变化 | 限制单批工作时间,查撤销/丢失回调 |
| read_committed 暂时读不到较新记录 | 未完成事务、LSO 与事务生产者状态 | 查悬挂事务及事务超时,不强行绕过隔离 |
| 消费者增加后仍有热点 | 分区数、key 分布、单分区处理时长 | 先处理热点或业务拆分,再增加并发 |
停止实验用 docker compose -f compose.kafka.yaml stop 保留日志。再次使用时启动所有四个服务并确认副本恢复。仅在确认 me12-kafka 全部为可丢弃实验数据后运行以下命令,删除本项目容器、网络和四个命名卷:
docker compose -f compose.kafka.yaml down --volumes