Spring Boot 4.1 | Kafka 4.0 | Share Groups (KIP-932) | 事务消息 | RAG 独立模块 | PgVector + Redis

在之前的博文中,我们介绍了 spring-cloud-samples 项目的 AI 智能体、gRPC 通信和 Seata 分布式事务三大模块。

时隔数月,项目迎来了又一轮重要升级

  • 📨 cloud-kafka-sample —— 基于 Kafka 4.0 的全新消息模块,演示 Share Groups (KIP-932)、事务消息、3 节点集群高可用
  • 🔄 cloud-stream-sample —— Spring Cloud Stream + RocketMQ 6 大场景深度演示,从基础消费到延迟/顺序/事务消息
  • 🧠 cloud-ai-rag-sample —— RAG 检索增强生成独立模块化,支持 PgVector + Redis 双向量库切换

下面逐一详解。


📨 模块一:cloud-kafka-sample —— Kafka 4.0 Share Groups 与事务消息实战

为什么要加 Kafka 模块?

Kafka 作为分布式消息系统的标杆,在微服务架构中扮演着至关重要的角色。然而:

  • Kafka 4.0 引入了 Share Groups (KIP-932),突破了传统消费者组"一个分区只能被一个消费者消费"的限制
  • 事务消息是保证分布式系统数据一致性的关键能力,但实际项目中很少看到完整的示例
  • 大多数 Kafka 示例只演示基础收发,缺乏生产级的高级特性演示

cloud-kafka-sample 正是为了填补这些空白。

技术栈

组件 版本
Spring Boot 4.1.0
Kafka 4.3.1 (3 节点集群)
端口 8768

功能全景

1️⃣ 传统消费者组 —— 基础消息收发

最经典的 Kafka 使用场景:一个生产者,一个消费者组,消息被组内消费者独占消费。

@Component
public class Consumer {
    @KafkaListener(topics = "${app.kafka.topic}")
    void processMessage(SampleMessage message) {
        log.info("Received sample message [{}]", message);
    }
}

启动时自动发送一条测试消息:

@Bean
ApplicationRunner runner(Producer producer) {
    return args -> producer.sendTraditional();
}

验证方式:

curl -X POST http://localhost:8768/kafka/traditional
2️⃣ Share Groups (KIP-932) —— 突破分区限制的并行消费

这是 Kafka 4.0 最重要的新特性之一。传统消费者组中,一个分区同一时刻只能被一个消费者消费;而 Share Groups 允许多个消费者从同一个分区并行消费不同的消息,非常适合任务分发和工作队列场景。

项目演示了两种确认模式:

隐式确认模式 —— 方法正常返回自动 ACCEPT,抛出异常自动 REJECT:

@KafkaListener(
    topics = "${app.kafka.share-topic}",
    containerFactory = "implicitShareKafkaListenerContainerFactory",
    groupId = "${app.kafka.share-group}"
)
void processImplicit(ConsumerRecord<String, SampleMessage> record) {
    log.info("[Share-Implicit] Received: {} from partition {} offset {}",
            record.value(), record.partition(), record.offset());
}

显式确认模式 —— 手动控制每条消息的确认,支持三种状态:

@KafkaListener(
    topics = "${app.kafka.share-topic-explicit}",
    containerFactory = "explicitShareKafkaListenerContainerFactory",
    groupId = "${app.kafka.share-group}",
    concurrency = "5"
)
void processExplicit(ConsumerRecord<String, SampleMessage> record, ShareAcknowledgment ack) {
    try {
        if (record.value().getId() % 5 == 0) {
            ack.release();  // 临时失败,消息将被重新投递
            return;
        }
        ack.acknowledge();  // 处理成功
    } catch (Exception e) {
        ack.reject();  // 永久失败,不再重试
    }
}

配置类的关键在于使用 ShareKafkaListenerContainerFactoryDefaultShareConsumerFactory

