这篇不先堆名词。我们把《我用大数据经验做了次 AI 项目,最先失效的是旧方法》拆成几级台阶,看完至少知道下一步该学什么、该练什么。

摘要

去年接了个内部需求,业务方要做一个"制度文档智能问答",底层走 RAG。我带队,三个大数据方向的工程师,两个做 ETL 的,一个做数据治理的。Demo 两周就跑通了,业务方很满意。结果真正上线第一周,权限问题、日志丢失、响应时间不可观测,三个问题全来了。

这篇文章复盘的就是那次项目,以及我作为数据工程师,在大模型工程化过程中踩过的坑和总结的方法论。

---

目录

  • 大数据与大模型的交叉点
  • 数据治理:RAG 系统的隐形门槛
  • 向量数据库:选型比技术更重要
  • RAG 数据管道:从 Demo 到生产的关键差异
  • 落地项目:权限、日志、可观测性的真实案例
  • 总结:数据工程师的差异化优势

大数据与大模型的交叉点

文章插图 1

很多人认为大数据转大模型,就是学一下 LangChain、调一下 API。这个思路本身没错,但有一个前提:你的核心价值是什么?

大数据工程师的核心能力,不是会写 SQL,也不是会调 Spark,而是——对数据血缘、质量、治理、可追溯性的理解。这些能力在大模型时代反而更值钱,因为 LLM 应用最大的问题,从来不是模型本身,而是数据从哪里来、去了哪里、出了什么问题。

我做 RAG 项目时,发现三个交叉点最值钱:

非结构化数据处理。传统数据工程师习惯处理结构化数据,但 LLM 应用的输入输出大量是非结构化的——文档、图片、日志。怎么清洗、分段、向量化,这是大数据工程师的延伸能力。

数据质量保障。RAG 系统的效果,70% 取决于知识库的质量。脏数据进去,垃圾答案出来。数据工程师的治理方法论,在这里直接复用。

可观测性。传统数据管道有 Airflow、有日志、有告警。RAG 系统呢?很多时候只有"答案对不对"。这就是为什么权限、日志、可观测性成为上线卡点。

---

数据治理:RAG 系统的隐形门槛

文章插图 2

业务方提需求时,只说了一句话:"员工问制度,系统能给准确答案。"

看起来简单,真正做起来,数据治理是第一个拦路虎。

案例:我们接入了公司 300+ 份制度文档,格式五花八门——PDF、Word、Excel、图片扫描件。最麻烦的是扫描件,OCR 之后文字识别错误率高达 15%,直接喂给 Embedding 模型,检索质量极差。

我们的处理方案:

原始文档 → 格式分类 → 图片走 OCR → 文本走解析 → 分段 → 去重 → 质量校验 → 入库

质量校验这一步,传统数据工程师会做数据采样、规则校验、异常检测。我们用同样的思路,对向量库做抽检:

1. 随机抽取 5% 的文档片段
2. 用高质量 Prompt 生成答案
3. 人工标注正确率
4. 低于 80% 的片段重新清洗

这个过程花了两周,但上线后效果显著——错误回答率从 35% 降到 8%。

失败原因分析:很多人跳过了数据治理,直接做 Embedding + 检索。这不是技术问题,是认知问题——认为 RAG 就是"把文档扔进向量库"。实际上,数据质量决定了上限,模型只是逼近这个上限的工具。

---

CSDN资料领取方式

向量数据库:选型比技术更重要

向量数据库的选型,大数据工程师有天然优势——你们见过太多存储方案了,MySQL、HBase、ClickHouse、Elasticsearch,每种方案的适用边界都清楚。

我们评估了四个选项:

| 方案 | 优势 | 劣势 | 适用场景 |
|------|------|------|----------|
| Milvus | 功能完整、社区活跃 | 部署重、运维成本高 | 大规模生产环境 |
| Chroma | 轻量、易用 | 不适合分布式 | 本地 Demo、小规模 |
| Weaviate | 内置混合检索 | 中文支持一般 | 多语言场景 |
| Elasticsearch + KNN | 团队熟悉、运维简单 | 向量检索性能一般 | 已有 ES 集群的团队 |

我们的选择:ES + KNN。

理由很数据工程师——团队熟悉度优先于技术先进性。公司已有 ES 集群,运维团队熟悉,监控告警体系完善。虽然向量检索性能不如 Milvus,但对于几千份制度文档的规模,完全够用。

适用边界:如果你的场景是百万级向量、毫秒级延迟要求,选 Milvus 或 Weaviate。如果规模在万级以下、团队已有 ES 基础,ES + KNN 是最优解。

---

RAG 数据管道:从 Demo 到生产的关键差异

Demo 版本和生產版本的差距,不在代码复杂度,而在工程化细节。

我们做了三件事,把 Demo 升级到了生产可用:

第一,权限过滤前置。不是检索完再过滤,而是在检索前就过滤。用户查询时,先根据部门、职级获取可访问的文档 ID 列表,再在这个子集里做向量检索。

第二,结构化日志。每个请求记录:

  • 输入查询
  • 检索到的文档片段(含来源)
  • 生成结果
  • 各阶段耗时
  • 权限过滤结果

第三,分阶段延迟监控。把 RAG 流程拆成四段:权限查询 → 向量检索 → Prompt 构建 → 模型生成。每段独立计时,快速定位瓶颈。

