Logo车 专
返回项目列表

异步架构解耦与基础设施

KafkaMinIOpgvectorRedisRedisson异步架构基础设施

基于 Kafka 异步处理文档解析、切块、向量化等耗时任务;MinIO 存储原始文档,pgvector 承载向量数据,Redis+Redisson 实现缓存与分布式限流,保障系统高并发稳定性

项目介绍

这是一个面向企业知识管理的AI智能体平台,支持多格式文档上传与解析、智能分块与向量化、知识库构建、多轮对话问答、ReActAgent 智能编排、RAG检索增强生成等全链路业务闭环。

技术栈:SpringBoot、SpringAI-Alibaba、MCP、Milvus、Kafka、MinIO、Redis、Redisson、RAG、Tavily等。

说明:本项目的向量存储生产环境采用 pgvector(PostgreSQL 扩展),Milvus 仅作为 ai-example 示例模块中的演示方案。


异步架构解耦与基础设施

基于 Kafka 异步处理文档解析、切块、向量化等耗时任务;MinIO 存储原始文档,pgvector 承载向量数据,Redis+Redisson 实现缓存与分布式限流,保障系统高并发稳定性。


一、架构全景

整个文档处理链路通过 Kafka 消息队列 实现完全异步解耦,用户上传文档后立即返回,后续的解析、切块、向量化全部在后台异步完成:

用户上传文档
  │
  ├─ ① 同步阶段(毫秒级返回)
  │   ├─ MinIO: 存储原始文件
  │   └─ MySQL: 写入文档记录
  │
  ├─ ② Kafka 消息:parse-topic ──────────────────────┐
  │                                                     │
  └─ ③ 异步消费:handleParseRoute()  ◄────────────────┘
      │
      ├─ MinIO: 下载原始文件
      ├─ Tika: 解析文本内容
      ├─ 提取文档结构节点(标题层级)
      ├─ MinIO: 上传解析后文本
      ├─ ES: 同步导航索引
      ├─ Neo4j: 同步图投影
      ├─ MySQL: 生成文档画像
      └─ 自动推荐切块策略 → 等待用户确认

用户确认策略
  │
  ├─ ④ Kafka 消息:index-topic ──────────────────────┐
  │                                                     │
  └─ ⑤ 异步消费:handleIndexBuild()  ◄────────────────┘
      │
      ├─ MinIO: 下载解析后文本
      ├─ 多级切块管线 (PARENT + CHILD 双流水线)
      ├─ MySQL: 保存 ParentBlock + Chunk
      ├─ EmbeddingModel: 批量向量化 (每批 10 个)
      ├─ pgvector: UPSERT 向量数据
      └─ ES: 关键词搜索索引

Redis + Redisson 作为横向基础设施贯穿全系统:

  • 分布式锁(@ServiceLock):防止并发冲突
  • 幂等控制(@RepeatExecuteLimit):防止重复提交
  • 延迟队列:定时任务调度
  • 数据缓存:热点数据加速
  • 租约管理:分布式协调

二、Kafka:文档处理异步管线

2.1 主题设计

DocumentManageProperties.Kafka 定义了两个主题,各自独立消费者组:

app.manage.kafka:
  parse-topic: super-agent-document-parse-route   # 解析路由
  index-topic: super-agent-document-index-build    # 索引构建
  group-id:    super-agent-document-manage

主题名通过 SpringUtil.getPrefixDistinctionName() 动态添加前缀,支持多环境隔离(如 dev-super-agent-document-parse-route)。

2.2 生产者

DocumentKafkaProducer 负责将文档处理任务投递到 Kafka:

public void sendParseRoute(DocumentParseRouteMessage message) {
    // 以 documentId 为 key 保证同一文档消息顺序
    String topic = prefix + "-" + properties.getKafka().getParseTopic();
    String payload = objectMapper.writeValueAsString(message);
    kafkaTemplate.send(topic, String.valueOf(message.getDocumentId()), payload).get();
    // .get() 同步阻塞,确保投递可靠性
}

