文档智能处理管道:用 AI Functions 激活非结构化数据
概述
企业文档(合同、发票、财报、临床记录等)是业务的核心载体,但传统数仓无法处理 PDF、扫描件、图片中的复杂信息。云器 Lakehouse 提供了一整套文档智能管道:
文档(PDF/图片) → Volume 存储 → 文档解析 → AI 提取 → 结构化入库 → 分析洞察
所有处理在 Lakehouse 内部完成,数据无需外传。
技术栈
| 能力 | 组件 |
|---|
| 文档存储 | Volume |
| 文档解析 | pdf_to_md / pdf_to_json |
| 结构化提取 | AI_EXTRACT |
| 文本分类 | AI_CLASSIFY |
| LLM 推理 | AI_COMPLETE |
| 向量化 | AI_EMBEDDING |
| 语义搜索 | AI_SIMILARITY |
| 管道编排 | Dynamic Table |
文档解析函数
pdf_to_md — PDF → Markdown 文本
将 PDF 解析为 Markdown 格式文本,保留表格结构和布局信息。
语法:
pdf_to_md(<presigned_url>, '<option>=<value>')
参数:
| 选项 | 值 | 说明 |
|---|
mode
mode | layout
layout | 保留表格结构、多栏布局、图片位置(推荐) |
mode
mode | ocr
ocr | 扫描件 OCR 识别 |
示例:
-- 解析 PDF 为 Markdown 文本
SELECT pdf_to_md(GET_PRESIGNED_URL(USER VOLUME, 'invoice.pdf', 36000), 'mode=layout') AS doc_text;
实测输出(铁路电子客票):
电子发票(铁路电子客票)
发票号码:26129116010001275986
天津西站 G3559 杭州东站
票价:¥626.00
pdf_to_json — PDF → JSON 结构化数据
将 PDF 解析为结构化 JSON,包含元数据、分页文本和布局信息。
语法:
pdf_to_json(<presigned_url>)
返回结构:
| 字段 | 类型 | 说明 |
|---|
page_count
page_count | INT | 总页数 |
metadata
metadata | OBJECT | PDF 元数据(格式、标题、作者、创建时间等) |
pages
pages | ARRAY | 每页内容:number
number (页码)、text
text (提取文本)、layout
layout (布局信息) |
示例:
-- 解析 PDF 为 JSON
SELECT pdf_to_json(GET_PRESIGNED_URL(USER VOLUME, 'invoice.pdf', 36000)) AS doc_json;
-- 提取元数据
SELECT JSON_EXTRACT_STRING(
pdf_to_json(GET_PRESIGNED_URL(USER VOLUME, 'invoice.pdf', 36000)),
'$.metadata.format'
) AS pdf_format;
-- 提取第一页文本
SELECT JSON_EXTRACT_STRING(
pdf_to_json(GET_PRESIGNED_URL(USER VOLUME, 'invoice.pdf', 36000)),
'$.pages[0].text'
) AS page_text;
选择建议
| 场景 | 推荐函数 | 原因 |
|---|
| 全文检索、摘要 | pdf_to_md
pdf_to_md | 输出纯文本,直接输入 AI Functions |
| 元数据提取、结构化分析 | pdf_to_json
pdf_to_json | 输出 JSON,含页面布局和元数据 |
| OCR 扫描件 | pdf_to_md
pdf_to_md + mode=ocr
mode=ocr | 图片转文字 |
三大应用模式
模式一:文档知识库(企业搜索)
场景:财务报表、技术文档、合同等 PDF 入库,支持语义检索。
文档 → Volume → pdf_to_md → AI_EMBEDDING 向量化 → 语义搜索
-- Step 1: 解析 PDF 文档(保留布局)
CREATE OR REPLACE FUNCTION public.parse_pdf(file_path STRING)
RETURNS STRING
AS
pdf_to_md(GET_PRESIGNED_URL(USER VOLUME, file_path, 36000), 'mode=layout'); -- layout 模式保留表格/多栏/图片
-- Step 2: 向量化并入库
CREATE DYNAMIC TABLE doc_knowledge_base
REFRESH INTERVAL 30 MINUTE
AS
SELECT
SPLIT_PART(RELATIVE_PATH, '/', -1) AS doc_name,
RELATIVE_PATH,
public.parse_pdf(RELATIVE_PATH) AS doc_text,
AI_EMBEDDING(
public.parse_pdf(RELATIVE_PATH),
JSON '{"input": "document", "dimensions": "1024"}'
) AS doc_vector
FROM (
SELECT RELATIVE_PATH FROM SHOW USER VOLUME DIRECTORY SUBDIRECTORY 'documents'
WHERE RELATIVE_PATH ILIKE '%.pdf'
);
-- Step 3: 语义搜索
WITH query_vec AS (
SELECT AI_EMBEDDING('什么是湖仓一体?', JSON '{"input": "query", "dimensions": "1024"}') AS q
)
SELECT doc_name,
LEFT(doc_text, 200) AS excerpt,
AI_SIMILARITY(
AI_EMBEDDING(LEFT(doc_text, 1000), JSON '{"input": "document", "dimensions": "1024"}'),
q.q,
JSON '{"input": "text"}'
) AS relevance
FROM doc_knowledge_base, query_vec
ORDER BY relevance DESC
LIMIT 5;
模式二:业务流程自动化(发票/合同提取)
场景:从发票、合同、订单中提取结构化字段,替代人工录入。
文档 → Volume → pdf_to_json → AI_EXTRACT 字段提取 → 结构化表
-- Step 1: 解析 PDF 为 JSON(含表格和元数据)
CREATE OR REPLACE FUNCTION public.parse_pdf_json(file_path STRING)
RETURNS STRING
AS
pdf_to_json(GET_PRESIGNED_URL(USER VOLUME, file_path, 36000));
-- Step 2: 提取结构化字段
CREATE DYNAMIC TABLE dt_invoice_extracted
REFRESH INTERVAL 30 MINUTE
AS
SELECT
RELATIVE_PATH,
public.parse_pdf_json(RELATIVE_PATH) AS raw_json,
AI_EXTRACT(
public.parse_pdf_json(RELATIVE_PATH),
JSON '{"vendor_name":"供应商名称", "invoice_date":"发票日期",
"total_amount":"总金额", "tax_amount":"税额",
"line_items":"明细项目"}',
JSON '{"output.behavior": "raw_string"}'
) AS extracted_info
FROM (
SELECT RELATIVE_PATH FROM SHOW USER VOLUME DIRECTORY SUBDIRECTORY 'invoices'
WHERE RELATIVE_PATH ILIKE '%.pdf'
);
-- Step 3: 分类路由
SELECT RELATIVE_PATH,
AI_CLASSIFY(
extracted_info,
ARRAY('增值税发票', '普通发票', '合同', '订单', '其他'),
JSON '{"output.behavior": "raw_string"}'
) AS doc_type
FROM dt_invoice_extracted;
模式三:跨文档深度分析
场景:分析师对比多份财报、研报,识别趋势和异常。
文档解析 → 字段提取 → AI_COMPLETE 摘要 → AI_EMBEDDING 聚类 → 主题发现
完整管道见下方。
完整管道实现:财报分析
以下管道示例使用云器能力重写:
Step 1:解析财报 PDF
-- 创建解析函数
CREATE OR REPLACE FUNCTION public.parse_pdf_text(file_path STRING)
RETURNS STRING
AS
pdf_to_md(GET_PRESIGNED_URL(USER VOLUME, file_path, 36000), 'mode=layout');
-- 动态表:增量解析 PDF
CREATE DYNAMIC TABLE dt_reports_parsed
REFRESH INTERVAL 30 MINUTE
AS
SELECT
SPLIT_PART(SPLIT_PART(RELATIVE_PATH, '/', -1), '.', 1) AS company,
RELATIVE_PATH,
public.parse_pdf_text(RELATIVE_PATH) AS raw_text
FROM (
SELECT RELATIVE_PATH FROM SHOW USER VOLUME DIRECTORY SUBDIRECTORY 'reports'
WHERE RELATIVE_PATH LIKE '%.pdf'
);
Step 2:提取结构化财务字段
CREATE DYNAMIC TABLE dt_reports_extracted
REFRESH INTERVAL 30 MINUTE
AS
SELECT
company,
RELATIVE_PATH,
LEFT(raw_text, 120000) AS doc_text, -- 截断至 12 万字符
AI_EXTRACT(
LEFT(raw_text, 120000),
JSON '{"title":"公司名称和财年",
"fiscal_year":"4位数字财年",
"revenue":"总收入(含货币单位)",
"strategy_themes":"核心战略重点",
"growth_initiatives":"主要增长项目",
"headwinds":"主要风险和挑战"}',
JSON '{"output.behavior": "raw_string"}'
) AS extracted
FROM dt_reports_parsed;
Step 3:生成摘要并向量化
CREATE DYNAMIC TABLE dt_reports_embedded
REFRESH INTERVAL 30 MINUTE
AS
WITH summarised AS (
SELECT
company,
RELATIVE_PATH,
extracted,
AI_COMPLETE(
CONCAT('根据以下财报信息,生成包含4个字段的JSON摘要:
strategy(战略), financial_performance(财务表现),
outlook(展望), differentiation(差异化优势)。
公司:', COALESCE(JSON_EXTRACT_STRING(extracted, '$.title'), company), '
主题:', COALESCE(JSON_EXTRACT_STRING(extracted, '$.strategy_themes'), ''), '
风险:', COALESCE(JSON_EXTRACT_STRING(extracted, '$.headwinds'), '')),
JSON '{"output.behavior": "raw_string", "model.params": {"enable_thinking": false}}'
) AS summary_json
FROM dt_reports_extracted
)
SELECT
company,
RELATIVE_PATH,
JSON_EXTRACT_STRING(extracted, '$.title') AS title,
JSON_EXTRACT_STRING(extracted, '$.fiscal_year') AS fiscal_year,
JSON_EXTRACT_STRING(extracted, '$.revenue') AS revenue,
summary_json AS summary_raw,
AI_EMBEDDING(
CONCAT(COALESCE(JSON_EXTRACT_STRING(extracted, '$.strategy_themes'), '')),
JSON '{"input": "document", "dimensions": "256"}'
) AS strategy_vector
FROM summarised;
Step 4:跨文档主题发现
-- 用 AI_COMPLETE 分析全量报告的共性主题
SELECT AI_COMPLETE(
CONCAT('分析以下公司的战略主题,找出4-5个共性战略方向,每个方向给出名称和一句话描述。
只输出JSON:{"themes":[{"name":"...","description":"..."}]}
公司列表:', ARRAY_JOIN(COLLECT_LIST(
CONCAT(company, ': ', COALESCE(JSON_EXTRACT_STRING(extracted, '$.strategy_themes'), ''))
), '\n')),
JSON '{"output.behavior": "raw_string", "model.params": {"enable_thinking": false}}'
) AS strategic_themes
FROM dt_reports_extracted;
性能优化
1. 文档截断
AI 模型有上下文窗口限制。建议策略:
| 文档类型 | 截断长度 | 说明 |
|---|
| 发票/订单(短) | 不截断 | 通常 < 5000 字符 |
| 合同(中) | 50000 字符 | 保留关键条款部分 |
| 财报(长) | 120000 字符 | 保留管理层讨论+财务数据 |
| 研报/白皮书 | 80000 字符 | 保留摘要+结论 |
-- 安全截断
LEFT(raw_text, 120000)
2. 增量处理
使用 Dynamic Table 的
REFRESH INTERVAL 30 MINUTE
REFRESH INTERVAL 30 MINUTE
,新文档入库时只处理增量,避免重复解析:
CREATE DYNAMIC TABLE dt_docs_parsed
REFRESH INTERVAL 30 MINUTE
AS
SELECT ... FROM file_log
WHERE RELATIVE_PATH ILIKE '%.pdf';
3. 并发控制
大批量文档处理时设置
task.concurrency
task.concurrency
:
SELECT AI_EXTRACT(
doc_text,
JSON '{"field":"description"}',
JSON '{"task.concurrency": "8", "output.behavior": "raw_string"}'
)
FROM large_doc_set;
产品能力
| 能力 | 说明 | 技术点 |
|---|
| 文档解析(PDF→文本) | PDF 转结构化文本 | pdf_to_md / pdf_to_json |
| 结构化字段提取 | 自然语言定义提取字段 | AI_EXTRACT |
| 文档分类 | 自动路由文档类型 | AI_CLASSIFY |
| LLM 推理 | 通用文本生成 | AI_COMPLETE |
| 向量化 | 文本转向量 | AI_EMBEDDING |
| 语义搜索 | 向量相似度检索 | AI_SIMILARITY + VECTOR INDEX |
| 管道编排 | 声明式增量刷新 | Dynamic Table |
| 增量刷新 | 仅处理新文件 | REFRESH INTERVAL 30 MINUTE |
| JSON Schema 提取 | 结构化字段定义 | AI_EXTRACT + JSON schema |
| 文档布局保留 | 保留表格/多栏结构 | pdf_to_md 的 layout 模式 |
| 存储 | 文档文件存储 | Volume |
最佳实践
1. 文档解析策略
| 场景 | 推荐解析方式 | 原因 |
|---|
| 全文检索 | pdf_to_md + layout | 保留表格和段落结构 |
| 字段提取 | pdf_to_json + AI_EXTRACT | JSON 更易定位目标字段 |
| OCR 扫描件 | pdf_to_md + OCR 模式 | 图片转文字 |
| 混合场景 | 两者都跑 | md 做搜索,json 做提取 |
2. 提取 schema 设计
-- 好的 schema:描述清晰,让 AI 理解意图
JSON '{"vendor_name":"供应商全称(如"北京云器科技有限公司")",
"invoice_date":"发票开具日期(格式:YYYY-MM-DD)",
"total_amount":"含税总金额(含货币单位)"}'
-- 不好的 schema:过于模糊
JSON '{"name":"名称", "date":"日期", "amount":"金额"}'
3. 错误处理
-- 兜底:提取失败时返回默认值
SELECT COALESCE(
NULLIF(AI_EXTRACT(doc_text, JSON '{"amount":"总金额"}',
JSON '{"output.behavior": "raw_string"}'), ''),
'提取失败'
) AS amount;
-- 人工审核:记录低置信度结果
INSERT INTO review_queue
SELECT * FROM extracted_docs
WHERE confidence < 0.6;
4. 成本优化
| 优化策略 | 方法 | 效果 |
|---|
| 文档截断 | LEFT(text, 120000)
LEFT(text, 120000) | 减少 token 消耗 |
| 增量处理 | Dynamic Table INCREMENTAL | 避免重复解析 |
| 分类先行 | 先用 AI_CLASSIFY 分类型,再走不同管道 | 精准匹配提取模板 |
| 并发控制 | task.concurrency
task.concurrency 调参 | 平衡速度和稳定性 |
5. 完整生产管道示例
-- ════════════════════════════════════════════════
-- 文档智能处理生产管道
-- ════════════════════════════════════════════════
-- 1. 解析 PDF
CREATE DYNAMIC TABLE dt_docs_parsed REFRESH INTERVAL 30 MINUTE AS
SELECT ... pdf_to_md(GET_PRESIGNED_URL(USER VOLUME, relative_path, 36000), 'mode=layout') AS raw_text FROM file_log WHERE ...;
-- 2. 文档分类
CREATE DYNAMIC TABLE dt_docs_classified REFRESH INTERVAL 30 MINUTE AS
SELECT *, AI_CLASSIFY(raw_text, ARRAY('invoice','contract','report','other'),
JSON '{"output.behavior": "raw_string"}') AS doc_type
FROM dt_docs_parsed;
-- 3. 按类型提取
CREATE DYNAMIC TABLE dt_docs_extracted REFRESH INTERVAL 30 MINUTE AS
SELECT *,
CASE doc_type
WHEN 'invoice' THEN AI_EXTRACT(raw_text, invoice_schema)
WHEN 'contract' THEN AI_EXTRACT(raw_text, contract_schema)
WHEN 'report' THEN AI_EXTRACT(raw_text, report_schema)
END AS extracted_data
FROM dt_docs_classified;
-- 4. 向量化入知识库
CREATE DYNAMIC TABLE dt_docs_embedded REFRESH INTERVAL 30 MINUTE AS
SELECT *, AI_EMBEDDING(LEFT(raw_text,5000),
JSON '{"input":"document","dimensions":"256"}') AS vector
FROM dt_docs_extracted;
-- 5. 业务查询(语义搜索)
SELECT doc_name, extracted_data,
AI_SIMILARITY(AI_EMBEDDING(...), query_vector) AS relevance
FROM dt_docs_embedded
ORDER BY relevance DESC LIMIT 10;
注意事项
| 注意点 | 说明 |
|---|
| 文档大小 | 单份文档建议 ≤ 50MB,超长文档先拆分再处理 |
| token 限制 | AI_EXTRACT
AI_EXTRACT 和 AI_COMPLETE
AI_COMPLETE 受模型 context window 限制,长文档需截断 |
| OCR 质量 | 扫描件清晰度直接影响提取准确率,建议 300dpi 以上 |
| 增量处理 | 确保文件有 change_tracking
change_tracking 支持,否则全量刷新 |
| 成本监控 | 使用 AI_CONTEXT_LENGTH
AI_CONTEXT_LENGTH 在解析前预估 token 消耗 |
| 权限控制 | Volume 中的文档通过标准 RBAC 控制访问权限 |
| 结果审核 | 关键字段(如金额、日期)建议人工抽检 |