@Bean
public ShareKafkaListenerContainerFactory<String, SampleMessage>
        explicitShareKafkaListenerContainerFactory(
        ShareConsumerFactory<String, SampleMessage> factory) {
    var cf = new ShareKafkaListenerContainerFactory<>(factory);
    cf.getContainerProperties()
      .setShareAckMode(ContainerProperties.ShareAckMode.MANUAL);
    return cf;
}

验证方式:

# 发送 10 条 Share Group 隐式确认消息
curl -X POST "http://localhost:8768/kafka/share/implicit?count=10"

# 发送 10 条显式确认消息(id%5==0 的消息会被 release 重投递)
curl -X POST "http://localhost:8768/kafka/share/explicit?count=10"

这意味着:Share Groups 让 Kafka 从"分区独占"进化为"分区共享",大幅提升了消费并发度,特别适合任务分发场景。

3️⃣ 事务消息 —— 原子性发送与回滚

事务消息保证多条消息的发送是原子性的:要么全部成功对消费者可见,要么全部不可见。

生产者配置 —— 显式创建事务和非事务两个 KafkaTemplate:

@Bean
@Qualifier("txKafkaTemplate")
public KafkaTemplate<Object, SampleMessage> txKafkaTemplate() {
    Map<String, Object> props = new HashMap<>();
    props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "kafka-sample-tx");
    ProducerFactory<Object, SampleMessage> pf =
        new DefaultKafkaProducerFactory<>(props);
    return new KafkaTemplate<>(pf);
}

发送事务消息 —— 使用 executeInTransaction 保证原子性:

public void sendTransactional(int count, boolean commit) {
    try {
        this.txKafkaTemplate.executeInTransaction(operations -> {
            for (int i = 1; i <= count; i++) {
                operations.send(this.txTopic,
                    new SampleMessage(i, "tx-task-" + i));
            }
            if (!commit) {
                throw new RuntimeException("Simulated rollback");
            }
            return null;
        });
    } catch (Exception e) {
        log.warn("[TX] Transaction rolled back");
    }
}

消费者配置 —— 使用 read_committed 隔离级别:

props.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed");

验证方式:

# 提交事务 —— 消费者收到消息
curl -X POST "http://localhost:8768/kafka/tx/commit?count=5"

# 回滚事务 —— 消费者收不到消息
curl -X POST "http://localhost:8768/kafka/tx/rollback?count=5"

这意味着:事务消息让 Kafka 从"至少一次"进化为"精确一次"语义,保证分布式系统的数据一致性。

4️⃣ 3 节点 Kafka 集群 —— KRaft 高可用部署

项目手动部署 3 节点 Kafka 集群(KRaft 模式,无 ZooKeeper):下载 Kafka 安装包,准备 3 份配置文件(server-1/2/3.properties),分别格式化存储后在独立终端启动,生产者配置指向 3 个 Broker:

spring:
  kafka:
    bootstrap-servers: localhost:9092,localhost:9094,localhost:9096

这意味着:3 节点集群演示了真实生产环境的高可用部署,任何一个 Broker 宕机不影响消息收发。


🔄 模块二:cloud-stream-sample —— Spring Cloud Stream 6 大场景

为什么 Stream 模块值得单独介绍?

之前的博文只简单提到了 Stream 模块支持 RocketMQ 消息驱动,但实际上它演示了 6 个完整场景,覆盖了消息中间件最常用的高级特性。

技术栈

组件 版本
Spring Boot 4.1.0
Spring Cloud Stream 4.1.0
RocketMQ 5.x
端口 8767

6 大场景

场景 1:基础消费 —— Consumer 函数
@Bean
public Consumer<String> input() {
    return message -> log.info("Received message: {}", message);
}
场景 2:定时消息源 —— Supplier 函数

Supplier 每隔 1 秒自动生成一条消息:

@Bean
public Supplier<String> output2() {
    return () -> "你好";
}
output2-out-0:
  destination: stream-demo-topic2
  producer:
    poller:
      fixed-delay: 1000  # 每秒发送一次
场景 3:消息处理管道 —— Function 函数

Function 接收消息、转换后输出到新 Topic:

@Bean
public Function<String, String> transform() {
    return message -> "[PROCESSED] " + message.toUpperCase();
}
curl -X POST "http://localhost:8767/stream/send?message=hello"
# 日志:消息转换: hello -> [PROCESSED] HELLO
场景 4:延迟消息 —— RocketMQ 延迟级别

通过 Message Header 设置延迟级别(1-18),消息在指定延迟后被消费:

Message<String> msg = MessageBuilder.withPayload(message)
        .setHeader("DELAY", delayLevel)
        .build();
streamBridge.send("delayPublish-out-0", msg);
curl -X POST "http://localhost:8767/stream/delay?message=hello&delayLevel=2"
# 5 秒后消费者收到消息
场景 5:顺序消息 —— 分区键 + 顺序消费
Message<String> msg = MessageBuilder.withPayload(message)
        .setHeader("ORDER_KEY", orderKey)
        .build();
streamBridge.send("fifoPublish-out-0", msg);
fifo-in-0:
  consumer:
    orderly: true
fifoPublish-out-0:
  producer:
    orderly: true
    partitionKeyExpression: headers['ORDER_KEY']
curl -X POST "http://localhost:8767/stream/fifo?message=msg1&orderKey=order-1"
curl -X POST "http://localhost:8767/stream/fifo?message=msg2&orderKey=order-1"
# 消息按发送顺序被消费
场景 6:事务消息 —— RocketMQ 两阶段提交

事务消息流程:发送半消息 -> 执行本地事务 -> 根据结果 Commit/Rollback。

@Component("demoTransactionListener")
public class DemoTransactionListener implements TransactionListener {
    @Override
    public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
        String txArg = msg.getProperty("TX_ARG");
        if ("commit".equalsIgnoreCase(txArg))
            return LocalTransactionState.COMMIT_MESSAGE;
        if ("rollback".equalsIgnoreCase(txArg))
            return LocalTransactionState.ROLLBACK_MESSAGE;
        return ThreadLocalRandom.current().nextBoolean()
            ? LocalTransactionState.COMMIT_MESSAGE
            : LocalTransactionState.ROLLBACK_MESSAGE;
    }
}
# 提交事务 —— 消费者收到消息
curl -X POST "http://localhost:8767/stream/tx?message=hello&arg=commit"

# 回滚事务 —— 消费者收不到消息
curl -X POST "http://localhost:8767/stream/tx?message=hello&arg=rollback"

这意味着:Stream 模块覆盖了延迟、顺序、事务等生产环境最常用的 6 大场景,每个场景都有完整的代码和验证方式。


🧠 模块三:cloud-ai-rag-sample —— RAG 独立模块化

为什么 RAG 要独立成模块?

在之前的博文中,RAG 功能集成在 cloud-ai-sample 中。但随着 RAG 在企业级应用中的广泛落地,它已经足够复杂和重要,值得独立成一个模块:

  • 向量数据库选型:支持 PgVector 和 Redis 两种主流方案,通过 Profile 切换
  • 知识库管理:文档摄入、检索、删除的完整生命周期
  • 独立部署:RAG 服务可以独立扩展

技术栈

组件 版本
Spring Boot 4.1.0
Spring AI 2.0.0
向量数据库 PostgreSQL + pgvector 或 Redis
Embedding 模型 text-embedding-v3 (1024 维)
端口 8889

核心功能