public void sendIndexBuild(DocumentIndexBuildMessage message) {
    // 同理,planId 用于关联切块策略
    String topic = prefix + "-" + properties.getKafka().getIndexTopic();
    // ...
}

关键设计:

  • 同步发送.get() 阻塞等待确认,确保消息不丢失
  • documentId 作 Key:同一文档的消息路由到同一分区,保证顺序处理
  • JSON 序列化:通过 Jackson 将消息对象序列化为字符串

消息体:

消息类型字段说明
DocumentParseRouteMessagedocumentId, taskId解析路由消息,taskId 关联任务记录
DocumentIndexBuildMessagedocumentId, taskId, planId索引构建消息,planId 关联切块策略方案

2.3 消费者

DocumentKafkaConsumer 使用两个 @KafkaListener 分别监听两个主题:

@KafkaListener(topics = prefix + "-${app.manage.kafka.parse-topic}",
               groupId = "${app.manage.kafka.group-id}-parse")
public void consumeParseRoute(String payload) {
    DocumentParseRouteMessage message = objectMapper.readValue(payload, ...);
    asyncProcessService.handleParseRoute(message.getDocumentId(), message.getTaskId());
}

@KafkaListener(topics = prefix + "-${app.manage.kafka.index-topic}",
               groupId = "${app.manage.kafka.group-id}-index")
public void consumeIndexBuild(String payload) {
    DocumentIndexBuildMessage message = objectMapper.readValue(payload, ...);
    asyncProcessService.handleIndexBuild(message.getDocumentId(), message.getTaskId(), message.getPlanId());
}
  • 两个独立的消费者组(-parse-index),互不影响
  • 反序列化失败时仅记录日志,不会无限重试
  • 消费者直接委托给 DocumentAsyncProcessService 执行业务逻辑

三、异步处理核心:DocumentAsyncProcessServiceImpl

DocumentAsyncProcessServiceImpl(655 行)是整个异步管线的核心,实现了 DocumentAsyncProcessService 接口的两个方法。

3.1 handleParseRoute:解析路由

阶段一:CONTENT_PARSE 内容解析

1. 更新文档状态 → PARSING, 任务状态 → RUNNING
2. 从 MinIO 下载原始文件 (storageService.downloadObject)
3. Tika 解析文本 (parserService.parse)
   → DocumentAnalysisResult { parsedText, charCount, tokenCount,
       structureLevel, contentQualityLevel, structureNodes }
4. 上传解析后文本到 MinIO (storageService.uploadParsedText)
5. 替换文档结构节点 (structureNodeService.replaceDocumentNodes)
6. 同步导航产物:
   ├─ ES 导航索引 (navigationIndexService.reindexDocumentNodes)
   └─ Neo4j 图投影 (graphProjectionService.projectToGraph)
7. 生成文档画像 (documentProfileService.generateProfile)

阶段二:STRATEGY_ROUTE 策略路由

1. strategyService.recommendStrategy(document, analysisResult)
   → DocumentStrategyPlanDraft { parentSteps, childSteps, strategySnapshot, recommendReason }
2. 生成策略方案 ID,版本号递增
3. 持久化 Parent 管线步骤 + Child 管线步骤
4. 文档状态 → PARSE_SUCCESS, 策略状态 → RECOMMENDED
5. 等待用户在管理后台确认策略

3.2 handleIndexBuild:索引构建

阶段三:CHUNK_EXECUTE 切块执行

1. 从 MinIO 下载解析后文本 (storageService.downloadText)
2. strategyService.buildParentBlocks(document, plan, stepList, parsedText)
   → 执行 PARENT 管线 → 执行 CHILD 管线
   → List<ParentBlockCandidate> (每个包含 N 个 ChildChunk)

阶段四:CHUNK_POST_PROCESS 后处理

1. 过滤无效切块(空文本、无子切块的父块)
2. 构建实体对象: buildParentChildEntities()
   → 给每个 ParentBlock 和 Chunk 生成 UID
   → 设置 sectionPath, structureNodeId, canonicalPath, itemIndex
   → 估算 tokenCount
3. 批量 INSERT 到 MySQL

