Skip to content

RAG 文档导入与清洗:从异构数据到可追溯 Document ​

在 RAG 系统中,数据加载是第一步。我们需要将各种格式的非结构化数据(PDF, Word, Markdown, HTML)加载到内存中,并统一转换为标准格式。

1. 核心数据结构 ​

在 Golang 中,我们通常定义一个通用的 Document 结构体来承载数据:

go
type Document struct {
    PageContent string                 // 文档正文
    Metadata    map[string]interface{} // 元数据(如文件名、页码、来源URL)
}

这与 langchaingo 中的 schema.Document 是一致的。

2. 常见文件格式加载 ​

2.1 加载 TXT / Markdown 文件 ​

对于纯文本文件,使用 Golang 标准库 os 和 io 即可轻松搞定。

go
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 示例:

go
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 标签。

go
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。

go
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. 增量更新、删除与失败恢复 ​

数据管道至少要区分新增、内容更新、权限变化和删除:

  1. 先记录一次导入任务及其源版本,再异步解析、切片和向量化。
  2. 新索引版本完成并通过抽检后再切换读流量,避免查询读到半成品。
  3. 内容更新时使旧 Chunk 失效;权限变化应尽快传播到所有检索索引和缓存。
  4. 删除不仅要处理数据库记录,还要覆盖向量、关键词索引、缓存和派生评估样本。
  5. 对失败项保留可重试状态和错误分类,不把“任务结束”误报为“全部成功”。

稳定的文档 ID、Chunk ID 和索引版本是幂等重试与回滚的基础。切片数据契约和参数选择见 RAG 文本切分与 Chunk Size。

总结 ​

数据加载的核心是将异构数据转化为同构的 Document 对象。在 Golang 中,利用接口(Interface)和并发(Goroutine)可以构建出非常高效的数据处理管道(ETL Pipeline)。

但生产目标不只是吞吐量,而是可追溯、可重放、可删除和可验证。回到 RAG 专题可以继续完整学习路径;准备上线时,使用 RAG 生产化检查清单核对数据、权限、观测和回滚,并进入 RAG 检索策略验证导入结果是否真的可被找到。

AI 应用开发训练营 · 现在开始

别只收藏 AI 文章,今天就做出第一个能上线的项目

文章解决认知,项目才证明能力。把 Agent、RAG、MCP 变成可运行、可展示、可面试表达的成果,现在就从一次岗位准备度自测开始。

立即开始 AI 岗位准备度自测 查看训练营路线
真实学员成果先看成果,再决定是否开始 →
  1. 01选方向
  2. 02做项目
  3. 03出成果
  4. 04拿面试