异步架构解耦与基础设施
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 将消息对象序列化为字符串
消息体:
| 消息类型 | 字段 | 说明 |
|---|---|---|
DocumentParseRouteMessage | documentId, taskId | 解析路由消息,taskId 关联任务记录 |
DocumentIndexBuildMessage | documentId, 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_text、embedding(vector 类型)外,还存储:
| 字段类别 | 字段 | 用途 |
|---|---|---|
| 溯源 | 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_json | JSONB 格式的完整元数据 |
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=180 | maxChars=800, overlap=120 |
SEMANTIC | 语义分割 | 内容质量较高的文档 | maxChars=1600, minChars=480 | maxChars=700, minChars=240 |
LLM | LLM 辅助分割 | 低质量文档 | maxChars=3500 | maxChars=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:
- 应用启动时自动扫描所有
ConsumerTaskBean - 为每个 Bean 创建对应的延迟消费者
- 自动启动监听线程
7.5 分布式租约
RedisLeaseManager 使用 Lua 脚本 实现原子化的租约操作:
| 操作 | Lua 脚本 | 说明 |
|---|---|---|
acquire | if redis.call('set', KEYS[1], ARGV[1], 'PX', ARGV[2], 'NX') then return 1 else return 0 end | 原子抢占 |
renew | if redis.call('get', KEYS[1]) == ARGV[1] then return redis.call('pexpire', KEYS[1], ARGV[2]) else return 0 end | 身份校验后续期 |
release | if redis.call('get', KEYS[1]) == ARGV[1] then return redis.call('del', KEYS[1]) else return 0 end | 身份校验后释放 |
适用于分布式场景下的 Leader 选举、任务互斥等场景。
7.6 数据操作封装
RedissonDataHandle 对 RBucket 的简化封装:
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 中定义了三个专用线程池:
| 线程池 | 核心线程 | 队列容量 | 用途 |
|---|---|---|---|
chatRagExecutorService | 8 | 256 | 核心 RAG 检索:子问题并发 + 通道并发 |
chatMemorySummaryExecutorService | 2 | 32 | 会话记忆摘要的异步压缩 |
chatPostProcessExecutorService | 2 | 64 | 对话后处理(推荐问题等) |
所有线程池使用 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
| 参数 | 默认值 | 说明 |
|---|---|---|
threads | 16 | Redisson 内部线程数 |
nettyThreads | 32 | Netty 事件循环线程数 |
corePoolSize | 可配置 | 自定义线程池核心数 |
maximumPoolSize | 可配置 | 自定义线程池最大数 |
connectTimeout | 1000ms | Redis 连接超时 |
十一、数据流总览
┌─────────────────────────────────────────────────────────┐
│ 文档处理全链路 │
├─────────────────────────────────────────────────────────┤
│ │
│ 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.java | Kafka 配置 |
异步处理
| 文件 | 职责 |
|---|---|
manage/service/DocumentAsyncProcessService.java | 异步处理接口 |
manage/service/impl/DocumentAsyncProcessServiceImpl.java | 核心实现(655 行):解析→切块→向量化全流程 |
manage/service/DocumentManageService.java | 文档管理门面(upload/confirmStrategy/buildIndex) |
MinIO
| 文件 | 职责 |
|---|---|
manage/service/impl/MinioDocumentStorageService.java | MinIO 存储实现(152 行) |
manage/service/DocumentStorageService.java | 存储服务接口 |
manage/config/DocumentManageMinioConfiguration.java | MinIO Client 配置 + bucket 自动创建 |
pgvector
| 文件 | 职责 |
|---|---|
manage/service/impl/DefaultDocumentVectorGateway.java | 向量网关:批量 embedding + UPSERT(279 行) |
manage/service/DocumentVectorGateway.java | 向量网关接口 |
manage/config/DocumentManagePgVectorConfiguration.java | pgvector 独立连接池配置 |
manage/support/DocumentPgVectorConstants.java | 表名常量 |
切块策略
| 文件 | 职责 |
|---|---|
manage/service/impl/DocumentStrategyServiceImpl.java | 多级切块管线:4 种策略 + 双流水线 |
文档解析
| 文件 | 职责 |
|---|---|
manage/service/impl/TikaDocumentParserService.java | Apache Tika 文档解析 |
Redisson 框架
| 文件 | 职责 |
|---|---|
super-agent-redisson-common-framework/.../RedissonCommonAutoConfiguration.java | Redisson 自动装配 + 线程池 |
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.java | Redis 租约管理器(Lua 脚本) |
线程池
| 文件 | 职责 |
|---|---|
chatagent/rag/config/ChatRagExecutorConfiguration.java | 3 个专用线程池 |
十三、设计亮点总结
- Kafka 双主题解耦:解析和索引分两个独立主题 + 独立消费者组,阶段分明,互不阻塞
- 同步发送保证可靠性:生产者使用
.get()阻塞确认,确保消息不丢失;以 documentId 为 key 保证同文档顺序 - MinIO 双路径存储:原始文件(带时间戳防重名)与解析文本(固定命名)分离,各司其职
- pgvector 独立连接池:与业务数据库 HikariCP 隔离,防止向量操作影响业务查询
- UPSERT 幂等写入:
ON CONFLICT DO UPDATE支持索引重建时重复执行不产生脏数据 - 批量向量化:每批 10 个 chunk 调用一次 EmbeddingModel,平衡 API 开销与内存占用
- 多级切块管线:PARENT(语义单元)+ CHILD(检索粒度)双流水线,4 种策略按文档质量自动推荐
- 注解驱动的分布式锁:
@ServiceLock用 SpEL 动态提取锁键,零代码侵入 - 三层幂等防护:Redis 标志位 → 本地锁 → Redisson 分布式公平锁,层层递进
- 延迟队列双线程池:listen 线程与 execute 线程分离,避免慢任务阻塞消息投递
- 租约 Lua 原子化:acquire/renew/release 全部使用 Lua 脚本,带身份校验防误释放
- 检索引擎全异步:CompletableFuture 并发执行子问题 + 双通道,双重超时控制防止长尾延迟