阶段五:VECTORIZE 向量化

1. vectorGateway.vectorize(chunkEntityList)
   → EmbeddingModel 批量生成向量(每批 10 个)
   → UPSERT 到 pgvector 表
2. keywordSearchGateway.indexChunks(chunkEntityList)
   → 写入 Elasticsearch 关键词索引
3. 更新 chunk 状态为 VECTOR_SUCCESS

阶段六:STORE_COMPLETE 完成归档

1. 更新策略方案状态 → EXECUTED
2. 文档索引状态 → BUILD_SUCCESS
3. 任务完成,记录总耗时

四、MinIO:文档对象存储

4.1 配置

app.manage.minio:
  endpoint: "http://127.0.0.1:9000"
  access-key: "minioadmin"
  secret-key: "minioadmin"
  bucket-name: "super-agent-document"
  object-prefix: "rag/document"           # 原始文件路径前缀
  parsed-text-prefix: "rag/parsed-text"   # 解析文本路径前缀

4.2 初始化

DocumentManageMinioConfiguration 创建 MinioClient Bean,并通过 CommandLineRunner 在启动时自动检查并创建 bucket。

4.3 存储操作

MinioDocumentStorageService 实现了 DocumentStorageService 接口,提供五个操作:

方法路径格式说明
uploadOriginalFile()rag/document/{documentId}/{timestamp}-{filename}上传原始文档,返回 URL
uploadParsedText()rag/parsed-text/{documentId}.txt上传 Tika 解析后的纯文本
downloadObject()按 objectName下载二进制,用于 Tika 解析
downloadText()按 objectName下载为 UTF-8 文本,用于切块
deleteObjects()批量先检查 bucket 是否存在,再逐个删除

关键设计:

  • 上传使用 PutObjectArgs.stream()-1 表示未知流长度
  • 每次上传前检查 bucket 是否存在,不存在则自动创建
  • buildObjectUrl() 直接拼接 endpoint + bucket + objectName 返回可公开访问的 URL
  • 删除前先 bucketExists() 判断,避免不必要的报错

五、pgvector:向量存储与检索

项目生产环境使用 pgvector(PostgreSQL 向量扩展),而非独立的 Milvus 服务。Milvus 仅在 ai-example 示例模块中作为演示。

5.1 独立连接池

DocumentManagePgVectorConfiguration 为 pgvector 创建独立于业务数据库的 HikariCP 连接池:

poolName: "super-agent-manage-pgvector-hikari"
maxPoolSize: 5, minIdle: 1
// 连接参数: stringtype=unspecified (兼容 pgvector 类型)

提供专用的 documentManagePgVectorJdbcTemplate Bean,使用 @Qualifier 注入。

5.2 向量表结构

public.super_agent_document_embedding 表包含 20+ 个字段,除了标准的 chunk_textembeddingvector 类型)外,还存储:

字段类别字段用途
溯源document_id, task_id, plan_id, parent_block_id追溯切块来源
结构section_path, structure_node_id, canonical_path, item_index文档结构信息
统计char_count, token_count内容统计
模型embedding_model向量模型溯源
元数据metadata_jsonJSONB 格式的完整元数据

5.3 批量向量化写入

DefaultDocumentVectorGateway 通过 Spring AI 的 EmbeddingModel 生成向量,批量写入 pgvector:

public static final int EMBEDDING_BATCH_SIZE_LIMIT = 10;

public void vectorize(List<SuperAgentDocumentChunk> chunkList) {
    EmbeddingModel embeddingModel = requireEmbeddingModel();

    for (int start = 0; start < validChunkList.size(); start += 10) {
        List<SuperAgentDocumentChunk> batch = validChunkList.subList(start, end);

        // 批量调用 EmbeddingModel
        List<float[]> embeddingList = embeddingModel.embed(
            batch.stream().map(SuperAgentDocumentChunk::getChunkText).toList());

        // UPSERT 带 ON CONFLICT 支持幂等写入
        batchUpsert(upsertSql, batch, embeddingList, embeddingModelName);
        markSuccess(batch);
    }
}