下面是核心代码:

import logging
import time
from typing import List, Dict, Any

# 结构化日志配置
logging.basicConfig(
    format='%(asctime)s | %(levelname)s | %(message)s',
    handlers=[logging.FileHandler('rag_pipeline.log')]
)
logger = logging.getLogger('rag_pipeline')

class RAGPipeline:
    def __init__(self, es_client, llm_client, permission_service):
        self.es = es_client
        self.llm = llm_client
        self.perm = permission_service

    def query(self, user_id: str, question: str) -> Dict[str, Any]:
        start_time = time.time()
        trace = {
            'user_id': user_id,
            'question': question,
            'stages': {},
            'result': None
        }

        try:
            # 阶段1:权限过滤
            t0 = time.time()
            allowed_doc_ids = self.perm.get_accessible_docs(user_id)
            trace['stages']['permission'] = {
                'duration_ms': (time.time() - t0) * 1000,
                'doc_count': len(allowed_doc_ids)
            }
            logger.info(f"权限过滤完成: user={user_id}, docs={len(allowed_doc_ids)}")

            if not allowed_doc_ids:
                return {'answer': '您无权访问任何文档', 'trace': trace}

            # 阶段2:向量检索
            t1 = time.time()
            chunks = self.es.search(
                query=question,
                doc_ids=allowed_doc_ids,
                top_k=5
            )
            trace['stages']['retrieval'] = {
                'duration_ms': (time.time() - t1) * 1000,
                'chunk_count': len(chunks)
            }
            logger.info(f"检索完成: chunks={len(chunks)}")

            # 阶段3:Prompt 构建
            t2 = time.time()
            context = self._build_context(chunks)
            prompt = self._build_prompt(question, context)
            trace['stages']['prompt_build'] = {
                'duration_ms': (time.time() - t2) * 1000
            }

            # 阶段4:模型生成
            t3 = time.time()
            answer = self.llm.generate(prompt)
            trace['stages']['generation'] = {
                'duration_ms': (time.time() - t3) * 1000,
                'tokens': len(answer.split())
            }

            trace['result'] = answer
            trace['total_duration_ms'] = (time.time() - start_time) * 1000

        except Exception as e:
            trace['error'] = str(e)
            logger.error(f"查询失败: user={user_id}, error={e}")
            raise

        return trace

代码解释:

  • permission.get_accessible_docs:权限服务,根据用户 ID 返回可访问文档列表。这是生产版本和 Demo 的关键区别——Demo 通常跳过这一步。
  • es.search:ES 向量检索,传入 doc_ids 参数实现权限过滤后的检索。
  • 每个阶段独立计时,记录到 trace 中。上线后可以通过这些日志快速定位瓶颈:如果 retrieval 阶段耗时过长,说明 ES 索引需要优化;如果 generation 阶段超时,可能是模型响应慢或 Prompt 过长。
  • 异常处理:任何阶段失败都会记录错误日志并抛出,不会静默失败。

---

落地项目:权限、日志、可观测性的真实案例

上线第一周,我们遇到了三个问题:

问题一:权限穿透。用户 A 属于技术部,但通过某种方式查到了财务部的制度文档。排查后发现,权限服务返回的文档列表是正确的,但向量检索时没有把 doc_ids 传进去,导致检索范围不受限制。

教训:权限过滤必须在检索层生效,不能在结果层过滤。结果层过滤有性能问题,而且容易遗漏。

问题二:日志丢失。某个用户反馈回答错误,但我们查日志时发现,那个请求的日志缺失了。排查后发现,日志是异步写入的,当服务重启时,缓冲区的日志丢失。

教训:关键日志必须同步写入,或者使用持久化队列。对于 RAG 系统,每个请求的 trace 就是"数据血缘",丢了就是不可追溯。

问题三:延迟不可观测。用户抱怨响应慢,但我们不知道慢在哪里。排查后发现,generation 阶段平均耗时 3 秒,但没有任何监控告警。

教训:每个阶段都应该有 P99 延迟监控和告警。不是等用户投诉才发现。

---

总结:数据工程师的差异化优势

大数据转大模型,不是从零开始,而是能力迁移。

数据工程师的核心优势有三:

1. 数据治理方法论。RAG 系统的效果上限由数据质量决定,这是数据工程师的强项。
2. 可观测性意识。传统数据管道有完整的监控告警体系,这个思维可以直接迁移到 LLM 应用。
3. 工程化能力。权限、日志、容错、性能优化,这些是生产环境的必备能力,也是 Demo 和生产版本的最大差距。

Demo 能跑通,不代表能上线。真正卡住团队的,从来不是模型本身,而是权限、日志和可观测性。这三个问题,大数据工程师应该比任何人都熟悉。

资料展示

下面是我整理的AI大模型学习资料和工具包预览,适合收藏后按主题逐步学习。

AI大模型资料展示 1

AI大模型资料展示 2

AI大模型资料展示 3

AI大模型资料展示 4

需要这份AI大模型资料清单的话,在评论区回复「清单」即可;我会根据大家的问题继续补充对应的实战内容。

CSDN官方大礼包

Logo

有“AI”的1024 = 2048,欢迎大家加入2048 AI社区

更多推荐