1️⃣ 文档摄入 —— 自动分块与向量化
public int ingest(String content, String source) {
    Document doc = new Document(content);
    doc.getMetadata().put("source", source);
    List<Document> chunks = textSplitter.split(doc);
    vectorStore.add(chunks);
    return chunks.size();
}
curl -X POST http://localhost:8889/ai/rag/ingest \
  -H "Content-Type: application/json" \
  -d '{"content":"Spring AI 是 Spring 生态的 AI 集成框架...","source":"docs"}'
2️⃣ RAG 查询 —— 检索增强 LLM 回答
public String query(String question, int topK) {
    List<Document> docs = vectorStore.similaritySearch(
        SearchRequest.builder().query(question).topK(topK).build());
    String context = docs.stream()
        .map(Document::getText).collect(joining("\n"));
    String prompt = "基于以下参考资料回答:\n" + context
        + "\n\n问题:" + question;
    return chatClient.prompt().user(prompt).call().content();
}
curl "http://localhost:8889/ai/rag/query?question=Spring AI&topK=3"
3️⃣ 双向量库 —— PgVector 与 Redis 一键切换
# 使用 PgVector(默认)
java -jar cloud-ai-rag-sample.jar

# 使用 Redis
java -jar cloud-ai-rag-sample.jar --spring.profiles.active=redis

这意味着:RAG 模块独立部署、双向量库支持,企业可根据现有技术栈灵活选型。


📊 模块总览:16 个模块

模块 端口 说明 状态
cloud-gateway-sample 8764 Gateway + Sentinel 限流 原有
cloud-provider-sample 8765 Web 服务提供者 原有
cloud-consumer-sample 8766 Web 服务消费者 原有
cloud-provider-reactive-sample 8762 Reactive Web 提供者 原有
cloud-consumer-reactive-sample 8763 Reactive Web 消费者 原有
cloud-provider-dubbo-sample 50051 Dubbo RPC 提供者 原有
cloud-nacos-config-sample 8761 Nacos 动态配置 原有
cloud-stream-sample 8767 Stream(6 大场景) 🆕 增强
cloud-ai-sample 8888 Spring AI 2.0 全场景 原有
cloud-grpc-server-sample 9090 gRPC Server 原有
cloud-seata-sample 18081-84 Seata 分布式事务 原有
cloud-kafka-sample 8768 Kafka 4.0 🆕 新增
cloud-ai-rag-sample 8889 RAG(PgVector/Redis) 🆕 独立
cloud-nacos-discovery-sample 8760 Nacos 服务发现 原有
cloud-sample-api - 接口定义 + Proto 原有
cloud-commons - gRPC 服务发现桥接 原有

🎓 学习价值

cloud-kafka-sample

  • Kafka 4.0 Share Groups (KIP-932) 隐式/显式确认
  • 事务消息原子性发送与 read_committed 隔离
  • 3 节点 Kafka 集群(KRaft 模式)高可用
  • 非事务与事务 KafkaTemplate 分离设计

cloud-stream-sample

  • 函数式编程模型(Consumer/Function/Supplier)
  • 延迟消息(RocketMQ 18 级延迟)
  • 顺序消息(分区键 + orderly)
  • 事务消息(两阶段提交)
  • StreamBridge 编程式发布

cloud-ai-rag-sample

  • RAG 完整流程(摄入 -> 分块 -> 向量化 -> 检索 -> 增强)
  • PgVector 与 Redis 双向量库 Profile 切换
  • 基于元数据的文档过滤与删除

🚦 快速体验

Kafka 模块

# 1. 手动部署 3 节点 Kafka 集群(KRaft 模式)
# 下载 Kafka 安装包,参考 cloud-kafka-sample/README.md 准备 3 份配置并依次启动
KAFKA_CLUSTER_ID="$(bin/kafka-storage.sh random-uuid)"
bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c config/server-1.properties
bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c config/server-2.properties
bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c config/server-3.properties
bin/kafka-server-start.sh config/server-1.properties  # 终端1
bin/kafka-server-start.sh config/server-2.properties  # 终端2
bin/kafka-server-start.sh config/server-3.properties  # 终端3