UPSERT SQL 使用 ON CONFLICT (id) DO UPDATE,支持重复执行不产生脏数据:

  • 向量通过 CAST(? AS vector) 写入
  • 元数据通过 CAST(? AS jsonb) 写入
  • 每次写入记录 embedding_model 字段实现模型溯源

5.4 向量检索

检索引擎 DocumentKnowledgeServiceImpl 中使用余弦距离检索:

SELECT ..., 1 - (embedding <=> CAST(? AS vector)) AS similarity_score
FROM public.super_agent_document_embedding
WHERE document_id = ? AND status = 1
ORDER BY embedding <=> CAST(? AS vector)
LIMIT ?

支持多维度过滤:文档 ID、任务 ID、章节路径、结构节点、规范路径、条目索引。


六、文档切块策略管线

DocumentStrategyServiceImpl 实现了多级切块管线,PARENT(大语义单元)和 CHILD(检索粒度单元)分别走不同的策略组合。

6.1 四种切块策略

策略类型码适用场景Parent 参数Child 参数
STRUCTURE按标题结构带目录层级的文档 (PDF/DOC/MD/HTML)基于标题标签分割基于标题标签分割
RECURSIVE递归分割通用兜底策略maxChars=2200, overlap=180maxChars=800, overlap=120
SEMANTIC语义分割内容质量较高的文档maxChars=1600, minChars=480maxChars=700, minChars=240
LLMLLM 辅助分割低质量文档maxChars=3500maxChars=3500

6.2 双管线架构

PARENT 流水线(大语义单元)
  STRUCTURE (按标题分割) → RECURSIVE (控制最大长度)
    ↓ 输出 ParentBlockCandidate 列表

CHILD 流水线(检索粒度单元)
  LLM/SEMANTIC (智能分句) → RECURSIVE (长度兜底)
    ↓ 输出 ChunkCandidate 列表(挂在对应 ParentBlock 下)

6.3 语义分割算法

// 1. 按句子切分(中英文标点:。!?!?;;.)
// 2. 提取每个句子的 Token 集合(英文单词 + 中文字符)
// 3. 计算相邻句子的 Jaccard 相似度
// 4. 当 similarity < 0.18 时触发语义断点
// 5. 同时检查最大字符数限制

6.4 递归分割算法

分割优先级: 段落(\n\n) → 行(\n) → 句子(标点) → 固定窗口
支持 overlap: Parent 默认 180 字符, Child 默认 120 字符

七、Redisson:分布式基础设施

7.1 全局配置

RedissonCommonAutoConfiguration 自动装配 RedissonClient

@Bean
public RedissonClient redissonClient(RedisProperties redisProperties, ...) {
    Config config = new Config();
    // 自动检测 SSL
    String prefix = isSsl ? "rediss://" : "redis://";
    config.useSingleServer()
        .setAddress(prefix + host + ":" + port)
        .setConnectTimeout(1000)
        .setDatabase(database)
        .setPassword(password);

    // 自定义线程池
    config.setThreads(16);          // Redisson 内部线程
    config.setNettyThreads(32);     // Netty 事件循环线程
    if (corePoolSize != null) {
        config.setExecutor(new ThreadPoolExecutor(
            corePoolSize, maximumPoolSize, keepAliveTime, ...));
    }
    return Redisson.create(config);
}

同时注册 RedissonDataHandle(数据缓存操作)、LocalLockCache(本地锁缓存)、LockInfoHandleFactory(锁信息工厂)三个 Bean。

7.2 分布式锁:@ServiceLock

通过 AOP 注解驱动,零侵入实现分布式锁:

@Target({ElementType.TYPE, ElementType.METHOD})
public @interface ServiceLock {
    LockType lockType() default LockType.Reentrant;  // 锁类型
    String name() default "";                         // 锁名称
    String[] keys();                                   // SpEL 表达式提取锁键
    long waitTime() default 10;                       // 等待时间
    TimeUnit timeUnit() default TimeUnit.SECONDS;
    LockTimeOutStrategy lockTimeoutStrategy() default FAIL;  // 超时策略
    String customLockTimeoutStrategy() default "";    // 自定义超时回调
}

