Appearance
完整项目(一):灌库管道与数据源接入
1. 本节产出
企业知识库项目的灌库侧:支持本地上传、S3/OSS 拉取、Confluence/ Wiki 定时同步三种数据源,统一走「解析 → 切分 → 元数据 → 幂等写入」管道,并有失败队列与重投机制。
2. 前置依赖
3. 为什么数据源接入是项目的第一个坑
前面各节的示例都是「本地上传一个 PDF」。真实项目里:
| 数据源 | 特点 | 坑 |
|---|---|---|
| 本地上传 | 最简单 | 用户重复上传、超大文件 |
| 对象存储(S3/OSS) | 批量、可脚本化 | 权限、大文件、断点续传 |
| Confluence / Wiki | 持续变化 | 增量同步困难、HTML 结构复杂 |
| 数据库 / 工单系统 | 结构化内容 | 需要转成文本,字段映射 |
| 网盘同步 | 文件随时变 | 删除检测(源里删了,库里要同步删) |
最后一行是真正的难点:文档在源端被删除时,你的库里也必须删掉。大多数团队的增量同步只做了「新增和更新」,没做「删除」,导致用户搜到早就被废弃的文档。
4. 核心原理
4.1 统一管道抽象
任意数据源
│
▼ SourceConnector(拉取 + 列举 + 探测变更)
DocumentRef(文件 + 元数据 + 内容哈希)
│
▼ 统一管道
解析 → 切分 → 元数据注入 → 幂等写入
│
▼
向量库 + ES + 登记表关键是把「数据源」和「处理管道」解耦。加一种新数据源 = 实现一个 SourceConnector,管道完全不动。
4.2 三种同步模式
| 模式 | 触发 | 适用 |
|---|---|---|
| 推送(Webhook) | 源系统变更时通知你 | 最优,实时且省资源 |
| 定时轮询 | 每 N 分钟全量列举并比对哈希 | 通用兜底 |
| 手动触发 | 用户点按钮 | 补充 |
优先做 Webhook,轮询做兜底。纯轮询在文档量大时,列举本身就是负担。
4.3 删除检测
轮询时:
源端文档集合 S = {doc1, doc2, doc3}
库端文档集合 D = {doc1, doc2, doc4}
D - S = {doc4} → 源端已删除 → 库端标记 deleted 并停止检索不要物理删除。标记 deleted=true 并在检索时过滤,这样:误删可恢复、有审计记录、不影响引用了该文档的历史问答记录。
4.4 大文件与超时
| 问题 | 处理 |
|---|---|
| 500 页 PDF | 异步处理,上传只返回 taskId(02-16) |
| 单文件切片过多 | 限制单文档最大片段数(如 500),超出告警 |
| 解析超时 | 单文档超时(如 120s),失败进重试队列 |
| 重复上传 | 内容哈希去重(02-14) |
5. 代码走查
5.1 数据源接口
java
// src/main/java/com/example/rag/ingest/source/SourceConnector.java
public interface SourceConnector {
String type(); // local / s3 / confluence
/** 列举当前所有文档及其哈希,用于变更检测与删除检测 */
List<DocumentRef> list(String tenantId);
/** 拉取单个文档内容 */
InputStream fetch(DocumentRef ref);
}
public record DocumentRef(String sourceId, String displayName,
String contentHash, long sizeBytes,
Instant modifiedAt, Map<String, Object> extra) {}5.2 同步服务(含删除检测)
java
// src/main/java/com/example/rag/ingest/source/SyncService.java
@Service
public class SyncService {
public SyncReport sync(String tenantId, SourceConnector connector) {
List<DocumentRef> remote = connector.list(tenantId);
Map<String, SourceRecord> local = registry.allByTenant(tenantId);
// 1. 新增 / 变更
List<DocumentRef> todo = remote.stream()
.filter(r -> {
SourceRecord l = local.get(r.sourceId());
return l == null || !l.contentHash().equals(r.contentHash());
})
.toList();
// 2. 删除检测(源端没有、库端有)
Set<String> remoteIds = remote.stream().map(DocumentRef::sourceId)
.collect(Collectors.toSet());
List<String> deleted = local.keySet().stream()
.filter(id -> !remoteIds.contains(id))
.toList();
deleted.forEach(id -> registry.markDeleted(id)); // 软删除
// 3. 提交异步灌库
String taskId = taskService.create(tenantId, todo);
return new SyncReport(todo.size(), deleted.size(), taskId);
}
}5.3 S3 连接器
java
// src/main/java/com/example/rag/ingest/source/S3Connector.java
@Component
@ConditionalOnProperty(name = "ai.source.s3.enabled", havingValue = "true")
public class S3Connector implements SourceConnector {
private final S3Client s3;
@Override
public List<DocumentRef> list(String tenantId) {
String prefix = tenantId + "/"; // 按租户前缀隔离
return s3.listObjectsV2Paginator(r -> r.bucket(bucket).prefix(prefix))
.contents().stream()
.map(o -> new DocumentRef(
o.key(), // sourceId = 对象 key
o.key().substring(prefix.length()),
o.eTag(), // S3 的 ETag 就是内容哈希
o.size(),
o.lastModified(),
Map.of()))
.toList();
}
@Override
public InputStream fetch(DocumentRef ref) {
return s3.getObject(r -> r.bucket(bucket).key(ref.sourceId()));
}
}用 ETag 作为内容哈希是省事的技巧:不需要下载文件再算哈希,列举时就拿到了。
5.4 Webhook 触发
java
// src/main/java/com/example/rag/controller/IngestWebhookController.java
@RestController
@RequestMapping("/api/webhook/ingest")
public class IngestWebhookController {
@PostMapping
public ResponseEntity<?> onWebhook(@RequestBody WebhookEvent e,
@RequestHeader("X-Signature") String sig) {
// 签名校验:不校验等于给外部开了一个灌库入口
if (!signatureVerifier.verify(e, sig)) {
return ResponseEntity.status(401).build();
}
syncService.syncOne(e.tenantId(), e.sourceId());
return ResponseEntity.accepted().build();
}
public record WebhookEvent(String tenantId, String sourceId, String action) {}
}5.5 失败队列与重投
java
// src/main/java/com/example/rag/ingest/DeadLetterService.java
@Service
public class DeadLetterService {
/** 重试 3 次仍失败 → 进死信队列,等待人工处理 */
public void handleFailure(IngestItem item, Exception e) {
int attempts = taskRepo.attempts(item);
if (attempts < 3) {
taskRepo.scheduleRetry(item, backoff(attempts));
} else {
deadLetter.save(item, e.getMessage());
alert.notify("灌库死信", item.sourceId(), e.getMessage());
}
}
}死信队列必须有告警。否则失败文档静默堆积,用户搜不到但没人知道。
6. 跑起来
bash
git checkout ch02-20-project-ingest
docker compose up -d
mvn spring-boot:runbash
# 1. 本地上传
curl -X POST http://localhost:8080/api/ingest -F "files=@handbook.pdf"
# 2. 触发同步(模拟 S3)
curl -X POST http://localhost:8080/api/sync -d '{"tenantId":"acme","source":"s3"}'
# 期望:{"toIngest":12,"deleted":2,"taskId":"..."}
# 3. 删除检测验证
# 在 S3 里删掉一个文件,再同步一次
curl -X POST http://localhost:8080/api/sync -d '{"tenantId":"acme","source":"s3"}'
# 期望:deleted >= 1,且该文档检索不到
# 4. 死信验证
# 上传一个损坏的 PDF,重试 3 次后进死信
curl http://localhost:8080/api/ingest/dead-letters| 检查项 | 通过标准 |
|---|---|
| 多数据源 | 至少两种数据源可跑通 |
| 增量 | 未变更文档被跳过 |
| 删除检测 | 源端删除后库端不再被检索到 |
| 死信 | 持续失败进死信并告警 |
| 权限 | Webhook 无签名被拒绝 |
7. 生产避坑
- 必须做删除检测,且用软删除。只做新增更新的同步,会让废弃文档永远留在库里被搜到——用户拿到过期制度的后果比搜不到更糟。软删除(标记而非物理删)则保留了恢复和审计的能力。
- Webhook 必须校验签名。暴露一个无鉴权的灌库接口,等于让任何人可以往你的知识库里灌内容——这是内容注入攻击的直接入口(见 04-06 安全合规)。
- 失败不能静默。灌库失败是常态(格式不支持、加密文档、扫描件),必须进死信队列 + 告警,并提供人工处理入口。静默失败的表现是「用户说搜不到,你查了半天说文档在啊」——排查成本极高。
8. 延伸与锚点
- 思考题:文档灌进去了,但业务方想看「哪些文档被问得最多、哪些问题答不上来」。这个后台怎么做?(答案在下一课时)
- 代码锚点:
git checkout ch02-20-project-ingest - 下一课时:02-21 完整项目(二):管理后台与运营闭环
- 对应课件:L02-20 灌库管道