欢迎光临
我们一直在努力

Kafka简介及其核心概念

简介

Kafka 是一个分布式的、高吞吐量的、可持久化的消息队列系统,最初由 LinkedIn 开发,现在属于 Apache 顶级项目。

Producer(生产者) → Kafka Cluster(集群) → Consumer(消费者)

Broker(服务器) + Zookeeper(协调服务)

  • Kafka 的核心优势:
    • ✅ 高吞吐:适合大数据量场景
    • ✅ 持久化:消息不丢失,可重放
    • ✅ 分布式:高可用,易扩展
    • ✅ 生态丰富:与各种流处理框架集成
  • 适用场景:
    • 日志收集和分析
    • 实时数据管道
    • 事件驱动架构
    • 流式处理
  • 选择建议:
    • 需要极高吞吐量 → Kafka
    • 交易场景,需要低延迟 → RocketMQ
    • 日志处理,大数据场景 → Kafka

1. 基本组成

  • Producer(生产者):产生和发送消息的客户端
  • Consumer(消费者):接收和处理消息的客户端
  • Broker(代理):Kafka 服务器实例
  • Cluster(集群):多个 Broker 组成的集合

2. Topic(主题)

  • 消息的逻辑分类,类似于数据库中的表
  • 每个 Topic 都有一个唯一的名称
  • 生产者向特定 Topic 发送消息,消费者从 Topic 消费消息

3. Partition(分区)

  • 每个 Topic 可以被分成多个分区
  • 分区是 Kafka 并行处理的基本单位
  • 分区内的消息是有序的,但不同分区之间的顺序不保证
  • 分区允许 Topic 水平扩展到多个 Broker

4. Offset(偏移量)

  • 分区中每条消息的唯一标识
  • 在分区内严格递增且连续
  • 消费者通过维护偏移量来跟踪消费进度

5. Replication(副本)

  • 每个分区可以有多个副本
  • 副本分为 Leader 和 Follower
  • Leader 处理所有读写请求,Follower 同步数据
  • 提供数据冗余和高可用性

6. Consumer Group(消费者组)

  • 由多个消费者实例组成的逻辑组
  • 组内的消费者共同消费一个 Topic
  • 每个分区只能被组内的一个消费者消费
  • 实现负载均衡和水平扩展

下载

官网:https://kafka.apache.org/downloads

# 解压
tar -xvf kafka_2.13-3.7.0.tgz -C /usr/local/

注意:本地环境必须安装java8+;

Apache Kafka可以使用Zookeeper或KRaft启动,但只能用其中一种方式,不能同时使用;

KRaft:Apache Kafka的内置共识机制,用于取代Zookeeper

使用zookeeper启动

cd kafka_2.13-3.7.0/
# 启动zookeeper
bin/zookeeper-server-start.sh config/zookeeper.properties &
# 启动kafka
bin/kafka-server-start.sh config/server.properties &
# 关闭Kafka
bin/kafka-server-stop.sh
# 关闭zookeeper
bin/zookeeper-server-stop.sh
# 查看占用的网络端口
netstat -nlpt

使用KRaft启动

cd /usr/local/kafka_2.13-3.7.0/

# 生成Cluster Id(集群UUID) –可以自己指定id,这里直接用命令生成id
bin/kafka-storage.sh random-uuid
# 生成的uuid为SWpip9vUTOuBz0WkIsucTw

# 格式化 Kafka KRaft 模式的存储目录
bin/kafka-storage.sh format -t SWpip9vUTOuBz0WkIsucTw -c config/kraft/server.properties
#生成的集群配置在:dirs={/tmp/kraft-combined-logs: EMPTY} ,要想重新指定需要删除或修改该目录下的文件

# 启动Kafka
bin/kafka-server-start.sh config/kraft/server.properties &

# 关闭Kafka
bin/kafka-server-stop.sh

# 查看占用的网络端口
netstat -nlpt

使用Docker启动

👇总是会出问题,建议在网上找centos7安装docker使用kafka的最新教程。

#检查是否安装过docker
yum list installed | grep docker

#卸载docker
yum remover 名称 -y

#安装docker。注意centos7中的 yum install 会出问题,docker也不维护centos7,下面命令会报错,最好上网重新找安装教程。
yum install yum-utils -y #安装yum工具包
yum-config-manager –add-repo https://download.docker.com/linux/centos/docker-ce.repo #添加 Docker 官方仓库
#安装 Docker 组件
yum install docker-ce docker-ce-cli containerd.io docker-buildx-plugin docker-compose-plugin -y
docker -v #查看是否安装成功

# 1. 安装 yum-utils(如果网络有问题,先配置好仓库)yum install
yum install yum-utils -y –nogpgcheck

# 2. 手动下载并配置 Docker 仓库
curl -o /etc/yum.repos.d/docker-ce.repo https://download.docker.com/linux/centos/docker-ce.repo

# 3. 修改仓库配置,使用阿里云镜像(可选,如果官方源慢)
sed -i 's#download.docker.com#mirrors.aliyun.com/docker-ce#g' /etc/yum.repos.d/docker-ce.repo

# 4. 安装指定版本的 Docker(兼容 CentOS 7)
yum install docker-ce-20.10.24 docker-ce-cli-20.10.24 containerd.io-1.6.28 -y –nogpgcheck

#启动docker
systemctl start docker 或者 service docker start #下面同理,都有两种写法,为了方便下面只写第一种
systemctl stop docker #停止docker
systemctl restart docker # 重启
systemctl status docker # 检查docker状态

#拉取镜像
docker pull apache/kafka:3.7.0
#启动kafka
docker run -p 9092:9092 apache/kafka:3.7.0
#查看docker镜像
docker images
#移除镜像
docker rmi apache/kafka:3.7.0

集成boot

配置

spring:
kafka:
# ==================== 基础连接配置 ====================
bootstrap-servers: 192.168.152.129:9092 # Kafka服务器地址,多个用逗号分隔

# ==================== 生产者配置 ====================
producer:
# 序列化配置
key-serializer: org.apache.kafka.common.serialization.StringSerializer # Key序列化器
value-serializer: org.springframework.kafka.support.serializer.JsonSerializer # Value序列化器

# 可靠性配置
acks: all # 消息确认机制
# 可选值:
# 0 – 不等待确认(性能最高,可能丢失消息)
# 1 – 只等待Leader确认(平衡性能与可靠性)
# all/-1 – 等待所有副本确认(最可靠,性能最低)
retries: 3 # 失败重试次数(0-2147483647)
retry-backoff-ms: 1000 # 重试间隔时间(毫秒)
enable-idempotence: true # 启用幂等性,防止重复消息(true/false)

# 性能配置
batch-size: 16384 # 批量发送大小(字节),0表示禁用批量
linger-ms: 10 # 批量发送等待时间(毫秒)
# 可选值:
# 0 – 立即发送
# >0 – 等待指定时间或达到batch-size后发送
buffer-memory: 33554432 # 生产者缓冲区大小(字节),默认32MB
compression-type: snappy # 压缩算法
# 可选值:
# none – 不压缩
# gzip – 高压缩比,CPU开销大
# snappy – 平衡压缩比和性能(推荐)
# lz4 – 快速压缩,低CPU开销
# zstd – 新一代高压缩比算法
max-block-ms: 60000 # 发送阻塞超时时间(毫秒)

# 事务配置(可选)
transaction-id-prefix: tx # 事务ID前缀,设置后开启事务支持

# ==================== 消费者配置 ====================
consumer:
# 基础配置
group-id: ${spring.application.name}group # 消费组ID,同一组内消费者负载均衡
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer # Key反序列化器
value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer # Value反序列化器
auto-offset-reset: earliest # 当没有初始偏移量或偏移量失效时的策略
# 可选值:
# earliest – 从最早的消息开始消费
# latest – 从最新的消息开始消费(默认)
# none – 没有偏移量时抛出异常

# 提交配置
enable-auto-commit: false # 是否自动提交偏移量
# 可选值:
# true – 自动提交(简单但不保证可靠性)
# false – 手动提交(推荐,保证精确一次处理)
auto-commit-interval-ms: 1000 # 自动提交间隔(毫秒),enable-auto-commit=true时生效

# 性能配置
max-poll-records: 500 # 单次poll()调用返回的最大消息数(1-50000)
max-poll-interval-ms: 300000 # 消费者组协调超时时间(毫秒),默认5分钟
# 超过此时间未poll()会被认为消费者死亡,触发重平衡
fetch-min-size: 1 # 服务器应返回的最小数据量(字节)
fetch-max-wait-ms: 500 # 如果没有足够数据时的最大等待时间(毫秒)
fetch-max-bytes: 52428800 # 单次fetch请求返回的最大数据量(字节),默认50MB
max-partition-fetch-bytes: 1048576 # 每个分区返回的最大数据量(字节),默认1MB

# 安全配置
isolation-level: read_committed # 事务消息的隔离级别
# 可选值:
# read_uncommitted – 读取所有消息(包括未提交的事务消息)
# read_committed – 只读取已提交的事务消息(推荐)

# JSON 反序列化配置
properties:
spring.json.trusted.packages: "*" # 信任的包名,*表示信任所有包
spring.json.use.type.headers: false # 是否使用类型头信息(true/false)

# ==================== 监听器配置 ====================
listener:
# 确认模式配置
ack-mode: manual_immediate # 偏移量确认模式
# 可选值:
# RECORD – 每条消息处理后立即提交
# BATCH – 批量消息处理完后提交(默认)
# TIME – 按时间间隔提交
# COUNT – 按处理数量提交
# COUNT_TIME – 按数量或时间先到者提交
# MANUAL – 手动调用Acknowledgment.acknowledge()提交
# MANUAL_IMMEDIATE – 手动立即提交(推荐)

# 并发配置
concurrency: 3 # 监听器并发数,建议与分区数一致
poll-timeout: 5000 # poll()方法超时时间(毫秒)
type: single # 监听器类型
# 可选值:
# single – 单条消息处理(默认)
# batch – 批量消息处理

# 容错配置
ack-on-error: false # 处理出错时是否提交偏移量
# 可选值:
# true – 出错也提交(可能丢失消息)
# false – 出错不提交(消息会重试,推荐)
retry-interval: 1000 # 重试间隔时间(毫秒)
max-attempts: 3 # 最大重试次数(1-100)

# ==================== 模板配置 ====================
template:
default-topic: defaulttopic # 默认主题名称

# ==================== 高级属性配置 ====================
properties:
# 连接配置
connections.max.idle.ms: 540000 # 连接最大空闲时间(毫秒),默认9分钟
request.timeout.ms: 30000 # 请求超时时间(毫秒),默认30秒
metadata.max.age.ms: 300000 # 强制刷新元数据的时间间隔(毫秒),默认5分钟

# 网络配置
send.buffer.bytes: 131072 # TCP发送缓冲区大小(字节),默认128KB
receive.buffer.bytes: 32768 # TCP接收缓冲区大小(字节),默认32KB

# 安全配置(如果需要SSL/SASL)
# security.protocol: SSL # 安全协议
# 可选值: PLAINTEXT, SSL, SASL_PLAINTEXT, SASL_SSL
# ssl.truststore.location: /path/to/truststore # 信任证书路径
# ssl.truststore.password: password # 信任证书密码
# sasl.mechanism: PLAIN # SASL机制
# 可选值: PLAIN, SCRAM-SHA-256, SCRAM-SHA-512, OAUTHBEARER

# 监控配置
metric.reporters: com.example.MyMetricsReporter # 指标报告器
metrics.num.samples: 2 # 维护的指标样本数量
metrics.sample.window.ms: 30000 # 指标样本时间窗口(毫秒)