使用示例:

@ServiceLock(name = "documentBuildIndex", keys = {"#documentId"},
             lockType = LockType.Reentrant, waitTime = 10)
public void buildIndex(Long documentId) {
    // 业务逻辑:同一文档同时只能有一个构建索引的线程执行
}

切面 ServiceLockAspect@Order(-10))通过 SpEL 从方法参数中动态提取锁键,支持四种锁类型:

  • Reentrant:可重入锁
  • Fair:公平锁
  • Read:读锁
  • Write:写锁

7.3 幂等控制:@RepeatExecuteLimit

防止用户在短时间内重复提交相同操作:

@Target({ElementType.TYPE, ElementType.METHOD})
public @interface RepeatExecuteLimit {
    String name() default "";
    String[] keys();                    // 幂等键
    long durationTime() default 0L;     // 禁重复时间窗口(秒)
    String message() default "提交频繁,请稍后重试";
}

三层防护机制:

┌─ 第一层:Redis 标志位 (RedissonDataHandle.get/set)
│   快速判断是否已存在幂等键
│
├─ 第二层:本地 JVM 锁 (LocalLockCache)
│   避免同一 JVM 内的并发竞争
│
└─ 第三层:Redisson 分布式公平锁
    跨实例的最终一致性保证

切面优先级 @Order(-11),在 @ServiceLock 之前执行。执行成功后设置 Redis 标志位,在 durationTime 秒内阻止重复执行。

7.4 延迟队列

基于 Redisson RDelayedQueue 实现的分布式延迟队列:

生产者 DelayProduceQueue

public void offer(T content, long delayTime, TimeUnit timeUnit) {
    RDelayedQueue<T> delayedQueue = redissonClient.getDelayedQueue(blockingQueue);
    delayedQueue.offer(content, delayTime, timeUnit);
}

消费者 DelayConsumerQueue

  • listenStartThreadPool(1 线程):从 RBlockingQueue 阻塞获取到期消息
  • executeTaskThreadPool(可配置):实际执行业务逻辑 ConsumerTask.execute()
  • 支持 isolationRegionCount,每个主题可创建多个隔离分区

启动机制 DelayQueueInitHandler

  • 应用启动时自动扫描所有 ConsumerTask Bean
  • 为每个 Bean 创建对应的延迟消费者
  • 自动启动监听线程

7.5 分布式租约

RedisLeaseManager 使用 Lua 脚本 实现原子化的租约操作:

操作Lua 脚本说明
acquireif redis.call('set', KEYS[1], ARGV[1], 'PX', ARGV[2], 'NX') then return 1 else return 0 end原子抢占
renewif redis.call('get', KEYS[1]) == ARGV[1] then return redis.call('pexpire', KEYS[1], ARGV[2]) else return 0 end身份校验后续期
releaseif redis.call('get', KEYS[1]) == ARGV[1] then return redis.call('del', KEYS[1]) else return 0 end身份校验后释放

适用于分布式场景下的 Leader 选举、任务互斥等场景。

7.6 数据操作封装

RedissonDataHandleRBucket 的简化封装:

public <T> T get(String key);
public <T> void set(String key, T value);
public <T> void set(String key, T value, long time, TimeUnit timeUnit);

支持 SECONDS, MINUTES, HOURS, DAYS 四种 TTL 单位。


八、检索引擎线程池架构

ChatRagExecutorConfiguration 中定义了三个专用线程池:

线程池核心线程队列容量用途
chatRagExecutorService8256核心 RAG 检索:子问题并发 + 通道并发
chatMemorySummaryExecutorService232会话记忆摘要的异步压缩
chatPostProcessExecutorService264对话后处理(推荐问题等)

所有线程池使用 CallerRunsPolicy 饱和策略,防止任务堆积导致 OOM。


九、RAG 检索时的异步并行架构

RagRetrievalEngine 在执行检索时充分利用线程池实现并行:

