RAG 文档导入与清洗:从异构数据到可追溯 Document
在 RAG 系统中,数据加载是第一步。我们需要将各种格式的非结构化数据(PDF, Word, Markdown, HTML)加载到内存中,并统一转换为标准格式。
1. 核心数据结构
在 Golang 中,我们通常定义一个通用的 Document 结构体来承载数据:
type Document struct {
PageContent string // 文档正文
Metadata map[string]interface{} // 元数据(如文件名、页码、来源URL)
}这与 langchaingo 中的 schema.Document 是一致的。
2. 常见文件格式加载
2.1 加载 TXT / Markdown 文件
对于纯文本文件,使用 Golang 标准库 os 和 io 即可轻松搞定。
package main
import (
"context"
"fmt"
"os"
"github.com/tmc/langchaingo/documentloaders"
)
func LoadText() {
f, _ := os.Open("example.md")
defer f.Close()
// 使用 langchaingo 的 TextLoader
loader := documentloaders.NewText(f)
docs, err := loader.Load(context.Background())
if err != nil {
panic(err)
}
fmt.Println(docs[0].PageContent)
}2.2 加载 PDF 文件
PDF 解析稍微复杂一些,可以基于经过项目验证的 PDF 解析库封装 Loader。库的格式覆盖、维护状态和 API 会变化,选型时应使用实际语料验证文本顺序、表格、扫描件和错误处理,不要只看“能读取文件”。
封装 PDFLoader 示例:
package main
import (
"bytes"
"context"
"fmt"
"github.com/ledongthuc/pdf"
"github.com/tmc/langchaingo/schema"
)
// PDFLoader 实现 documentloaders.Loader 接口
type PDFLoader struct {
path string
}
func NewPDFLoader(path string) *PDFLoader {
return &PDFLoader{path: path}
}
func (l *PDFLoader) Load(ctx context.Context) ([]schema.Document, error) {
f, r, err := pdf.Open(l.path)
if err != nil {
return nil, err
}
defer f.Close()
var buf bytes.Buffer
b, err := r.GetPlainText() // 获取纯文本
if err != nil {
return nil, err
}
buf.ReadFrom(b)
return []schema.Document{
{
PageContent: buf.String(),
Metadata: map[string]any{
"source": l.path,
},
},
}, nil
}
func main() {
loader := NewPDFLoader("manual.pdf")
docs, _ := loader.Load(context.Background())
fmt.Printf("PDF 内容长度: %d\n", len(docs[0].PageContent))
}2.3 加载网页 (HTML)
可以使用 goquery 库来提取网页正文,去除 HTML 标签。
import (
"bytes"
"context"
"errors"
"fmt"
"io"
"net"
"net/http"
"net/url"
"strings"
"time"
"github.com/PuerkitoBio/goquery"
)
const maxHTMLBytes = 2 << 20
func isBlockedIP(ip net.IP) bool {
return ip.IsLoopback() || ip.IsPrivate() || ip.IsLinkLocalUnicast() ||
ip.IsLinkLocalMulticast() || ip.IsUnspecified() || ip.IsMulticast()
}
func newPublicHTTPClient() *http.Client {
dialer := &net.Dialer{Timeout: 5 * time.Second}
transport := &http.Transport{}
transport.DialContext = func(ctx context.Context, network, address string) (net.Conn, error) {
host, port, err := net.SplitHostPort(address)
if err != nil {
return nil, err
}
ips, err := net.DefaultResolver.LookupIP(ctx, "ip", host)
if err != nil {
return nil, fmt.Errorf("resolve host: %w", err)
}
if len(ips) == 0 {
return nil, errors.New("host resolved without an address")
}
for _, ip := range ips {
if isBlockedIP(ip) {
return nil, errors.New("private or special-use address is not allowed")
}
}
return dialer.DialContext(ctx, network, net.JoinHostPort(ips[0].String(), port))
}
return &http.Client{
Transport: transport,
Timeout: 10 * time.Second,
CheckRedirect: func(req *http.Request, via []*http.Request) error {
if len(via) >= 5 {
return errors.New("too many redirects")
}
if req.URL.Scheme != "http" && req.URL.Scheme != "https" {
return errors.New("redirect scheme is not allowed")
}
return nil // 新地址仍会经过受限 DialContext 解析。
},
}
}
func FetchURL(ctx context.Context, rawURL string) (string, error) {
parsed, err := url.Parse(rawURL)
if err != nil || (parsed.Scheme != "http" && parsed.Scheme != "https") || parsed.Hostname() == "" {
return "", errors.New("only absolute http/https URLs are allowed")
}
req, err := http.NewRequestWithContext(ctx, http.MethodGet, parsed.String(), nil)
if err != nil {
return "", err
}
res, err := newPublicHTTPClient().Do(req)
if err != nil {
return "", err
}
defer res.Body.Close()
if res.StatusCode < 200 || res.StatusCode >= 300 {
return "", fmt.Errorf("unexpected HTTP status: %s", res.Status)
}
body, err := io.ReadAll(io.LimitReader(res.Body, maxHTMLBytes+1))
if err != nil {
return "", err
}
if len(body) > maxHTMLBytes {
return "", errors.New("response body exceeds size limit")
}
doc, err := goquery.NewDocumentFromReader(bytes.NewReader(body))
if err != nil {
return "", err
}
// 移除 script 和 style 标签
doc.Find("script,style").Each(func(i int, s *goquery.Selection) {
s.Remove()
})
return strings.TrimSpace(doc.Text()), nil
}这个示例在实际连接时校验解析后的 IP,并让重定向后的地址再次经过同一限制;仅在字符串层检查域名无法防止 DNS 解析到内网地址。生产环境还应配置出站网络策略,将云元数据和内部网段在基础设施层一并阻断。
3. 并发加载优化
Golang 的最大优势在于并发。当需要加载数千个文档时,一定要利用 Goroutine。
func LoadDir(dir string) []schema.Document {
files, _ := os.ReadDir(dir)
ch := make(chan schema.Document, len(files))
// 启动并发加载
for _, file := range files {
go func(f os.DirEntry) {
// ... 加载逻辑 ...
ch <- doc
}(file)
}
// 收集结果
var allDocs []schema.Document
for range files {
allDocs = append(allDocs, <-ch)
}
return allDocs
}上面的代码只演示并发结构。生产实现还需要限制并发数、传播 context 取消、回收文件和网络资源、收集每个文件的错误,并为重试设置幂等键;不能因为某个 Loader 卡住而无限占用 Goroutine。
4. 导入后的清洗与质量门禁
加载成功不等于内容可用于检索。清洗阶段应保留原始文档,只对派生文本执行可重放的转换:
- 去除可识别的重复页眉、页脚和导航噪声,但不要删除可能影响语义的标题、表头和条件说明。
- 统一编码和换行,识别空页、乱码、异常短文本与重复正文。
- PDF、网页和表格按来源类型抽样,检查阅读顺序、字段名、页码与链接是否保留。
- 扫描文档需要 OCR 时,记录 OCR 工具与版本,并把低置信或布局复杂页面送入人工复核。
- 将清洗规则版本写入元数据,使错误规则可以定位并重建。
建议为每批导入生成质量报告,包括成功、失败、跳过、重复和待复核的文档列表。具体通过条件应由语料风险和业务容忍度决定,不应从其他项目复制一个固定比例。
5. 可追溯的 Document 契约
除正文外,生产中的 Document 通常还需要:
| 字段 | 作用 |
|---|---|
document_id | 关联业务主键,支持更新和删除 |
source_uri | 返回原始来源并排查解析问题 |
source_version | 区分文档修订,避免新旧内容混用 |
title_path / page | 支持切片语义和精确引用 |
tenant_id / access_scope | 为后续服务端权限过滤提供依据 |
content_hash | 识别未变化内容,减少重复处理 |
parser_version / ingested_at | 追踪处理链路并支持重建 |
敏感级别、保留期限和删除状态也应随文档进入索引链路。权限元数据只是检索过滤的输入,不是授权事实本身;查询时仍要根据可信身份执行服务端鉴权。
6. 增量更新、删除与失败恢复
数据管道至少要区分新增、内容更新、权限变化和删除:
- 先记录一次导入任务及其源版本,再异步解析、切片和向量化。
- 新索引版本完成并通过抽检后再切换读流量,避免查询读到半成品。
- 内容更新时使旧 Chunk 失效;权限变化应尽快传播到所有检索索引和缓存。
- 删除不仅要处理数据库记录,还要覆盖向量、关键词索引、缓存和派生评估样本。
- 对失败项保留可重试状态和错误分类,不把“任务结束”误报为“全部成功”。
稳定的文档 ID、Chunk ID 和索引版本是幂等重试与回滚的基础。切片数据契约和参数选择见 RAG 文本切分与 Chunk Size。
总结
数据加载的核心是将异构数据转化为同构的 Document 对象。在 Golang 中,利用接口(Interface)和并发(Goroutine)可以构建出非常高效的数据处理管道(ETL Pipeline)。
但生产目标不只是吞吐量,而是可追溯、可重放、可删除和可验证。回到 RAG 专题可以继续完整学习路径;准备上线时,使用 RAG 生产化检查清单核对数据、权限、观测和回滚,并进入 RAG 检索策略验证导入结果是否真的可被找到。