# 2. 启动 Kafka 模块
cd cloud-kafka-sample && mvn spring-boot:run

# 3. 验证
curl -X POST http://localhost:8768/kafka/traditional
curl -X POST "http://localhost:8768/kafka/share/implicit?count=10"
curl -X POST "http://localhost:8768/kafka/tx/commit?count=5"
curl -X POST "http://localhost:8768/kafka/tx/rollback?count=5"

Stream 模块

# 1. 启动 RocketMQ
bin/mqnamesrv && bin/mqbroker -n localhost:9876

# 2. 启动 Stream 模块
cd cloud-stream-sample && mvn spring-boot:run

# 3. 验证 6 大场景
curl -X POST "http://localhost:8767/stream/send?message=hello"
curl -X POST "http://localhost:8767/stream/delay?message=hello&delayLevel=2"
curl -X POST "http://localhost:8767/stream/fifo?message=msg1&orderKey=order-1"
curl -X POST "http://localhost:8767/stream/tx?message=hello&arg=commit"

RAG 模块

# 1. 启动 PostgreSQL + pgvector
psql -U postgres -c "CREATE DATABASE ai_demo;"
psql -U postgres -d ai_demo -f init_ai_demo.sql

# 2. 启动 RAG 模块
cd cloud-ai-rag-sample && mvn spring-boot:run

# 3. 验证
curl -X POST http://localhost:8889/ai/rag/ingest \
  -H "Content-Type: application/json" \
  -d '{"content":"Spring AI 支持 Chat 和 Embedding...","source":"docs"}'
curl "http://localhost:8889/ai/rag/query?question=Spring AI&topK=3"

💡 最佳实践总结

1. Kafka Share Groups vs 传统消费者组

特性 传统消费者组 Share Groups
分区消费 一个分区只能被一个消费者消费 多个消费者可并行消费同一分区
并发度 受限于分区数 不受分区数限制
适用场景 流式处理、事件溯源 任务分发、工作队列
确认模式 自动/手动 offset 提交 隐式/显式 ACK

2. Kafka 事务消息关键配置

  • 生产者必须配置 transactional.id
  • 消费者必须设置 isolation.level=read_committed
  • 使用 executeInTransaction() 保证原子性
  • 非事务和事务 KafkaTemplate 应分离

3. Spring Cloud Stream 函数式模型

  • Consumer:终端消费者,处理消息后不输出
  • Function:消息处理管道,转换后输出到新 Topic
  • Supplier:定时消息源,自动生成消息
  • StreamBridge:编程式消息发布

4. RAG 的 temperature 调优

  • RAG 场景建议 temperature=0.2,忠于检索事实
  • 通用对话可用 temperature=0.7,发挥创造力

🔗 相关链接


📝 结语

从最初的 10 个模块到如今的 16 个,spring-cloud-samples 始终围绕一个目标:提供生产环境可参考的微服务全栈示例

本次新增的三大模块,分别对应了消息驱动和 AI 工程化领域最核心的能力:

  • Kafka 4.0 —— Share Groups 突破分区限制,事务消息保证数据一致性
  • Stream 6 大场景 —— 从基础消费到延迟/顺序/事务,覆盖生产环境最常用的消息模式
  • RAG 独立模块 —— 双向量库支持,企业级知识库检索增强生成方案

如果你正在:

  • 📨 评估 Kafka 4.0 的新特性(Share Groups、事务消息)
  • 🔄 学习 Spring Cloud Stream 的高级用法
  • 🧠 构建企业级 RAG 应用
  • 📚 寻找一站式 Spring Cloud 学习项目

Star ⭐ 这个项目,三大模块即刻上手!

git clone https://github.com/javahongxi/spring-cloud-samples.git
cd spring-cloud-samples

© hongxi.org | 以生产环境可参考为目标,持续打造最完整的 Spring Cloud 示例项目

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