RAG 检索
  │
  ├─ 子问题 1 ──┐
  │   ├─ VectorRetrievalChannel  ──┐
  │   └─ KeywordRetrievalChannel ──┤  CompletableFuture 并行
  │                                 │
  ├─ 子问题 2 ──┐                  │
  │   ├─ VectorRetrievalChannel  ──┤
  │   └─ KeywordRetrievalChannel ──┤
  │                                 │
  └─ 子问题 N ──┘                  │
                                    │
  CompletableFuture.allOf() ◄───────┘
      │
      ├─ 证据闸门 (vector 相似度 / keyword 相对分)
      ├─ RRF 融合 (K=60)
      ├─ Parent 块提升
      ├─ HTTP Rerank 精排
      └─ TopK 截断
  • 子问题并发CompletableFuture.supplyAsync() 提交到 chatRagExecutorService
  • 通道并发:每个子问题内部,向量和关键词通道同时发起
  • 超时控制:子问题级 subQuestionTimeoutMs + 通道级 channelTimeoutMs 双重超时

十、完整配置参数

10.1 Kafka

app.manage.kafka:
  parse-topic: super-agent-document-parse-route
  index-topic: super-agent-document-index-build
  group-id: super-agent-document-manage

10.2 MinIO

app.manage.minio:
  endpoint: "http://127.0.0.1:9000"
  access-key: "minioadmin"
  secret-key: "minioadmin"
  bucket-name: "super-agent-document"
  object-prefix: "rag/document"
  parsed-text-prefix: "rag/parsed-text"

10.3 pgvector

app.manage.pg-vector:
  host: 127.0.0.1
  port: 5432
  database: super_agent_pgvector
  schema: public
  pool-name: super-agent-manage-pgvector-hikari
  max-pool-size: 5
  min-idle: 1

10.4 Redisson

参数默认值说明
threads16Redisson 内部线程数
nettyThreads32Netty 事件循环线程数
corePoolSize可配置自定义线程池核心数
maximumPoolSize可配置自定义线程池最大数
connectTimeout1000msRedis 连接超时

十一、数据流总览

┌─────────────────────────────────────────────────────────┐
│                    文档处理全链路                          │
├─────────────────────────────────────────────────────────┤
│                                                         │
│  upload()         kafka.send          kafka.consume      │
│  ─────────► MinIO ────────► parseRoute ────────►         │
│  │          存储原始       消息投递        消费解析        │
│  │                                                       │
│  │  ┌───────── 异步解析管线 ─────────┐                   │
│  │  │ MinIO下载 → Tika解析 → 结构提取 │                  │
│  │  │   → MinIO上传文本 → ES/Neo4j   │                   │
│  │  │   → 策略推荐 → 等待确认         │                   │
│  │  └──────────────────────────────┘                    │
│  │                                                       │
│  └─ confirmStrategy() ──► kafka.send ──► kafka.consume  │
│                           indexBuild     消费索引构建     │
│                                                         │
│      ┌───────── 异步索引构建管线 ────────┐               │
│      │ MinIO下载 → 多级切块(PARENT+CHILD)│               │
│      │  → MySQL → EmbeddingModel批量向量 │               │
│      │  → pgvector UPSERT → ES关键词索引 │               │
│      └──────────────────────────────────┘               │
│                                                         │
├─────────────────────────────────────────────────────────┤
│                     RAG 检索时                            │
├─────────────────────────────────────────────────────────┤
│                                                         │
│  用户提问 → EmbeddingModel → pgvector 向量检索            │
│           └─ ES 关键词检索 ─┤                            │
│                             ├─ RRF融合 → Rerank → TopK   │
│           Redis 缓存 ◄──────┘                            │
│           Redisson 限流                                   │
│                                                         │
└─────────────────────────────────────────────────────────┘

十二、完整文件清单

Kafka 消息

文件职责
manage/mq/DocumentKafkaProducer.java生产者:发送解析/索引消息
manage/mq/DocumentKafkaConsumer.java消费者:两个 @KafkaListener
manage/mq/message/DocumentParseRouteMessage.java解析路由消息体
manage/mq/message/DocumentIndexBuildMessage.java索引构建消息体
manage/config/DocumentManageKafkaConfiguration.javaKafka 配置