# ==================== 应用自定义配置 ====================
app:
kafka:
# 主题配置
topics:
order-topic: "orders" # 订单主题
payment-topic: "payments" # 支付主题
notification-topic: "notifications" # 通知主题

# 消费者组配置
consumer-groups:
order-group: "order-service-group" # 订单服务消费者组
payment-group: "payment-service-group" # 支付服务消费者组
notification-group: "notification-service-group" # 通知服务消费者组

# 重试配置
retry:
max-attempts: 3 # 最大重试次数
backoff-initial-interval: 1000 # 初始重试间隔(毫秒)
backoff-multiplier: 2.0 # 重试间隔乘数

生产者

package org.ldustu.Producer;

import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.header.Headers;
import org.apache.kafka.common.header.internals.RecordHeader;
import org.apache.kafka.common.header.internals.RecordHeaders;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.kafka.support.SendResult;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component;

import java.nio.charset.StandardCharsets;
import java.util.HashMap;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;

/**
* Kafka 消息生产者
* 功能:提供多种消息发送方式,包含同步/异步发送、消息追踪、异常处理等
*
* 最佳实践:
* 1. 使用唯一消息ID便于追踪
* 2. 添加消息头记录元数据
* 3. 合理处理同步/异步发送
* 4. 完善的异常处理和重试机制
*/

@Component
public class EventProducer {

private static final Logger logger = LoggerFactory.getLogger(EventProducer.class);

@Autowired
private KafkaTemplate<String, Object> kafkaTemplate;

/**
* 综合演示方法 – 展示所有发送方式
*/

public void sendAllMessageTypes() {
logger.info("🚀 开始演示所有消息发送方式…");

// 1. 基础发送
sendBasicMessage();

// 2. 通过Message对象发送
sendMessageWithMessageBuilder();

// 3. 通过Headers发送消息
sendMessageWithHeaders();

// 4. 通过多个参数发送消息
sendMessageWithMultipleParams();

// 5. 通过sendDefault发送消息
sendDefaultMessage();

// 6. 同步发送消息
sendSyncMessage();

// 7. 异步发送消息
sendAsyncMessage();

// 8. 发送对象类型消息
sendObjectMessage();

// 9. 带事务的消息发送
sendTransactionalMessage();

logger.info("✅ 所有消息发送方式演示完成");
}

/**
* 1. 基础消息发送 – 最简单的方式
* 适用场景:快速发送简单消息
*/

public void sendBasicMessage() {
String messageId = generateMessageId();
String messageContent = "基础消息 – " + System.currentTimeMillis();

logger.info("1️⃣ 发送基础消息 – ID: {}, 内容: {}", messageId, messageContent);

// 添加追踪头信息
Map<String, Object> headers = new HashMap<>();
headers.put("messageId", messageId);
headers.put("messageType", "BASIC");
headers.put("timestamp", System.currentTimeMillis());

kafkaTemplate.send("topic1", messageContent);

logger.debug("📤 基础消息已发送 – ID: {}", messageId);
}

/**
* 2. 通过MessageBuilder构建消息 – 推荐方式
* 优点:可以灵活设置消息头和属性
*/

public void sendMessageWithMessageBuilder() {
String messageId = generateMessageId();
String messageContent = "MessageBuilder消息 – " + System.currentTimeMillis();

logger.info("2️⃣ 发送MessageBuilder消息 – ID: {}, 内容: {}", messageId, messageContent);

Message<String> message = MessageBuilder
.withPayload(messageContent)
.setHeader(KafkaHeaders.TOPIC, "topic1")
.setHeader(KafkaHeaders.KEY, "msg-key-" + messageId)
.setHeader("messageId", messageId)
.setHeader("messageType", "MESSAGE_BUILDER")
.setHeader("timestamp", System.currentTimeMillis())
.setHeader("source", "EventProducer")
.build();

kafkaTemplate.send(message);
logger.debug("📤 MessageBuilder消息已发送 – ID: {}", messageId);
}

/**
* 3. 通过Headers发送消息 – 最灵活的方式
* 适用场景:需要丰富元数据的消息
*/

public void sendMessageWithHeaders() {
String messageId = generateMessageId();
String messageContent = "Headers消息 – " + System.currentTimeMillis();

logger.info("3️⃣ 发送Headers消息 – ID: {}, 内容: {}", messageId, messageContent);

// 创建自定义头部信息
Headers headers = new RecordHeaders();
headers.add(new RecordHeader("messageId", messageId.getBytes(StandardCharsets.UTF_8)));
headers.add(new RecordHeader("messageType", "HEADERS".getBytes(StandardCharsets.UTF_8)));
headers.add(new RecordHeader("timestamp", String.valueOf(System.currentTimeMillis()).getBytes(StandardCharsets.UTF_8)));
headers.add(new RecordHeader("version", "1.0".getBytes(StandardCharsets.UTF_8)));
headers.add(new RecordHeader("businessType", "ORDER".getBytes(StandardCharsets.UTF_8)));

ProducerRecord<String, Object> record = new ProducerRecord<>(
"topic1", // topic
0, // partition (明确指定分区)
System.currentTimeMillis(), // timestamp
"header-key-" + messageId, // key
messageContent, // value
headers // headers
);

kafkaTemplate.send(record);
logger.debug("📤 Headers消息已发送 – ID: {}", messageId);
}

/**
* 4. 通过多个参数发送消息 – 精确控制
* 适用场景:需要精确控制分区、时间戳等参数
*/

public void sendMessageWithMultipleParams() {
String messageId = generateMessageId();
String messageContent = "多参数消息 – " + System.currentTimeMillis();

logger.info("4️⃣ 发送多参数消息 – ID: {}, 内容: {}", messageId, messageContent);

kafkaTemplate.send(
"topic1", // topic
1, // partition (指定分区1)
System.currentTimeMillis(), // timestamp
"multi-param-key-" + messageId, // key
messageContent // value
);

logger.debug("📤 多参数消息已发送 – ID: {}", messageId);
}

/**
* 5. 通过sendDefault发送消息 – 需要配置默认topic
* 注意:需要在配置文件中配置 spring.kafka.template.default-topic
*/

public void sendDefaultMessage() {
String messageId = generateMessageId();
String messageContent = "sendDefault消息 – " + System.currentTimeMillis();

logger.info("5️⃣ 发送sendDefault消息 – ID: {}, 内容: {}", messageId, messageContent);

kafkaTemplate.sendDefault(
0, // partition
System.currentTimeMillis(), // timestamp
"default-key-" + messageId, // key
messageContent // value
);

logger.debug("📤 sendDefault消息已发送 – ID: {}", messageId);
}

/**
* 6. 同步发送消息 – 阻塞等待结果
* 适用场景:需要确保消息发送成功的业务
* 优点:可靠性高,可以立即知道发送结果
* 缺点:性能较低,会阻塞线程
*/

public SendResult<String, Object> sendSyncMessage() {
String messageId = generateMessageId();
String messageContent = "同步消息 – " + System.currentTimeMillis();

logger.info("6️⃣ 发送同步消息 – ID: {}, 内容: {}", messageId, messageContent);

try {
CompletableFuture<SendResult<String, Object>> future = kafkaTemplate.send(
"topic1",
"sync-key-" + messageId,
messageContent
);

// 设置超时时间,避免无限等待
SendResult<String, Object> result = future.get(5, TimeUnit.SECONDS);

if (result.getRecordMetadata() != null) {
logger.info("✅ 同步消息发送成功 – ID: {}, 元数据: {}",
messageId, formatRecordMetadata(result.getRecordMetadata()));
}

return result;

} catch (InterruptedException e) {
Thread.currentThread().interrupt();
logger.error("❌ 同步消息发送被中断 – ID: {}", messageId, e);
throw new RuntimeException("消息发送被中断", e);
} catch (ExecutionException e) {
logger.error("❌ 同步消息发送执行失败 – ID: {}", messageId, e);
throw new RuntimeException("消息发送执行失败", e);
} catch (TimeoutException e) {
logger.error("❌ 同步消息发送超时 – ID: {}", messageId, e);
throw new RuntimeException("消息发送超时", e);
}
}

/**
* 7. 异步发送消息 – 非阻塞,通过回调处理结果
* 适用场景:高性能要求,不要求立即知道结果的业务
* 优点:性能高,不阻塞线程
* 缺点:不能立即知道发送结果
*/

public void sendAsyncMessage() {
String messageId = generateMessageId();
String messageContent = "异步消息 – " + System.currentTimeMillis();

logger.info("7️⃣ 发送异步消息 – ID: {}, 内容: {}", messageId, messageContent);

CompletableFuture<SendResult<String, Object>> future = kafkaTemplate.send(
"topic1",
"async-key-" + messageId,
messageContent
);

future.thenAccept(result -> {
if(result.getRecordMetadata()!=null) {
System.out.println("非阻塞等待发送消息成功:" + result.getRecordMetadata().toString());
}
}).exceptionally(t -> {
//失败处理
t.printStackTrace();
return null;
});

logger.debug("📤 异步消息已提交 – ID: {}", messageId);
}

/**
* 8. 发送对象类型消息 – 需要配置Json序列化器
* 配置:spring.kafka.producer.value-serializer=org.springframework.kafka.support.serializer.JsonSerializer
*/

public void sendObjectMessage() {
String messageId = generateMessageId();

logger.info("8️⃣ 发送对象消息 – ID: {}", messageId);

Map<String, Object> businessData = new HashMap<>();
businessData.put("orderId", messageId);
businessData.put("userId", "user123");
businessData.put("amount", 199.99);
businessData.put("productName", "测试商品");
businessData.put("timestamp", System.currentTimeMillis());
businessData.put("status", "CREATED");

// 添加业务相关的头信息
Message<Map<String, Object>> message = MessageBuilder
.withPayload(businessData)
.setHeader(KafkaHeaders.TOPIC, "topic1")
.setHeader(KafkaHeaders.KEY, "order-" + messageId)
.setHeader("messageId", messageId)
.setHeader("messageType", "BUSINESS_OBJECT")
.setHeader("businessType", "ORDER_CREATED")
.setHeader("timestamp", System.currentTimeMillis())
.build();

kafkaTemplate.send(message);
logger.debug("📤 对象消息已发送 – ID: {}", messageId);
}

/**
* 9. 带事务的消息发送
* 配置:spring.kafka.producer.transaction-id-prefix=tx-
*/

public void sendTransactionalMessage() {
String messageId = generateMessageId();
String messageContent = "事务消息 – " + System.currentTimeMillis();

logger.info("9️⃣ 发送事务消息 – ID: {}, 内容: {}", messageId, messageContent);

kafkaTemplate.executeInTransaction(operations -> {
// 在事务中发送多条消息
operations.send("topic1", "tx-key-1-" + messageId, messageContent + "-1");
operations.send("topic1", "tx-key-2-" + messageId, messageContent + "-2");
operations.send("topic1", "tx-key-3-" + messageId, messageContent + "-3");

// 这里可以添加数据库操作等其他事务性操作

logger.debug("📤 事务消息已发送 – ID: {}", messageId);
return null;
});
}

/**
* ==================== 工具方法 ====================
*/

/**
* 生成唯一消息ID
*/

private String generateMessageId() {
return UUID.randomUUID().toString().replace("-", "").substring(0, 16);
}

/**
* 格式化记录元数据,便于日志记录
*/

private String formatRecordMetadata(org.apache.kafka.clients.producer.RecordMetadata metadata) {
return String.format("topic=%s, partition=%d, offset=%d, timestamp=%d",
metadata.topic(),
metadata.partition(),
metadata.offset(),
metadata.timestamp());
}

/**
* 处理发送成功逻辑
*/

private void handleSendSuccess(String messageId, SendResult<String, Object> result) {
// 这里可以添加成功后的业务逻辑
// 例如:更新消息状态、发送通知、记录审计日志等
logger.debug("🎯 处理消息发送成功 – ID: {}", messageId);

// 示例:记录消息发送成功审计日志
// auditService.recordMessageAudit(messageId, "SEND_SUCCESS", result);
}

/**
* 处理发送失败逻辑
*/

private void handleSendFailure(String messageId, Throwable ex) {
// 这里可以添加失败后的处理逻辑
// 例如:记录失败日志、发送告警、将消息存入死信队列等
logger.error("💀 处理消息发送失败 – ID: {}", messageId);

// 示例:记录失败消息到数据库,便于后续处理
// failedMessageService.saveFailedMessage(messageId, ex.getMessage());
}

/**
* 消息发送重试逻辑
*/

private void retrySendMessage(String messageId, String messageContent) {
logger.warn("🔄 尝试重试发送消息 – ID: {}", messageId);

// 这里可以实现重试逻辑
// 例如:使用指数退避策略进行重试
// retryTemplate.execute(context -> {
// kafkaTemplate.send("topic1", "retry-key-" + messageId, messageContent);
// return null;
// });
}

/**
* 发送带重试的消息
*/

public void sendMessageWithRetry(String topic, String key, Object message, int maxRetries) {
String messageId = generateMessageId();
int retryCount = 0;

while (retryCount < maxRetries) {
try {
SendResult<String, Object> result = kafkaTemplate.send(topic, key, message).get(10, TimeUnit.SECONDS);
logger.info("✅ 消息发送成功(重试{}次) – ID: {}", retryCount, messageId);
return;
} catch (Exception e) {
retryCount++;
logger.warn("⚠️ 消息发送失败,进行第{}次重试 – ID: {}", retryCount, messageId, e);

if (retryCount >= maxRetries) {
logger.error("❌ 消息发送失败,已达最大重试次数 – ID: {}", messageId, e);
throw new RuntimeException("消息发送失败,重试次数耗尽", e);
}

// 指数退避
try {
Thread.sleep(1000L * (long) Math.pow(2, retryCount));
} catch (InterruptedException ie) {
Thread.currentThread().interrupt();
throw new RuntimeException("重试被中断", ie);
}
}
}
}
}

