Skip to content

完整项目(一):灌库管道与数据源接入 ​

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:run
bash
# 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. 生产避坑 ​

  1. 必须做删除检测,且用软删除。只做新增更新的同步,会让废弃文档永远留在库里被搜到——用户拿到过期制度的后果比搜不到更糟。软删除(标记而非物理删)则保留了恢复和审计的能力。
  2. Webhook 必须校验签名。暴露一个无鉴权的灌库接口,等于让任何人可以往你的知识库里灌内容——这是内容注入攻击的直接入口(见 04-06 安全合规)。
  3. 失败不能静默。灌库失败是常态(格式不支持、加密文档、扫描件),必须进死信队列 + 告警,并提供人工处理入口。静默失败的表现是「用户说搜不到,你查了半天说文档在啊」——排查成本极高。

8. 延伸与锚点 ​