异步处理

文件职责
manage/service/DocumentAsyncProcessService.java异步处理接口
manage/service/impl/DocumentAsyncProcessServiceImpl.java核心实现(655 行):解析→切块→向量化全流程
manage/service/DocumentManageService.java文档管理门面(upload/confirmStrategy/buildIndex)

MinIO

文件职责
manage/service/impl/MinioDocumentStorageService.javaMinIO 存储实现(152 行)
manage/service/DocumentStorageService.java存储服务接口
manage/config/DocumentManageMinioConfiguration.javaMinIO Client 配置 + bucket 自动创建

pgvector

文件职责
manage/service/impl/DefaultDocumentVectorGateway.java向量网关:批量 embedding + UPSERT(279 行)
manage/service/DocumentVectorGateway.java向量网关接口
manage/config/DocumentManagePgVectorConfiguration.javapgvector 独立连接池配置
manage/support/DocumentPgVectorConstants.java表名常量

切块策略

文件职责
manage/service/impl/DocumentStrategyServiceImpl.java多级切块管线:4 种策略 + 双流水线

文档解析

文件职责
manage/service/impl/TikaDocumentParserService.javaApache Tika 文档解析

Redisson 框架

文件职责
super-agent-redisson-common-framework/.../RedissonCommonAutoConfiguration.javaRedisson 自动装配 + 线程池
super-agent-redisson-common-framework/.../RedissonDataHandle.java数据缓存封装
super-agent-redisson-common-framework/.../LocalLockCache.java本地锁缓存
super-agent-service-lock-framework/.../ServiceLock.java分布式锁注解
super-agent-service-lock-framework/.../ServiceLockAspect.java锁 AOP 切面 (@Order(-10))
super-agent-repeat-execute-limit-framework/.../RepeatExecuteLimit.java幂等控制注解
super-agent-repeat-execute-limit-framework/.../RepeatExecuteLimitAspect.java幂等 AOP 切面 (@Order(-11))
super-agent-service-delay-queue-framework/.../DelayProduceQueue.java延迟队列生产者
super-agent-service-delay-queue-framework/.../DelayConsumerQueue.java延迟队列消费者
super-agent-service-lease-framework/.../RedisLeaseManager.javaRedis 租约管理器(Lua 脚本)

线程池

文件职责
chatagent/rag/config/ChatRagExecutorConfiguration.java3 个专用线程池

十三、设计亮点总结

  1. Kafka 双主题解耦:解析和索引分两个独立主题 + 独立消费者组,阶段分明,互不阻塞
  2. 同步发送保证可靠性:生产者使用 .get() 阻塞确认,确保消息不丢失;以 documentId 为 key 保证同文档顺序
  3. MinIO 双路径存储:原始文件(带时间戳防重名)与解析文本(固定命名)分离,各司其职
  4. pgvector 独立连接池:与业务数据库 HikariCP 隔离,防止向量操作影响业务查询
  5. UPSERT 幂等写入ON CONFLICT DO UPDATE 支持索引重建时重复执行不产生脏数据
  6. 批量向量化:每批 10 个 chunk 调用一次 EmbeddingModel,平衡 API 开销与内存占用
  7. 多级切块管线:PARENT(语义单元)+ CHILD(检索粒度)双流水线,4 种策略按文档质量自动推荐
  8. 注解驱动的分布式锁@ServiceLock 用 SpEL 动态提取锁键,零代码侵入
  9. 三层幂等防护:Redis 标志位 → 本地锁 → Redisson 分布式公平锁,层层递进
  10. 延迟队列双线程池:listen 线程与 execute 线程分离,避免慢任务阻塞消息投递
  11. 租约 Lua 原子化:acquire/renew/release 全部使用 Lua 脚本,带身份校验防误释放
  12. 检索引擎全异步:CompletableFuture 并发执行子问题 + 双通道,双重超时控制防止长尾延迟
返回项目列表