消费者

package org.ldustu.Consumer;

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.annotation.PartitionOffset;
import org.springframework.kafka.annotation.TopicPartition;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.messaging.handler.annotation.Payload;
import org.springframework.messaging.handler.annotation.SendTo;
import org.springframework.stereotype.Component;

import java.util.List;

/**
* Kafka 消息消费者示例
* 功能:多种消息消费模式,包含单条消费、批量消费、消息转发等
*/

@Component
public class EventConsumer {

/**
* 1. 基础消息消费 – 单条消息处理
* 使用场景:普通消息处理,需要获取消息元信息
*
* @param message 消息内容(使用 @Payload 获取消息体)
* @param topic 消息主题(使用 @Header 获取消息头)
* @param partition 消息分区
* @param record 完整的 ConsumerRecord 对象
* @param ack 确认对象(用于手动提交偏移量)
*/

@KafkaListener(topics = {"topic1"}, groupId = "myConsumerGroup")
public void listenSingleMessage(
@Payload String message,
@Header(KafkaHeaders.RECEIVED_TOPIC) String topic,
@Header(KafkaHeaders.RECEIVED_PARTITION) Integer partition,
@Header(KafkaHeaders.RECEIVED_KEY) String key,
@Header(KafkaHeaders.RECEIVED_TIMESTAMP) Long timestamp,
ConsumerRecord<String, String> record,
Acknowledgment ack) {

try {
// 业务处理逻辑
System.out.println("📥 收到消息: " +
"主题=" + topic +
", 分区=" + partition +
", 键=" + key +
", 时间戳=" + timestamp +
", 消息体=" + message);

// 打印完整消息记录
System.out.println("📋 消息详情: " + record.toString());

// 业务处理(这里可以调用你的业务服务)
processBusinessLogic(message);

// 手动确认消息(配置 ack-mode: manual_immediate 时使用)
ack.acknowledge();
System.out.println("✅ 消息处理完成并确认");

} catch (Exception e) {
System.err.println("❌ 消息处理失败: " + e.getMessage());
// 不调用 ack.acknowledge(),消息会重新投递
// 可以根据异常类型决定是否重试
}
}

/**
* 2. 批量消息消费
* 使用场景:高性能处理,一次处理多条消息
* 配置:spring.kafka.listener.type=batch
*
* @param records 消息记录列表
*/

@KafkaListener(topics = {"topic1"}, groupId = "myConsumerGroup")
public void listenBatchMessages(List<ConsumerRecord<String, String>> records) {
System.out.println("📦 收到批量消息,数量: " + records.size());

for (ConsumerRecord<String, String> record : records) {
try {
System.out.println("🔹 处理批量消息: " +
"主题=" + record.topic() +
", 分区=" + record.partition() +
", 偏移量=" + record.offset() +
", 消息=" + record.value());

// 批量处理逻辑
processBusinessLogic(record.value());

} catch (Exception e) {
System.err.println("❌ 批量消息处理失败: " + e.getMessage());
// 记录失败消息,可以单独处理或重试
}
}
System.out.println("✅ 批量消息处理完成");
}

/**
* 3. 消息转发(处理后将结果发送到另一个主题)
* 使用场景:消息处理管道,ETL处理
*
* @param record 原始消息记录
* @return 转发到 topic2 的消息内容
*/

@KafkaListener(topics = {"topic1"}, groupId = "myConsumerGroup")
@SendTo("topic2") // 将返回值发送到 topic2
public String processAndForwardMessage(ConsumerRecord<String, String> record) {
try {
System.out.println("🔄 处理并转发消息: " + record.value());

// 业务处理逻辑
String processedMessage = processBusinessLogic(record.value());

// 返回处理后的消息,会自动发送到 topic2
String result = "处理完成: " + processedMessage + " | 原始消息: " + record.value();
System.out.println("➡️ 转发消息到 topic2: " + result);

return result;

} catch (Exception e) {
System.err.println("❌ 消息处理转发失败: " + e.getMessage());
return "处理失败: " + e.getMessage();
}
}

/**
* 4. 指定分区消费
* 使用场景:精确控制消费哪些分区
*/

@KafkaListener(topicPartitions = {
@TopicPartition(
topic = "topic1",
partitions = {"0", "1", "2"}, // 消费分区 0,1,2
partitionOffsets = {
@PartitionOffset(partition = "3", initialOffset = "100"), // 从偏移量100开始消费分区3
@PartitionOffset(partition = "4", initialOffset = "0") // 从开头消费分区4
}
)
})
public void listenSpecificPartitions(ConsumerRecord<String, String> record) {
System.out.println("🎯 指定分区消费: " +
"分区=" + record.partition() +
", 偏移量=" + record.offset() +
", 消息=" + record.value());
}

/**
* 5. 异常处理消费者
* 使用场景:处理消费失败的消息
*/

@KafkaListener(topics = {"topic1.DLT"}, groupId = "myConsumerGroup") // DLT: Dead Letter Topic
public void listenDeadLetterMessages(ConsumerRecord<String, String> record) {
System.err.println("💀 死信队列消息: " +
"主题=" + record.topic() +
", 原始主题=" + record.headers().lastHeader("original-topic") +
", 异常=" + record.headers().lastHeader("exception-message") +
", 消息=" + record.value());

// 死信消息处理逻辑(记录日志、人工干预等)
}

/**
* 业务处理逻辑(示例)
* @param message 消息内容
* @return 处理结果
*/

private String processBusinessLogic(String message) {
// 这里实现你的业务逻辑
// 例如:数据转换、数据库操作、调用外部服务等

// 模拟业务处理
try {
Thread.sleep(100); // 模拟处理时间
return message.toUpperCase(); // 示例:转换为大写
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new RuntimeException("业务处理被中断", e);
}
}
}

Replica副本

副本(Replica)的核心作用是提供数据冗余,提高系统的容错能力和数据可靠性。通过将数据复制到多个节点,即使部分节点发生故障,数据依然可用,从而保证Kafka服务的持续性和数据不丢失。

副本配置策略

  • Kafka默认创建1个副本,这意味着没有数据冗余,一旦Broker宕机,该副本上的数据将不可用。
  • 生产环境中,通常建议将副本因子(replication.factor)配置为 至少2,常见为3。配置为2可以提供一份数据备份,配置为3能提供更高的容错能力(允许同时两台Broker宕机而数据不丢失)。
  • 权衡:副本数量并非越多越好。过多的副本会:
    • 消耗更多的磁盘存储空间。
    • 增加集群内部网络数据传输的负载。
    • 在数据写入和同步时可能降低整体吞吐量。

副本角色与工作机制

Kafka集群中的副本分为两种角色:

  • Leader副本:每个分区都有一个Leader副本。它负责处理所有的客户端读写请求(生产者的写入和消费者的读取)。
  • Follower副本:其他副本均为Follower。它们不处理客户端请求,唯一的任务就是从Leader副本异步拉取数据,与Leader保持数据同步。
  • 副本集合术语

    Kafka分区中的所有副本被统称为 AR(Assigned Replicas)。

    AR被进一步划分为两个子集:

    • ISR(In-Sync Replicas):与Leader副本保持同步的副本集合(包括Leader自己)。这里的“同步”是一个动态过程,只要Follower在指定的时间内成功拉取到Leader的最新数据,它就被认为是“同步中”的。
    • OSR(Out-of-Sync Replicas):与Leader副本同步滞后过多的Follower副本集合。

    ISR动态管理机制

    • 剔除机制:如果一个Follower副本在 replica.lag.time.max.ms 参数规定的时间(默认30秒)内没有向Leader发起拉取请求或未能追赶上Leader的最新进度,它就会被Leader从ISR中移除,转移到OSR中。
    • 加入机制:如果OSR中的副本重新追上了Leader的进度,并且持续健康同步一段时间,它会被重新加入ISR。
    • Leader选举:当Leader副本发生故障时,Kafka不会从AR中随机选举,而是优先从ISR集合中选举新的Leader。因为ISR中的副本拥有最完整的数据,这样可以最大限度地保证数据的一致性,避免数据丢失。

    核心要点

    • 一主多从:Kafka 分区的副本机制遵循“一主多从”的原则。一个分区有且仅有一个 Leader,可以有多个 Follower。
    • 读写分离:所有读写流量都由 Leader 承担,Follower 只负责备份。这是为了实现简单而高效的一致性模型。
    • 高可用:正是因为有了 Follower 副本,当 Leader 失效时,系统才能无缝地切换到一个最新的 Follower 上,从而保证服务的持续可用性和数据的可靠性。

    #mermaid-svg-MytNG7FzQW9ZDywS {font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}#mermaid-svg-MytNG7FzQW9ZDywS .error-icon{fill:#552222;}#mermaid-svg-MytNG7FzQW9ZDywS .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-MytNG7FzQW9ZDywS .edge-thickness-normal{stroke-width:2px;}#mermaid-svg-MytNG7FzQW9ZDywS .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-MytNG7FzQW9ZDywS .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-MytNG7FzQW9ZDywS .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-MytNG7FzQW9ZDywS .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-MytNG7FzQW9ZDywS .marker{fill:#333333;stroke:#333333;}#mermaid-svg-MytNG7FzQW9ZDywS .marker.cross{stroke:#333333;}#mermaid-svg-MytNG7FzQW9ZDywS svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-MytNG7FzQW9ZDywS .label{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;color:#333;}#mermaid-svg-MytNG7FzQW9ZDywS .cluster-label text{fill:#333;}#mermaid-svg-MytNG7FzQW9ZDywS .cluster-label span{color:#333;}#mermaid-svg-MytNG7FzQW9ZDywS .label text,#mermaid-svg-MytNG7FzQW9ZDywS span{fill:#333;color:#333;}#mermaid-svg-MytNG7FzQW9ZDywS .node rect,#mermaid-svg-MytNG7FzQW9ZDywS .node circle,#mermaid-svg-MytNG7FzQW9ZDywS .node ellipse,#mermaid-svg-MytNG7FzQW9ZDywS .node polygon,#mermaid-svg-MytNG7FzQW9ZDywS .node path{fill:#ECECFF;stroke:#9370DB;stroke-width:1px;}#mermaid-svg-MytNG7FzQW9ZDywS .node .label{text-align:center;}#mermaid-svg-MytNG7FzQW9ZDywS .node.clickable{cursor:pointer;}#mermaid-svg-MytNG7FzQW9ZDywS .arrowheadPath{fill:#333333;}#mermaid-svg-MytNG7FzQW9ZDywS .edgePath .path{stroke:#333333;stroke-width:2.0px;}#mermaid-svg-MytNG7FzQW9ZDywS .flowchart-link{stroke:#333333;fill:none;}#mermaid-svg-MytNG7FzQW9ZDywS .edgeLabel{background-color:#e8e8e8;text-align:center;}#mermaid-svg-MytNG7FzQW9ZDywS .edgeLabel rect{opacity:0.5;background-color:#e8e8e8;fill:#e8e8e8;}#mermaid-svg-MytNG7FzQW9ZDywS .cluster rect{fill:#ffffde;stroke:#aaaa33;stroke-width:1px;}#mermaid-svg-MytNG7FzQW9ZDywS .cluster text{fill:#333;}#mermaid-svg-MytNG7FzQW9ZDywS .cluster span{color:#333;}#mermaid-svg-MytNG7FzQW9ZDywS div.mermaidTooltip{position:absolute;text-align:center;max-width:200px;padding:2px;font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:12px;background:hsl(80, 100%, 96.2745098039%);border:1px solid #aaaa33;border-radius:2px;pointer-events:none;z-index:100;}#mermaid-svg-MytNG7FzQW9ZDywS :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}分区 P1分区 P0Kafka集群Leader 副本Follower 副本Leader 副本Follower 副本Broker 1Broker 2Broker 3

    Leader 选举流程

    在理解选举之前,必须牢记一个核心概念:Leader 选举并非从所有副本中随机挑选,而是优先从 ISR(In-Sync Replicas)列表中进行。 因为 ISR 中的副本拥有最完整的数据,这能最大限度地保证数据一致性。

    选举的两种主要触发场景

    场景一:Broker 宕机(最常见)

    假设我们有一个分区,它有 3 个副本,分布在 3 个 Broker 上。

    • Leader: Broker1
    • Follower (ISR): Broker2, Broker3
  • 故障发生:Broker1 由于网络问题、硬件故障或进程崩溃而宕机。
  • 检测故障:Kafka 集群通过 ZooKeeper 或 KRaft 模式(新版本)的心跳机制,检测到与 Broker1 的连接丢失。
  • 触发选举:集群的控制器(Controller,一个特殊的 Broker)检测到这一变化,并自动发起对该分区的 Leader 选举。
  • 选举决策:控制器查看该分区的 ISR 列表。当前的 ISR 是 [Broker2, Broker3]。
  • 选择新 Leader:控制器会从 ISR 列表中选择一个副本作为新的 Leader。通常,它是选择 ISR 列表中的第一个副本(但这并非绝对,策略可以配置)。假设它选择了 Broker2。
  • 元数据更新:控制器将新的 Leader 信息(Broker2)更新到集群的元数据中。
  • 流量切换:所有生产和消费该分区的客户端,会从集群元数据中得知 Leader 已变更为 Broker2,随后将流量切换到 Broker2。
  • 此过程对客户端是透明的,可能会有几秒到几十秒的不可用时间。

    为了更直观地理解这一过程,下图描绘了 Broker1 宕机后,Broker2 被选举为新 Leader 的完整流程:

    #mermaid-svg-hq9AnDR6Ci9zoUAY {font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}#mermaid-svg-hq9AnDR6Ci9zoUAY .error-icon{fill:#552222;}#mermaid-svg-hq9AnDR6Ci9zoUAY .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-hq9AnDR6Ci9zoUAY .edge-thickness-normal{stroke-width:2px;}#mermaid-svg-hq9AnDR6Ci9zoUAY .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-hq9AnDR6Ci9zoUAY .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-hq9AnDR6Ci9zoUAY .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-hq9AnDR6Ci9zoUAY .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-hq9AnDR6Ci9zoUAY .marker{fill:#333333;stroke:#333333;}#mermaid-svg-hq9AnDR6Ci9zoUAY .marker.cross{stroke:#333333;}#mermaid-svg-hq9AnDR6Ci9zoUAY svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-hq9AnDR6Ci9zoUAY .label{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;color:#333;}#mermaid-svg-hq9AnDR6Ci9zoUAY .cluster-label text{fill:#333;}#mermaid-svg-hq9AnDR6Ci9zoUAY .cluster-label span{color:#333;}#mermaid-svg-hq9AnDR6Ci9zoUAY .label text,#mermaid-svg-hq9AnDR6Ci9zoUAY span{fill:#333;color:#333;}#mermaid-svg-hq9AnDR6Ci9zoUAY .node rect,#mermaid-svg-hq9AnDR6Ci9zoUAY .node circle,#mermaid-svg-hq9AnDR6Ci9zoUAY .node ellipse,#mermaid-svg-hq9AnDR6Ci9zoUAY .node polygon,#mermaid-svg-hq9AnDR6Ci9zoUAY .node path{fill:#ECECFF;stroke:#9370DB;stroke-width:1px;}#mermaid-svg-hq9AnDR6Ci9zoUAY .node .label{text-align:center;}#mermaid-svg-hq9AnDR6Ci9zoUAY .node.clickable{cursor:pointer;}#mermaid-svg-hq9AnDR6Ci9zoUAY .arrowheadPath{fill:#333333;}#mermaid-svg-hq9AnDR6Ci9zoUAY .edgePath .path{stroke:#333333;stroke-width:2.0px;}#mermaid-svg-hq9AnDR6Ci9zoUAY .flowchart-link{stroke:#333333;fill:none;}#mermaid-svg-hq9AnDR6Ci9zoUAY .edgeLabel{background-color:#e8e8e8;text-align:center;}#mermaid-svg-hq9AnDR6Ci9zoUAY .edgeLabel rect{opacity:0.5;background-color:#e8e8e8;fill:#e8e8e8;}#mermaid-svg-hq9AnDR6Ci9zoUAY .cluster rect{fill:#ffffde;stroke:#aaaa33;stroke-width:1px;}#mermaid-svg-hq9AnDR6Ci9zoUAY .cluster text{fill:#333;}#mermaid-svg-hq9AnDR6Ci9zoUAY .cluster span{color:#333;}#mermaid-svg-hq9AnDR6Ci9zoUAY div.mermaidTooltip{position:absolute;text-align:center;max-width:200px;padding:2px;font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:12px;background:hsl(80, 100%, 96.2745098039%);border:1px solid #aaaa33;border-radius:2px;pointer-events:none;z-index:100;}#mermaid-svg-hq9AnDR6Ci9zoUAY :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}Broker1 原始 Leader 宕机控制器检测到故障控制器检查该分区的 ISR 列表发现 Broker2 和 Broker3控制器从 ISR 中选举 Broker2 为新 Leader控制器更新集群元数据生产者/消费者将流量切换到新 Leader Broker2

    场景二:优雅重启或维护

  • 主动触发:管理员通过工具执行 Broker 的重启或关闭。
  • Leader 卸载:该 Broker 上的所有 Leader 分区会主动将其领导权移交给其 ISR 中的另一个副本。这个过程称为 LeaderAndIsr 请求。
  • 无缝切换:因为这个过程是受控的,副本同步状态良好,所以选举和切换会非常快速和平滑,对客户端的影响最小。
  • 极端情况与数据一致性保障

    如果 ISR 列表为空怎么办?

    当所有 Follower 副本都因同步滞后而被踢出 ISR,导致 ISR 列表中只剩下 Leader 自己时,如果此时 Leader 也宕机了,就会发生 ISR 为空的情况。

    此时,Kafka 提供了两种配置策略:

  • unclean.leader.election.enable = false(默认,推荐)
    • 禁止不洁选举:控制器将不会从 ISR 以外的副本中选举 Leader。
    • 结果:该分区会不可用,直到原 Leader 重新上线。这是一种 “宁可停止服务,也不丢失数据” 的保守策略,保证了数据一致性。
  • unclean.leader.election.enable = true
    • 允许不洁选举:控制器可以从 OSR(不同步的副本)中选举一个新的 Leader。
    • 风险:这个新的 Leader 可能缺少原 Leader 已确认的最新消息,导致数据丢失。
    • 适用场景:对可用性要求高于数据一致性的场景(如日志收集)。
  • LEO 和 HW 的详细解析

    1. LEO – 日志末端位移

    • 定义:每个副本都有自己的LEO。它指向当前副本日志中下一条待写入消息的位置。
    • 特性:
      • 副本私有:每个副本(Leader和Follower)都独立维护自己的LEO。
      • 动态变化:
        • 当生产者向Leader写入一条新消息时,Leader的LEO会增加。
        • 当Follower从Leader拉取到新数据后,Follower的LEO也会增加。
    • 作用:它标识了一个副本当前已存储的最新数据的位置。

    2. HW – 高水位线

    • 定义:HW是 所有副本(具体是指ISR中的所有副本)中最小的LEO。
    • 特性:
      • 分区全局:对于一个分区而言,只有一个HW值,所有副本看到的HW都是一样的。
      • 一致性保证:它代表了一份已被所有ISR副本成功复制的数据的边界。
    • 核心作用:
    • 定义消息可见性:消费者只能拉取到HW之前的消息。HW之后的消息对消费者是不可见的,即使Leader副本已经存储了这些消息。
    • 定义数据安全线:HW之前的消息被认为是“已提交”的,即使Leader副本宕机,这些消息也不会丢失,因为ISR中至少有一个其他副本也拥有了这部分数据。

    工作流程示例

    假设有一个分区,副本因子为3(1个Leader L,2个Follower F1, F2),初始状态所有LEO和HW都是0。

  • 生产者发送消息:
    • 生产者向Leader发送两条消息,M1和M2。
    • Leader将其写入本地日志。
    • Leader LEO = 2 (有了两条消息,下一条消息的位置是2)
    • HW = 0 (因为Follower还没同步)
  • Follower开始拉取数据:
    • F1和F2向Leader发送抓取请求。
    • Leader不仅返回数据,还会在响应中带上当前的HW。
    • F1成功拉取了两条消息,F1 LEO = 2。
    • F2网络较慢,只成功拉取了一条消息,F2 LEO = 1。
  • 更新HW:
    • Leader会计算新的HW:HW = min(LEO_L, LEO_F1, LEO_F2) = min(2, 2, 1) = 1
    • Leader将新的HW=1通知给所有副本。
  • 此时的状态:
    • HW = 1
    • 消费者:只能消费到offset=0的消息(M1)。M2虽然已在Leader和F1上,但对消费者不可见。
  • F2追赶上进度:
    • F2拉取到M2,F2 LEO 更新为2。
    • Leader重新计算HW:HW = min(2, 2, 2) = 2
    • 消费者:现在可以消费到M1和M2了。
  • 为了更直观地展示这一过程,下图模拟了上述步骤中LEO和HW的动态变化:

    flowchart TD
    subgraph S1 [第1步: 生产者发送消息]
    A1[生产者发送 M1, M2] –> A2[Leader 写入本地日志<br>Leader LEO=2]
    A2 –> A3[计算 HW=min2,2,1=0<br>HW 保持为0]
    end

    subgraph S2 [第2步: Follower 拉取数据]
    B1[F1 拉取 M1,M2<br>F1 LEO=2]
    B2[F2 拉取 M1<br>F2 LEO=1]
    B1 & B2 –> B3[计算 HW=min2,2,1=1<br>HW 更新为1]
    end

    subgraph S3 [第3步: 消费者可见性]
    C1[HW=1] –> C2[消费者只能看到 M1<br>offset=0]
    end

    subgraph S4 [第4步: F2 追赶进度]
    D1[F2 拉取 M2<br>F2 LEO=2] –> D2[计算 HW=min2,2,2=2<br>HW 更新为2]
    D2 –> D3[消费者可以看到 M1 和 M2<br>offset=0 和 1]
    end

    S1 –> S2 –> S3 –> S4

    为什么需要HW和LEO?

    • 可靠性:确保了即使Leader宕机,只有被充分复制的数据才会被消费者看到,从而避免了数据丢失。
    • 性能:Follower副本的同步是异步的,允许Leader在Follower尚未完全同步的情况下继续接收新消息,提高了吞吐量。HW机制则是在性能和可靠性之间提供了一个完美的平衡点。

    与ISR的关系

    • ISR 是一个副本列表,定义了哪些Follower是“活跃的、同步的”。
    • LEO 是每个副本的内部状态。
    • HW 是根据ISR中所有副本的LEO计算出来的,是副本同步状态的最终体现。

    如果一个Follower的LEO持续远低于Leader的LEO,它就会被踢出ISR,而HW的计算也将不再考虑它。

    发送消息(生产者)

    分配策略

    生产者决定消息发送到哪个分区的策略由 partitioner.class 参数控制。以下是常见的分区策略:

    1. 默认策略

    默认策略是 粘性分区策略,它在保证高效批次发送的同时,兼顾了负载均衡。

    • 当消息指定 Key 时:
      • 对 Key 进行哈希计算,然后根据哈希值对分区总数取模,得到目标分区。
      • 核心特点:同一个 Key 的消息总是被发送到同一个分区。这保证了同一业务实体的事件在分区内有序。
    • 当消息未指定 Key 时:
      • 生产者会随机选择一个分区,并在一段时间内(或直到该批次被填满并发送)“粘性”地将所有无Key消息发送到该分区。
      • 核心特点:与纯随机或轮询相比,这种“粘性”能构建更大的数据批次,显著减少网络请求次数,提升吞吐量。批次发送后,会重新选择下一个粘性分区。

    2. 轮询策略

    • 工作机制:将消息依次、循环地发送到所有可用分区。
    • 效果:在分区级别实现了最极致的负载均衡,每个分区接收到的消息数量几乎相等。
    • 适用场景:当消息没有逻辑分组需求,且追求绝对均衡的负载时。

    3. 自定义策略

    • 实现方式:实现 org.apache.kafka.clients.producer.Partitioner 接口,并重写 partition() 方法,在其中编写自定义的分区逻辑。
    • 配置:在生产者配置中指定 partitioner.class 为您的自定义实现类。
    • 适用场景:需要根据业务逻辑(如特定字段、数据范围、服务区域等)进行复杂分区的场景。

    生产环境建议:

    • 优先考虑无序:大部分场景下,Kafka不保证全局消息有序。应优先考虑使用无Key的粘性策略来获得最佳性能和负载均衡。
    • 谨慎使用 Key:仅在必须保证分区内有序(即局部有序)时,才为消息设置Key。同时,要监控分区数据量,确保不会因少数热点Key导致严重的数据倾斜。
    • 自定义策略的潜力:在高级用法中,可以通过自定义分区器来缓解数据倾斜。例如,对热点Key添加随机后缀进行“打散”,在消费端再做聚合,但这会增加系统的复杂性。

    流程

    生产者客户端的核心架构是一个双线程(主线程 + Sender线程)与一个共享缓冲区的协作模型。整个发送流程可以清晰地分为以下几个阶段:

    #mermaid-svg-CHKtWPyiY0Ndelr8 {font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}#mermaid-svg-CHKtWPyiY0Ndelr8 .error-icon{fill:#552222;}#mermaid-svg-CHKtWPyiY0Ndelr8 .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-CHKtWPyiY0Ndelr8 .edge-thickness-normal{stroke-width:2px;}#mermaid-svg-CHKtWPyiY0Ndelr8 .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-CHKtWPyiY0Ndelr8 .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-CHKtWPyiY0Ndelr8 .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-CHKtWPyiY0Ndelr8 .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-CHKtWPyiY0Ndelr8 .marker{fill:#333333;stroke:#333333;}#mermaid-svg-CHKtWPyiY0Ndelr8 .marker.cross{stroke:#333333;}#mermaid-svg-CHKtWPyiY0Ndelr8 svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-CHKtWPyiY0Ndelr8 .label{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;color:#333;}#mermaid-svg-CHKtWPyiY0Ndelr8 .cluster-label text{fill:#333;}#mermaid-svg-CHKtWPyiY0Ndelr8 .cluster-label span{color:#333;}#mermaid-svg-CHKtWPyiY0Ndelr8 .label text,#mermaid-svg-CHKtWPyiY0Ndelr8 span{fill:#333;color:#333;}#mermaid-svg-CHKtWPyiY0Ndelr8 .node rect,#mermaid-svg-CHKtWPyiY0Ndelr8 .node circle,#mermaid-svg-CHKtWPyiY0Ndelr8 .node ellipse,#mermaid-svg-CHKtWPyiY0Ndelr8 .node polygon,#mermaid-svg-CHKtWPyiY0Ndelr8 .node path{fill:#ECECFF;stroke:#9370DB;stroke-width:1px;}#mermaid-svg-CHKtWPyiY0Ndelr8 .node .label{text-align:center;}#mermaid-svg-CHKtWPyiY0Ndelr8 .node.clickable{cursor:pointer;}#mermaid-svg-CHKtWPyiY0Ndelr8 .arrowheadPath{fill:#333333;}#mermaid-svg-CHKtWPyiY0Ndelr8 .edgePath .path{stroke:#333333;stroke-width:2.0px;}#mermaid-svg-CHKtWPyiY0Ndelr8 .flowchart-link{stroke:#333333;fill:none;}#mermaid-svg-CHKtWPyiY0Ndelr8 .edgeLabel{background-color:#e8e8e8;text-align:center;}#mermaid-svg-CHKtWPyiY0Ndelr8 .edgeLabel rect{opacity:0.5;background-color:#e8e8e8;fill:#e8e8e8;}#mermaid-svg-CHKtWPyiY0Ndelr8 .cluster rect{fill:#ffffde;stroke:#aaaa33;stroke-width:1px;}#mermaid-svg-CHKtWPyiY0Ndelr8 .cluster text{fill:#333;}#mermaid-svg-CHKtWPyiY0Ndelr8 .cluster span{color:#333;}#mermaid-svg-CHKtWPyiY0Ndelr8 div.mermaidTooltip{position:absolute;text-align:center;max-width:200px;padding:2px;font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:12px;background:hsl(80, 100%, 96.2745098039%);border:1px solid #aaaa33;border-radius:2px;pointer-events:none;z-index:100;}#mermaid-svg-CHKtWPyiY0Ndelr8 :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}Sender线程主线程 Main ThreadYesNo从累积器获取就绪的批次批次是否已满/超时?创建ClientRequest通过Selector将请求发送到网络I/O分区器 Partitioner序列化器 Serializer拦截器 InterceptorsProducer: 创建ProducerRecord记录累积器 RecordAccumulator消息存入缓冲区按Topic-Partition分批Broker: 接收并处理请求Producer: 收到响应调用用户指定的Callback可能触发拦截器的onAcknowledgement

    阶段一:主线程处理流程(图中上半部分)

    此阶段在发送 KafkaTemplate.send() 的调用线程中执行。

  • 拦截器

    • 作用:在消息发送前以及收到服务端应答后,提供可插拔的定制化逻辑,如消息审计、属性增强、链路追踪、指标收集等。

    • 实现:

      java

      // 1. 实现 org.apache.kafka.clients.producer.ProducerInterceptor 接口
      public class CustomProducerInterceptor implements ProducerInterceptor<String, String> {

      // 发送前调用,可以对消息进行修改或统计
      @Override
      public ProducerRecord<String, String> onSend(ProducerRecord<String, String> record) {
      System.out.println("准备发送消息至主题: " + record.topic());
      // 例如: 为消息头添加追踪ID
      record.headers().add("trace-id", UUID.randomUUID().toString().getBytes());
      return record; // 可以返回null来过滤掉此消息
      }

      // 在收到服务端响应或发送失败时调用,用于审计或统计
      @Override
      public void onAcknowledgement(RecordMetadata metadata, Exception exception) {
      if (exception == null) {
      System.out.println("消息发送成功: " + metadata);
      } else {
      System.err.println("消息发送失败: " + exception.getMessage());
      }
      }

      @Override
      public void close() { /* 关闭拦截器,释放资源 */ }

      @Override
      public void configure(Map<String, ?> configs) { /* 读取配置 */ }
      }

    • 配置:

      @Configuration
      public class KafkaConfig {

      @Bean
      public Map<String, Object> producerConfigs() {
      Map<String, Object> props = new HashMap<>();
      // … 其他配置(如bootstrap.servers, key.serializer等)
      // 指定自定义拦截器,多个拦截器用逗号分隔,按顺序执行
      props.put(ProducerConfig.INTERCEPTOR_CLASSES_CONFIG,
      "com.yourpackage.CustomProducerInterceptor");
      return props;
      }
      }

  • 序列化器

    • 作用:将消息的 Key 和 Value 从 Java 对象转换为字节数组,以便在网络中传输。
    • 常用序列化器:
      • StringSerializer: 用于字符串序列化。
      • IntegerSerializer, LongSerializer: 用于基本类型。
      • ByteArraySerializer: 用于已有的字节数组。
      • 推荐:对于复杂对象,推荐使用 JsonSerializer 或高效的二进制序列化框架(如 Apache Avro, Protocol Buffers)。
  • 分区器

    • 作用:确定消息应该被发送到主题的哪个分区。
    • 策略:
      • 指定了 Key:默认对 Key 进行哈希( murmur2 算法),然后对分区总数取模,确保相同 Key 的消息总是进入同一分区,实现分区内有序。
      • 未指定 Key:采用粘性分区策略。随机选择一个分区,并在一段时间内或直到批次被填满前,将所有无 Key 消息都发送到该分区,以形成更大的批次,提升吞吐量。批次发送后,会随机切换到下一个分区。
  • 阶段二:Sender线程与网络通信(图中下半部分)

    此阶段由Kafka生产者内部的后台Sender线程异步完成。

  • 记录累加器
    • 作用:主线程处理完的消息并不会立即发送,而是被追加到一个名为 RecordAccumulator 的双端队列缓冲区中。
    • 批处理:累加器会为每个主题分区维护一个 ProducerBatch(消息批次)。新消息会追加到对应的批次中。这种按分区批处理的机制是Kafka高吞吐量的关键之一。
  • Sender线程
    • 作用:一个后台守护线程,负责将累加器中已就绪的批次发送到Kafka集群。
    • 触发条件:当一个批次满足以下任一条件时,Sender线程会将其取出并发送:
      • batch.size:批次大小达到该阈值。
      • linger.ms:批次创建后等待的时间达到该阈值。这是为了在吞吐量和延迟之间取得平衡(默认为0,表示立即发送,但粘性分区策略会使其行为更智能)。
    • 网络通信:Sender线程将多个批次按照目标Node(Broker)进行分组,构建 ClientRequest,并通过 Selector 网络选择器以非阻塞I/O的方式批量发送出去。
  • 响应处理与回调
    • Broker处理请求后,会返回一个 ProducerResponse。
    • Sender线程收到响应后,会解析它,并根据结果(成功或失败):
      • 调用我们通过 .addCallback() 设置的回调函数。
      • 触发拦截器的 onAcknowledgement() 方法。
      • 无论成功与否,都会释放批次占用的内存空间,以备复用。
  • 消费消息

    分区策略

    RangeAssignor(范围分配器)- 默认策略

    工作原理:

    • 按主题分区,对每个主题独立进行分配
    • 按字典序对消费者排序,对分区排序
    • 计算每个消费者应该消费的分区数量:分区数 / 消费者数
    • 多余的分区按顺序分配给前面的消费者

    示例:

    // 假设:主题A有3个分区,主题B有2个分区,消费者组有2个消费者
    消费者C1:主题AP0, 主题AP1, 主题BP0
    消费者C2:主题AP2, 主题BP1

    配置:

    spring:
    kafka:
    consumer:
    properties:
    partition.assignment.strategy: org.apache.kafka.clients.consumer.RangeAssignor

    RoundRobinAssignor(轮询分配器)

    工作原理:

    • 将所有主题的所有分区放在一起
    • 按字典序对消费者和分区排序
    • 轮询方式将分区分配给消费者

    示例:

    // 假设:主题A有3个分区,主题B有2个分区,消费者组有2个消费者
    消费者C1:主题AP0, 主题AP2, 主题BP1
    消费者C2:主题AP1, 主题BP0

    StickyAssignor(粘性分配器)

    工作原理:

    • 在轮询的基础上,尽量减少分区重新分配
    • 重平衡时尽可能保持原有的分配关系
    • 提供更均衡的分配,减少"停止世界"的影响

    CooperativeStickyAssignor(协作粘性分配器)

    工作原理:

    • Kafka 2.4+ 引入,支持增量协同重平衡
    • 消费者可以继续处理未被重新分配的分区
    • 减少重平衡期间的停机时间

    流程

    #mermaid-svg-7PAjw5fvNkRRwJmN {font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}#mermaid-svg-7PAjw5fvNkRRwJmN .error-icon{fill:#552222;}#mermaid-svg-7PAjw5fvNkRRwJmN .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-7PAjw5fvNkRRwJmN .edge-thickness-normal{stroke-width:2px;}#mermaid-svg-7PAjw5fvNkRRwJmN .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-7PAjw5fvNkRRwJmN .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-7PAjw5fvNkRRwJmN .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-7PAjw5fvNkRRwJmN .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-7PAjw5fvNkRRwJmN .marker{fill:#333333;stroke:#333333;}#mermaid-svg-7PAjw5fvNkRRwJmN .marker.cross{stroke:#333333;}#mermaid-svg-7PAjw5fvNkRRwJmN svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-7PAjw5fvNkRRwJmN .label{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;color:#333;}#mermaid-svg-7PAjw5fvNkRRwJmN .cluster-label text{fill:#333;}#mermaid-svg-7PAjw5fvNkRRwJmN .cluster-label span{color:#333;}#mermaid-svg-7PAjw5fvNkRRwJmN .label text,#mermaid-svg-7PAjw5fvNkRRwJmN span{fill:#333;color:#333;}#mermaid-svg-7PAjw5fvNkRRwJmN .node rect,#mermaid-svg-7PAjw5fvNkRRwJmN .node circle,#mermaid-svg-7PAjw5fvNkRRwJmN .node ellipse,#mermaid-svg-7PAjw5fvNkRRwJmN .node polygon,#mermaid-svg-7PAjw5fvNkRRwJmN .node path{fill:#ECECFF;stroke:#9370DB;stroke-width:1px;}#mermaid-svg-7PAjw5fvNkRRwJmN .node .label{text-align:center;}#mermaid-svg-7PAjw5fvNkRRwJmN .node.clickable{cursor:pointer;}#mermaid-svg-7PAjw5fvNkRRwJmN .arrowheadPath{fill:#333333;}#mermaid-svg-7PAjw5fvNkRRwJmN .edgePath .path{stroke:#333333;stroke-width:2.0px;}#mermaid-svg-7PAjw5fvNkRRwJmN .flowchart-link{stroke:#333333;fill:none;}#mermaid-svg-7PAjw5fvNkRRwJmN .edgeLabel{background-color:#e8e8e8;text-align:center;}#mermaid-svg-7PAjw5fvNkRRwJmN .edgeLabel rect{opacity:0.5;background-color:#e8e8e8;fill:#e8e8e8;}#mermaid-svg-7PAjw5fvNkRRwJmN .cluster rect{fill:#ffffde;stroke:#aaaa33;stroke-width:1px;}#mermaid-svg-7PAjw5fvNkRRwJmN .cluster text{fill:#333;}#mermaid-svg-7PAjw5fvNkRRwJmN .cluster span{color:#333;}#mermaid-svg-7PAjw5fvNkRRwJmN div.mermaidTooltip{position:absolute;text-align:center;max-width:200px;padding:2px;font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:12px;background:hsl(80, 100%, 96.2745098039%);border:1px solid #aaaa33;border-radius:2px;pointer-events:none;z-index:100;}#mermaid-svg-7PAjw5fvNkRRwJmN :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}Heartbeat线程Fetcher线程主线程 Main Thread周期性发送心跳请求检测消费组状态防止超时触发Rebalance从RecordAccumulator中读取assigned分区Fetcher线程启动向Broker发送FetchRequest接收批量消息并写入本地缓冲区拉取消息结果ConsumerRecords执行分区分配Rebalance向GroupCoordinator注册并加入消费组Consumer: 调用poll发起拉取请求反序列化器 Deserializer业务处理逻辑提交Offset自动或手动可能触发拦截器的onConsume回调

    阶段一:主线程拉取与消费流程(上半部分)

    ① 加入消费组(Join Group)

    • 目标:让消费者与集群协调器(GroupCoordinator)建立连接,加入对应的 Consumer Group。
    • 流程:
    • 消费者启动时,会发送 JoinGroupRequest。
    • 由 GroupCoordinator 完成组成员注册、选出组协调者(leader)。
    • leader 消费者根据分区分配策略(Range、RoundRobin、Sticky 等)分配分区。
    • 每个消费者接收分配结果,开始各自负责的分区消费。

    ② 拉取消息(poll)

    poll() 是消费的核心入口:

    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));

    内部执行:

  • 检查消费组状态(确保未在 rebalance)。
  • 若有可用分区 → 从 Fetcher线程 缓冲区中获取已拉取的消息。
  • 若缓冲区为空 → 阻塞等待或发起新的拉取请求。
  • 返回消息给上层应用处理。

  • ③ 反序列化(Deserializer)

    Kafka 消息从 Broker 拉取时是二进制的,Consumer 通过反序列化器将其还原为 Java 对象:

    props.put("key.deserializer", StringDeserializer.class.getName());
    props.put("value.deserializer", JsonDeserializer.class.getName());

    支持多种类型:StringDeserializer、LongDeserializer、ByteArrayDeserializer 等。


    ④ 业务处理逻辑(Processing)

    应用程序在此阶段处理业务逻辑,比如:

    for (ConsumerRecord<String, String> record : records) {
    System.out.printf("topic=%s, partition=%d, offset=%d, value=%s%n",
    record.topic(), record.partition(), record.offset(), record.value());
    // 执行业务处理
    }

    注意:Kafka Consumer 是非线程安全的,如需多线程消费,需谨慎管理 offset。


    ⑤ Offset 提交(Commit Offset)

    Kafka 通过偏移量(Offset)记录消费进度。
    提交 offset 可以自动或手动进行。

    🔹 自动提交

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

    优点:简单
    缺点:可能在程序异常时重复或丢失消息。

    🔹 手动提交

    consumer.commitSync(); // 同步提交,可靠但阻塞
    consumer.commitAsync(); // 异步提交,性能好但可能丢失少量偏移

    生产环境常结合两者:

    • 正常时用 commitAsync()
    • 关闭前或 rebalance 前用 commitSync()

    阶段二:后台线程与协作机制(下半部分)

    Kafka Consumer 内部维护多个后台线程,保证消费的高效与稳定:


    ① Fetcher 拉取线程(Fetcher Thread)

    职责:从 Broker 拉取消息,写入本地缓冲区。

    • 采用 批量拉取机制,按分区维护 fetchPosition(当前读取 offset)。
    • 当主线程调用 poll() 时,Fetcher 会提前准备好一批数据。
    • 支持参数:
      • fetch.min.bytes
      • max.partition.fetch.bytes
      • fetch.max.wait.ms

    这种异步拉取 + 本地缓冲策略,让消费过程更顺滑。


    ② 心跳线程(Heartbeat Thread)

    职责:维持与协调器的会话活性,防止被判定为“失活”从而触发 rebalance。

    心跳周期受参数控制:

    heartbeat.interval.ms = 3000
    session.timeout.ms = 10000

    若超过 session.timeout.ms 未发送心跳,协调器会认为该消费者已失联,并触发再均衡。


    ③ 重平衡(Rebalance)

    触发条件:

    • 消费者加入或离开组;
    • 订阅的 topic 分区变化;
    • 心跳超时。

    影响:

    • 当前分区暂停消费;
    • 提交 offset;
    • 重新分配新分区;
    • 消费者从新位置恢复。

    新版使用 StickyAssignor,能尽量保证分区分配稳定性。


    ④ Offset 存储与恢复

    消费者的 offset 信息存储在 Kafka 内部主题 __consumer_offsets。

    当消费者重启后:

    • 从该主题读取上次提交的 offset;
    • 从对应分区继续消费;
    • 实现 “断点续消费”。

    消息存储

    #mermaid-svg-AGTpoIOSOP4SCRub {font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}#mermaid-svg-AGTpoIOSOP4SCRub .error-icon{fill:#552222;}#mermaid-svg-AGTpoIOSOP4SCRub .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-AGTpoIOSOP4SCRub .edge-thickness-normal{stroke-width:2px;}#mermaid-svg-AGTpoIOSOP4SCRub .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-AGTpoIOSOP4SCRub .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-AGTpoIOSOP4SCRub .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-AGTpoIOSOP4SCRub .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-AGTpoIOSOP4SCRub .marker{fill:#333333;stroke:#333333;}#mermaid-svg-AGTpoIOSOP4SCRub .marker.cross{stroke:#333333;}#mermaid-svg-AGTpoIOSOP4SCRub svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-AGTpoIOSOP4SCRub .label{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;color:#333;}#mermaid-svg-AGTpoIOSOP4SCRub .cluster-label text{fill:#333;}#mermaid-svg-AGTpoIOSOP4SCRub .cluster-label span{color:#333;}#mermaid-svg-AGTpoIOSOP4SCRub .label text,#mermaid-svg-AGTpoIOSOP4SCRub span{fill:#333;color:#333;}#mermaid-svg-AGTpoIOSOP4SCRub .node rect,#mermaid-svg-AGTpoIOSOP4SCRub .node circle,#mermaid-svg-AGTpoIOSOP4SCRub .node ellipse,#mermaid-svg-AGTpoIOSOP4SCRub .node polygon,#mermaid-svg-AGTpoIOSOP4SCRub .node path{fill:#ECECFF;stroke:#9370DB;stroke-width:1px;}#mermaid-svg-AGTpoIOSOP4SCRub .node .label{text-align:center;}#mermaid-svg-AGTpoIOSOP4SCRub .node.clickable{cursor:pointer;}#mermaid-svg-AGTpoIOSOP4SCRub .arrowheadPath{fill:#333333;}#mermaid-svg-AGTpoIOSOP4SCRub .edgePath .path{stroke:#333333;stroke-width:2.0px;}#mermaid-svg-AGTpoIOSOP4SCRub .flowchart-link{stroke:#333333;fill:none;}#mermaid-svg-AGTpoIOSOP4SCRub .edgeLabel{background-color:#e8e8e8;text-align:center;}#mermaid-svg-AGTpoIOSOP4SCRub .edgeLabel rect{opacity:0.5;background-color:#e8e8e8;fill:#e8e8e8;}#mermaid-svg-AGTpoIOSOP4SCRub .cluster rect{fill:#ffffde;stroke:#aaaa33;stroke-width:1px;}#mermaid-svg-AGTpoIOSOP4SCRub .cluster text{fill:#333;}#mermaid-svg-AGTpoIOSOP4SCRub .cluster span{color:#333;}#mermaid-svg-AGTpoIOSOP4SCRub div.mermaidTooltip{position:absolute;text-align:center;max-width:200px;padding:2px;font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:12px;background:hsl(80, 100%, 96.2745098039%);border:1px solid #aaaa33;border-radius:2px;pointer-events:none;z-index:100;}#mermaid-svg-AGTpoIOSOP4SCRub :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}TopicPartition 0Partition 1Partition 2Segment 000000.logSegment 000001.logSegment 000002.log消息集合RecordBatch消息集合RecordBatch

    层级内容特点
    Topic 消息分类 逻辑层
    Partition 并行单元 内部消息有序
    Segment 存储单元 由 .log/.index/.timeindex 组成
    Offset 位移标识 消费进度控制
    Cleanup Policy delete / compact 控制保留策略
    Replica & ISR 多副本同步 保证可靠性

    offset

    #mermaid-svg-TBCWlZJjTBrwWJM2 {font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}#mermaid-svg-TBCWlZJjTBrwWJM2 .error-icon{fill:#552222;}#mermaid-svg-TBCWlZJjTBrwWJM2 .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-TBCWlZJjTBrwWJM2 .edge-thickness-normal{stroke-width:2px;}#mermaid-svg-TBCWlZJjTBrwWJM2 .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-TBCWlZJjTBrwWJM2 .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-TBCWlZJjTBrwWJM2 .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-TBCWlZJjTBrwWJM2 .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-TBCWlZJjTBrwWJM2 .marker{fill:#333333;stroke:#333333;}#mermaid-svg-TBCWlZJjTBrwWJM2 .marker.cross{stroke:#333333;}#mermaid-svg-TBCWlZJjTBrwWJM2 svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-TBCWlZJjTBrwWJM2 .actor{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:#ECECFF;}#mermaid-svg-TBCWlZJjTBrwWJM2 text.actor>tspan{fill:black;stroke:none;}#mermaid-svg-TBCWlZJjTBrwWJM2 .actor-line{stroke:grey;}#mermaid-svg-TBCWlZJjTBrwWJM2 .messageLine0{stroke-width:1.5;stroke-dasharray:none;stroke:#333;}#mermaid-svg-TBCWlZJjTBrwWJM2 .messageLine1{stroke-width:1.5;stroke-dasharray:2,2;stroke:#333;}#mermaid-svg-TBCWlZJjTBrwWJM2 #arrowhead path{fill:#333;stroke:#333;}#mermaid-svg-TBCWlZJjTBrwWJM2 .sequenceNumber{fill:white;}#mermaid-svg-TBCWlZJjTBrwWJM2 #sequencenumber{fill:#333;}#mermaid-svg-TBCWlZJjTBrwWJM2 #crosshead path{fill:#333;stroke:#333;}#mermaid-svg-TBCWlZJjTBrwWJM2 .messageText{fill:#333;stroke:#333;}#mermaid-svg-TBCWlZJjTBrwWJM2 .labelBox{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:#ECECFF;}#mermaid-svg-TBCWlZJjTBrwWJM2 .labelText,#mermaid-svg-TBCWlZJjTBrwWJM2 .labelText>tspan{fill:black;stroke:none;}#mermaid-svg-TBCWlZJjTBrwWJM2 .loopText,#mermaid-svg-TBCWlZJjTBrwWJM2 .loopText>tspan{fill:black;stroke:none;}#mermaid-svg-TBCWlZJjTBrwWJM2 .loopLine{stroke-width:2px;stroke-dasharray:2,2;stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);}#mermaid-svg-TBCWlZJjTBrwWJM2 .note{stroke:#aaaa33;fill:#fff5ad;}#mermaid-svg-TBCWlZJjTBrwWJM2 .noteText,#mermaid-svg-TBCWlZJjTBrwWJM2 .noteText>tspan{fill:black;stroke:none;}#mermaid-svg-TBCWlZJjTBrwWJM2 .activation0{fill:#f4f4f4;stroke:#666;}#mermaid-svg-TBCWlZJjTBrwWJM2 .activation1{fill:#f4f4f4;stroke:#666;}#mermaid-svg-TBCWlZJjTBrwWJM2 .activation2{fill:#f4f4f4;stroke:#666;}#mermaid-svg-TBCWlZJjTBrwWJM2 .actorPopupMenu{position:absolute;}#mermaid-svg-TBCWlZJjTBrwWJM2 .actorPopupMenuPanel{position:absolute;fill:#ECECFF;box-shadow:0px 8px 16px 0px rgba(0,0,0,0.2);filter:drop-shadow(3px 5px 2px rgb(0 0 0 / 0.4));}#mermaid-svg-TBCWlZJjTBrwWJM2 .actor-man line{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:#ECECFF;}#mermaid-svg-TBCWlZJjTBrwWJM2 .actor-man circle,#mermaid-svg-TBCWlZJjTBrwWJM2 line{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:#ECECFF;stroke-width:2px;}#mermaid-svg-TBCWlZJjTBrwWJM2 :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}BrokerConsumer__consumer_offsets拉取消息(offset=100~105)处理消息commitSync(offset=106)提交确认下次拉取从 offset=106 开始BrokerConsumer__consumer_offsets
    flowchart LR
    subgraph KafkaCluster[Kafka 集群]
    direction TB

    subgraph TopicA[Topic: user-events]
    direction LR
    P0[Partition 0<br>offset 0..999]
    P1[Partition 1<br>offset 0..850]
    P2[Partition 2<br>offset 0..1200]
    end

    subgraph Offsets[内部主题: __consumer_offsets]
    O1[(GroupA<br>user-events-0<br>offset=500)]
    O2[(GroupA<br>user-events-1<br>offset=400)]
    O3[(GroupA<br>user-events-2<br>offset=600)]
    end
    end

    subgraph CG[Consumer Group: GroupA]
    direction LR
    C1[Consumer-1<br>订阅 Partition-0]
    C2[Consumer-2<br>订阅 Partition-1]
    C3[Consumer-3<br>订阅 Partition-2]
    end

    C1 — 拉取消息(offset=500~505) –> P0
    C2 — 拉取消息(offset=400~405) –> P1
    C3 — 拉取消息(offset=600~605) –> P2

    C1 -. 提交消费位移 .-> O1
    C2 -. 提交消费位移 .-> O2
    C3 -. 提交消费位移 .-> O3

    O1 -. offset 持久化 .-> __consumer_offsets
    O2 -. offset 持久化 .-> __consumer_offsets
    O3 -. offset 持久化 .-> __consumer_offsets

    Kafka事务

    🧩 一、为什么 Kafka 需要事务?

    Kafka 最早只能做到:

    • At-least-once(至少一次) → 不丢,但可能重复
    • At-most-once(至多一次) → 不重复,但可能丢

    而真正的业务(例如转账、扣库存)需要:

    “消息既不能重复,也不能丢”

    这就要求 Kafka 必须支持:

    Exactly Once Processing(EoP,精准一次处理)

    要实现 EoP,Kafka 必须解决一个原子性问题:

    “我从 input-topic 消费一条消息,处理后写到 output-topic,同时提交 offset”
    —— 要么全成功,要么全失败。

    Kafka事务消息核心应用场景:

  • 精确一次流处理
    • 场景:从A Topic消费,处理后再写入B Topic。
    • 作用:确保“处理结果输出”与“消费位移提交”原子同步,防止数据重复或丢失。这是Kafka Streams的基石。
  • 跨分区/主题原子写入
    • 场景:一个业务需要向多个Topic或分区发送多条消息。
    • 作用:保证这些消息要么全部对下游可见,要么全部不可见,避免产生部分更新的中间状态。
  • 数据库与缓存/搜索同步
    • 场景:应用在写数据库后,需要发消息通知缓存(如Redis)或搜索索引(如Elasticsearch)更新。
    • 作用:将“数据库提交”和“消息发送”包装成一个分布式事务,确保两者状态一致,避免数据不一致。
  • 金融与交易系统
    • 场景:转账、订单处理等。
    • 作用:确保扣款和入账、订单创建和库存扣减等关联操作的消息,作为一个原子单元执行,满足高一致性要求。

  • ⚙️ 二、事务能解决什么问题?

    Kafka 事务的核心目标是:

    把 “生产消息” 和 “提交 offset” 变成一个原子操作(Atomic Commit)。

    也就是说:

    • 不可能出现只写了消息但没提交 offset;
    • 也不可能出现提交了 offset 但消息没写入。

    这样可以防止由重复消费引发的副作用。

    🧱 三、Kafka 事务的基础:幂等性(Idempotence)

    Kafka 事务的底层是基于**幂等生产(Idempotent Producer)**的。

    在 Kafka 0.11+ 中,Producer 默认支持幂等:

    enable.idempotence=true

    这意味着即使发送重试,也不会产生重复消息。
    Broker 根据 Producer 的 PID(Producer ID) + Sequence Number 来判断是否是重复的消息。

    🧠 但幂等性只保证单分区单会话内不重复。
    想要跨分区、跨主题、跨重启仍然一致,就需要——事务(Transactions)。


    🧩 四、事务的核心概念

    Kafka 的事务机制引入了以下关键概念:

    名称含义
    PID(Producer ID) 唯一标识一个生产者实例,由 Broker 分配
    Transactional ID 持久化事务身份,允许生产者重启后继续事务
    Transaction Coordinator 事务协调者,管理事务的提交和中止
    Transaction Log 内部主题 __transaction_state,记录事务状态

    流程图

    #mermaid-svg-sAiHbqSoGXkV4YqU {font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}#mermaid-svg-sAiHbqSoGXkV4YqU .error-icon{fill:#552222;}#mermaid-svg-sAiHbqSoGXkV4YqU .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-sAiHbqSoGXkV4YqU .edge-thickness-normal{stroke-width:2px;}#mermaid-svg-sAiHbqSoGXkV4YqU .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-sAiHbqSoGXkV4YqU .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-sAiHbqSoGXkV4YqU .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-sAiHbqSoGXkV4YqU .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-sAiHbqSoGXkV4YqU .marker{fill:#333333;stroke:#333333;}#mermaid-svg-sAiHbqSoGXkV4YqU .marker.cross{stroke:#333333;}#mermaid-svg-sAiHbqSoGXkV4YqU svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-sAiHbqSoGXkV4YqU .actor{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:#ECECFF;}#mermaid-svg-sAiHbqSoGXkV4YqU text.actor>tspan{fill:black;stroke:none;}#mermaid-svg-sAiHbqSoGXkV4YqU .actor-line{stroke:grey;}#mermaid-svg-sAiHbqSoGXkV4YqU .messageLine0{stroke-width:1.5;stroke-dasharray:none;stroke:#333;}#mermaid-svg-sAiHbqSoGXkV4YqU .messageLine1{stroke-width:1.5;stroke-dasharray:2,2;stroke:#333;}#mermaid-svg-sAiHbqSoGXkV4YqU #arrowhead path{fill:#333;stroke:#333;}#mermaid-svg-sAiHbqSoGXkV4YqU .sequenceNumber{fill:white;}#mermaid-svg-sAiHbqSoGXkV4YqU #sequencenumber{fill:#333;}#mermaid-svg-sAiHbqSoGXkV4YqU #crosshead path{fill:#333;stroke:#333;}#mermaid-svg-sAiHbqSoGXkV4YqU .messageText{fill:#333;stroke:#333;}#mermaid-svg-sAiHbqSoGXkV4YqU .labelBox{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:#ECECFF;}#mermaid-svg-sAiHbqSoGXkV4YqU .labelText,#mermaid-svg-sAiHbqSoGXkV4YqU .labelText>tspan{fill:black;stroke:none;}#mermaid-svg-sAiHbqSoGXkV4YqU .loopText,#mermaid-svg-sAiHbqSoGXkV4YqU .loopText>tspan{fill:black;stroke:none;}#mermaid-svg-sAiHbqSoGXkV4YqU .loopLine{stroke-width:2px;stroke-dasharray:2,2;stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);}#mermaid-svg-sAiHbqSoGXkV4YqU .note{stroke:#aaaa33;fill:#fff5ad;}#mermaid-svg-sAiHbqSoGXkV4YqU .noteText,#mermaid-svg-sAiHbqSoGXkV4YqU .noteText>tspan{fill:black;stroke:none;}#mermaid-svg-sAiHbqSoGXkV4YqU .activation0{fill:#f4f4f4;stroke:#666;}#mermaid-svg-sAiHbqSoGXkV4YqU .activation1{fill:#f4f4f4;stroke:#666;}#mermaid-svg-sAiHbqSoGXkV4YqU .activation2{fill:#f4f4f4;stroke:#666;}#mermaid-svg-sAiHbqSoGXkV4YqU .actorPopupMenu{position:absolute;}#mermaid-svg-sAiHbqSoGXkV4YqU .actorPopupMenuPanel{position:absolute;fill:#ECECFF;box-shadow:0px 8px 16px 0px rgba(0,0,0,0.2);filter:drop-shadow(3px 5px 2px rgb(0 0 0 / 0.4));}#mermaid-svg-sAiHbqSoGXkV4YqU .actor-man line{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:#ECECFF;}#mermaid-svg-sAiHbqSoGXkV4YqU .actor-man circle,#mermaid-svg-sAiHbqSoGXkV4YqU line{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:#ECECFF;stroke-width:2px;}#mermaid-svg-sAiHbqSoGXkV4YqU :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}生产者(Producer)事务协调者(Transaction Coordinator)事务日志(Transaction Log)PID元数据(__producer_epochs)Leader1. 初始化事务INIT_PRODUCER_ID(transactional.id)查询/分配新的PID和Epoch返回PID+N, Epoch+M返回PID+N, Epoch+M生产者获得唯一身份PID+N (身份)Epoch+M (世代号)2. 开始事务beginTransaction()3. 发送数据ADD_PARTITIONS_TO_TXN(主题分区)记录分区到事务确认发送数据(携带PID, Epoch, 序列号)loop[对于每条消息]4. 提交事务 (两阶段)END_TXN(提交请求)写入PREPARE_COMMIT状态发送COMMIT提交标记确认写入COMPLETE_COMMIT状态事务提交成功生产者(Producer)事务协调者(Transaction Coordinator)事务日志(Transaction Log)PID元数据(__producer_epochs)Leader

    关键交互流程说明

  • 初始化事务 (initTransactions)
    • 生产者向协调器注册自己的transactional.id
    • 协调器分配新的(PID, Epoch)对,并持久化到__producer_epochs
    • 僵尸防护:任何旧的相同transactional.id的生产者都会因Epoch过低而被拒绝
  • 数据发送阶段
    • 生产者将数据发送到各分区Leader
    • 每条消息都携带(PID, Epoch, Sequence Number)实现幂等性
    • 协调器在事务日志中记录该事务涉及的分区
  • 事务提交(两阶段提交)
    • 阶段1:协调器在事务日志中写入PREPARE_COMMIT状态
    • 阶段2:向所有涉及的分区Leader发送COMMIT标记,使消息对消费者可见
    • 容错:如果协调器在阶段1后崩溃,新协调器能从事务日志恢复状态并继续提交
  • 🧠 五、Kafka 事务工作原理(流程)

    我们来完整看一遍事务从开始到提交的过程👇

    #mermaid-svg-OwZqVE03R3peLwFm {font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}#mermaid-svg-OwZqVE03R3peLwFm .error-icon{fill:#552222;}#mermaid-svg-OwZqVE03R3peLwFm .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-OwZqVE03R3peLwFm .edge-thickness-normal{stroke-width:2px;}#mermaid-svg-OwZqVE03R3peLwFm .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-OwZqVE03R3peLwFm .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-OwZqVE03R3peLwFm .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-OwZqVE03R3peLwFm .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-OwZqVE03R3peLwFm .marker{fill:#333333;stroke:#333333;}#mermaid-svg-OwZqVE03R3peLwFm .marker.cross{stroke:#333333;}#mermaid-svg-OwZqVE03R3peLwFm svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-OwZqVE03R3peLwFm .label{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;color:#333;}#mermaid-svg-OwZqVE03R3peLwFm .cluster-label text{fill:#333;}#mermaid-svg-OwZqVE03R3peLwFm .cluster-label span{color:#333;}#mermaid-svg-OwZqVE03R3peLwFm .label text,#mermaid-svg-OwZqVE03R3peLwFm span{fill:#333;color:#333;}#mermaid-svg-OwZqVE03R3peLwFm .node rect,#mermaid-svg-OwZqVE03R3peLwFm .node circle,#mermaid-svg-OwZqVE03R3peLwFm .node ellipse,#mermaid-svg-OwZqVE03R3peLwFm .node polygon,#mermaid-svg-OwZqVE03R3peLwFm .node path{fill:#ECECFF;stroke:#9370DB;stroke-width:1px;}#mermaid-svg-OwZqVE03R3peLwFm .node .label{text-align:center;}#mermaid-svg-OwZqVE03R3peLwFm .node.clickable{cursor:pointer;}#mermaid-svg-OwZqVE03R3peLwFm .arrowheadPath{fill:#333333;}#mermaid-svg-OwZqVE03R3peLwFm .edgePath .path{stroke:#333333;stroke-width:2.0px;}#mermaid-svg-OwZqVE03R3peLwFm .flowchart-link{stroke:#333333;fill:none;}#mermaid-svg-OwZqVE03R3peLwFm .edgeLabel{background-color:#e8e8e8;text-align:center;}#mermaid-svg-OwZqVE03R3peLwFm .edgeLabel rect{opacity:0.5;background-color:#e8e8e8;fill:#e8e8e8;}#mermaid-svg-OwZqVE03R3peLwFm .cluster rect{fill:#ffffde;stroke:#aaaa33;stroke-width:1px;}#mermaid-svg-OwZqVE03R3peLwFm .cluster text{fill:#333;}#mermaid-svg-OwZqVE03R3peLwFm .cluster span{color:#333;}#mermaid-svg-OwZqVE03R3peLwFm div.mermaidTooltip{position:absolute;text-align:center;max-width:200px;padding:2px;font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:12px;background:hsl(80, 100%, 96.2745098039%);border:1px solid #aaaa33;border-radius:2px;pointer-events:none;z-index:100;}#mermaid-svg-OwZqVE03R3peLwFm :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}否或出现异常消费者Poll一批消息Producer.beginTransaction循环处理本批消息producer.send输出消息到目标Topic所有消息处理/发送成功?producer.sendOffsetsToTransaction将消费位移纳入事务Producer.abortTransactionProducer.commitTransaction原子提交: 消息+位移事务成功消息可见 & 位移更新事务回滚消息不可见 & 位移不变

    六、代码相关

    // ===== 配置和初始化 =====
    // 生产者配置(启用事务)
    Properties producerProps = new Properties();
    producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
    producerProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
    producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
    producerProps.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "true"); // 启用幂等性
    producerProps.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "my-transactional-id-1"); // 事务ID

    KafkaProducer<String, String> producer = new KafkaProducer<>(producerProps);

    // 消费者配置(只读已提交的消息)
    Properties consumerProps = new Properties();
    consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
    consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, "my-exactly-once-app");
    consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
    consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
    consumerProps.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed"); // 只读已提交
    consumerProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); // 关闭自动提交

    KafkaConsumer<String, String> consumer = new KafkaConsumer<>(consumerProps);
    consumer.subscribe(Arrays.asList("input-topic"));

    // 初始化事务
    producer.initTransactions();

    // ===== 主处理循环 =====
    try {
    while (true) {
    // 步骤1: 消费者Poll一批消息
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));

    if (!records.isEmpty()) {
    try {
    // 步骤2: 开始事务
    producer.beginTransaction();

    // 步骤3: 处理并发送消息
    Map<TopicPartition, OffsetAndMetadata> offsetsToCommit = new HashMap<>();

    for (TopicPartition partition : records.partitions()) {
    List<ConsumerRecord<String, String>> partitionRecords = records.records(partition);

    for (ConsumerRecord<String, String> record : partitionRecords) {
    // 业务处理逻辑
    String processedValue = processMessage(record.value());

    // 发送到输出Topic
    producer.send(new ProducerRecord<>(
    "output-topic",
    record.key(),
    processedValue
    ));
    }

    // 计算该分区要提交的位移(最后一条消息的offset + 1)
    long lastOffset = partitionRecords.get(partitionRecords.size() 1).offset();
    offsetsToCommit.put(partition, new OffsetAndMetadata(lastOffset + 1));
    }

    // 步骤4: 将消费位移纳入事务
    producer.sendOffsetsToTransaction(offsetsToCommit, consumer.groupMetadata());

    // 步骤5: 提交事务 – 原子操作!
    producer.commitTransaction();

    System.out.println("事务提交成功: 处理 " + records.count() + " 条消息");

    } catch (Exception e) {
    // 步骤6: 中止事务(回滚)
    producer.abortTransaction();
    System.out.println("事务中止: " + e.getMessage());
    // 在实际应用中,这里可能需要重试或报警逻辑
    }
    }
    }
    } finally {
    producer.close();
    consumer.close();
    }

    // 业务处理函数
    private String processMessage(String value) {
    // 这里是你的实际业务逻辑
    // 例如: 数据转换、过滤、 enrichment等
    return value.toUpperCase() + "-PROCESSED";
    }

    赞(0)
    未经允许不得转载:171主机测评 » Kafka简介及其核心概念
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址