大数据转大模型:Demo能跑只是开始,权限日志才是真正的分水岭
2024年下半年,我所在的团队决定把内部知识库从传统检索升级成大模型问答系统。我当时想的是,大数据工程师做这个有什么难的?数据清洗、特征工程、模型训练,哪样没干过?
结果上线第一周,系统就崩了三次。
不是因为模型不好,不是因为Prompt写得差,是因为权限问题——某个用户通过API拿到了不该看的数据,日志追踪不到来源,可观测性为零。
这篇文章复盘的就是这个踩坑过程,以及大数据工程师转型大模型时真正需要补的能力。
---
摘要
大数据工程师转型大模型,最常见的误区是觉得缺的是算法能力。实际上,Demo阶段和工程化阶段之间隔着一道门槛:权限控制、日志追踪、可观测性。本文结合一个真实项目踩坑经历,梳理数据治理、向量数据库选型、RAG数据管道构建、以及生产环境上线的关键经验。
---
目录
- 大数据与大模型的交叉点
- 数据治理:从ETL到数据质量闭环
- 向量数据库:选型比技术栈更重要
- RAG数据管道:Demo和生产之间隔着一道坎
- 落地项目:权限日志才是真正门槛
- 总结
---
大数据与大模型的交叉点
大数据工程师和大模型工程师的交集,其实比想象中多。
核心能力重叠在三个地方:
第一,数据处理。 大模型需要高质量训练数据,需要清洗、去重、格式化。这和做数据仓库的ETL逻辑一脉相承。
第二, pipeline构建。 无论是Spark批处理还是Flink流处理,思路都是"输入→处理→输出"。RAG的数据管道本质上也是这个逻辑,只是输入变成了非结构化文档,输出变成了向量。
第三,可观测性思维。 大数据系统天生就需要监控——任务失败率、数据延迟、资源占用。大模型应用同样需要,只是监控的对象从"数据流"变成了"请求流"。
我的认知转变发生在一个具体场景:
团队让我负责知识库RAG系统的后端数据管道。我以为就是把文档切块、向量化、存入向量库。做起来才发现,真正卡住的是:文档从哪个系统来、谁有权访问、更新频率是多少、失败怎么重试、向量库挂了怎么降级。
这些都不是算法问题,是工程问题。
---
数据治理:从ETL到数据质量闭环
大数据工程师做数据治理是基本功,但大模型场景下的数据治理有几个新挑战。
传统ETL关注的是结构化数据的质量。 字段缺失、类型错误、重复记录,这些都有成熟的校验逻辑。
大模型场景需要关注的是非结构化数据的质量。 一段文档内容是否准确?是否存在敏感信息?切块后语义是否完整?
我们踩的第一个坑是文档来源管理。
知识库的文档来自三个系统:Wiki、Confluence、和手动上传的PDF。Wiki和Confluence有版本控制,但PDF没有。结果某次模型回答了一个过期的政策,追溯发现是有人手动上传了旧版PDF,覆盖了新版本。
解决方案是建立数据血缘:
# 文档数据血缘记录 class DocumentLineage: def __init__(self, doc_id, source_system, source_url, upload_time, version, creator, hash_value): self.doc_id = doc_id self.source_system = source_system # wiki/confluence/manual self.source_url = source_url self.upload_time = upload_time self.version = version self.creator = creator self.hash_value = hash_value # 用于检测重复和变更 def is_same_content(self, new_hash): """检测文档是否被修改""" return self.hash_value == new_hash def get_access_level(self): """根据来源系统返回默认权限级别""" access_map = { 'wiki': 'internal', 'confluence': 'team', 'manual': 'restricted' } return access_map.get(self.source_system, 'internal')这段代码的核心思想是:每条文档进入系统时,记录它的来源、创建者、版本、内容哈希。后续任何查询都可以追溯到原始数据。
权限级别根据来源系统自动分配。 Wiki文档默认内部可见,Confluence文档根据团队权限继承,手动上传的PDF默认受限。
这个设计看起来简单,但很多Demo项目根本不做这一步。
---
向量数据库:选型比技术栈更重要
向量数据库选型,大数据工程师有天然优势。
为什么?因为选型逻辑和数据仓库选型逻辑一样:看场景,而不是看参数。
我们团队调研过四种主流方案:
| 方案 | 优势 | 劣势 | 适用场景 |
|------|------|------|----------|
| Milvus | 功能最全,支持复杂查询 | 部署重,运维成本高 | 大规模生产环境 |
| Chroma | 轻量,Python友好 | 不适合分布式部署 | 个人项目、小规模应用 |
| Weaviate | 内置向量+图查询 | 生态相对小众 | 需要混合查询的场景 |
| pgvector | PostgreSQL扩展,熟悉度高 | 性能上限有限 | 已有PostgreSQL基础设施的团队 |
我们的选择是pgvector。
原因很实际:团队已经有PostgreSQL集群,运维熟悉,不想再引入一个新的存储系统。而且知识库规模不大,百万级向量,pgvector完全够用。
如果你们团队已经有大数据基础设施(HDFS、Spark、Flink),选型逻辑也一样:优先考虑和现有栈的集成成本,而不是单项性能指标。
代码层面,pgvector的集成非常简单:
from pgvector.sqlalchemy import Vector from sqlalchemy import create_engine, Column, String from sqlalchemy.ext.declarative import declarative_base Base = declarative_base() class DocumentChunk(Base): __tablename__ = 'document_chunks' id = Column(String, primary_key=True) content = Column(Text, nullable=False) embedding = Column(Vector(1536)) # OpenAI ada-002 维度 doc_id = Column(String, nullable=False) source_system = Column(String, nullable=False) access_level = Column(String, nullable=False) def __repr__(self): return f"<DocumentChunk(id={self.id}, doc_id={self.doc_id})>" engine = create_engine("postgresql+psycopg2://user:pass@host/db") Base.metadata.create_all(engine)这段代码的核心是:向量存储和元数据放在同一个表里。查询时先过滤权限,再做向量检索,避免拿到不该看的数据。
---
RAG数据管道:Demo和生产之间隔着一道坎
RAG的数据管道,理论上很简单:文档→切块→向量化→存储→检索→生成。
但实际上,每个环节都有坑。
切块不是简单的字符串分割。
我们最初用固定长度切块,每500字一块。结果一个问题被切到了两块里,模型只能看到一半,回答错误。
后来改成语义切块,用句子边界+长度约束:
import re from typing import List, Tuple def semantic_chunk(text: str, max_length: int = 500, min_length: int = 100) -> List[str]: """ 语义切块:优先在句子边界切分,兼顾长度约束 """ # 按句子分割(中文句号、问号、感叹号) sentences = re.split(r'([。!?])', text) chunks = [] current_chunk = "" for i in range(0, len(sentences), 2): sentence = sentences[i] # 保留标点 punct = sentences[i+1] if i+1 < len(sentences) else "" # 如果当前块+新句子超过最大长度,先输出当前块 if len(current_chunk) + len(sentence + punct) > max_length: if len(current_chunk) >= min_length: chunks.append(current_chunk) current_chunk = sentence + punct else: current_chunk += sentence + punct # 处理最后一个块 if current_chunk and len(current_chunk) >= min_length: chunks.append(current_chunk) return chunks这个切块策略的核心是:优先保证语义完整性,其次满足长度约束。
向量化不是调个API就完事。
我们最初用OpenAI的embedding API,效果不错,但有两个问题:成本高、延迟高。
后来换成本地部署的text2vec模型,成本降了90%,延迟从500ms降到50ms。代价是精度略有下降,但对内部知识库场景完全够用。
检索不是相似度越高越好。
我们有一个案例:用户问"2024年Q3的销售目标是多少",系统返回了2023年的销售目标,因为相似度更高。
问题出在检索时没有过滤时间范围。后来在检索层加了元数据过滤,先按时间、权限过滤,再做向量检索。
from langchain_community.vectorstores import PGVector from langchain_openai import OpenAIEmbeddings from langchain_core.documents import Document def retrieve_with_filters(query: str, access_level: str, time_range: dict = None) -> List[Document]: """ 带权限和时间过滤的检索 """ # 构建过滤条件 filter_dict = {"access_level": access_level} if time_range: filter_dict["year"] = time_range.get("year") filter_dict["quarter"] = time_range.get("quarter") # 向量检索 + 元数据过滤 results = vectorstore.similarity_search_with_score( query=query, k=5, filter=filter_dict ) return [doc for doc, score in results]这段代码的关键是:过滤条件在向量检索之前应用,而不是检索之后再过滤。这样既保证了结果正确,也避免了不必要的计算。
---
落地项目:权限日志才是真正门槛
回到开头说的那个项目。
Demo阶段一切顺利:文档上传、向量化、检索、生成,流程跑得通。
生产环境上线后,问题一个接一个:
问题一:权限越界。
用户A登录系统后,通过API直接调用了用户B的文档。原因:检索时没有校验当前用户的权限级别,只校验了文档本身的权限。
修复方案:在检索层加入用户权限校验,所有查询必须带上用户ID,系统自动过滤超出权限的文档。
from functools import wraps from fastapi import Request, HTTPException def require_access_level(min_level: str): """权限装饰器:校验用户访问级别""" def decorator(func): @wraps(func) async def wrapper(request: Request, *args, **kwargs): user = request.state.user user_level = user.get_access_level() # 权限级别映射 level_map = { 'public': 0, 'internal': 1, 'team': 2, 'restricted': 3 } if level_map.get(user_level, 0) < level_map.get(min_level, 0): raise HTTPException( status_code=403, detail="权限不足" ) return await func(request, *args, **kwargs) return wrapper return decorator @app.post("/query") @require_access_level("internal") async def query_document(request: Request): # 业务逻辑 pass问题二:日志缺失。
某个用户反馈回答错误,但日志里没有他的查询记录。原因:日志只在生成阶段记录,检索阶段的查询被跳过了。
修复方案:统一日志格式,所有关键操作都记录请求ID、用户ID、查询内容、检索结果、生成结果。
import logging import uuid from contextlib import contextmanager logger = logging.getLogger("rag_pipeline") @contextmanager def trace_request(): """请求追踪上下文管理器""" request_id = str(uuid.uuid4()) start_time = time.time() logger.info(f"[{request_id}] 请求开始") try: yield request_id logger.info(f"[{request_id}] 请求成功, 耗时: {time.time()-start_time:.2f}s") except Exception as e: logger.error(f"[{request_id}] 请求失败: {e}, 耗时: {time.time()-start_time:.2f}s") raise finally: logger.info(f"[{request_id}] 请求结束") # 使用示例 @app.post("/query") async def query_document(request: Request): with trace_request() as request_id: # 所有操作都带request_id result = await process_query(request, request_id) return result问题三:可观测性为零。
系统挂了没人知道,用户投诉了才发现。原因:没有监控关键指标。
修复方案:接入Prometheus,监控以下指标:
- 请求成功率
- 平均响应时间
- 检索命中率
- 向量库连接数
- 错误类型分布
from prometheus_client import Counter, Histogram, start_http_server import time # 定义监控指标 REQUEST_COUNT = Counter( 'rag_request_total', 'RAG请求总数', ['endpoint', 'status'] ) REQUEST_LATENCY = Histogram( 'rag_request_latency_seconds', 'RAG请求延迟', ['endpoint'] ) ERROR_COUNT = Counter( 'rag_error_total', 'RAG错误总数', ['error_type'] ) # 在请求处理中记录指标 @app.middleware("http") async def monitor_requests(request: Request, call_next): start_time = time.time() response = await call_next(request) duration = time.time() - start_time REQUEST_LATENCY.labels(endpoint=request.url.path).observe(duration) REQUEST_COUNT.labels( endpoint=request.url.path, status=response.status_code ).inc() if response.status_code >= 500: ERROR_COUNT.labels(error_type="server_error").inc() return response # 启动PrometheusExporter start_http_server(8000)这三个问题,任何一个单独出现都不会致命。但组合在一起,系统就无法上线。
---
总结
大数据工程师转型大模型,真正需要补的不是算法,是工程化能力。
具体来说,有三件事比调参更重要:
第一,权限控制。 大模型应用处理的是企业数据,权限问题不是技术细节,是合规红线。Demo阶段可以忽略,生产环境必须解决。
第二,日志追踪。 没有日志就没有可观测性,没有可观测性就没有问题定位能力。建议从第一天就统一日志格式,不要等到出问题再补。
第三,可观测性。 监控不是上线后的事,是设计时就该考虑的事。Prometheus + Grafana 是标配,不要嫌麻烦。
最后给一个学习顺序建议:
1. 先掌握RAG的基础架构(检索+生成)
2. 再深入数据管道构建(切块、向量化、存储)
3. 最后补齐工程化能力(权限、日志、监控)
很多工程师的顺序是反的:先学Prompt工程,再学向量数据库,最后发现权限日志才是真正卡住的地方。
Demo能跑只是开始,权限日志才是真正的分水岭。
---
参考资料
- LangChain官方文档:https://python.langchain.com/
- pgvector文档:https://github.com/pgvector/pgvector
- Prometheus官方文档:https://prometheus.io/docs/
资料展示
下面是我整理的AI大模型学习资料和工具包预览,适合收藏后按主题逐步学习。
如果你想看完整资料目录,可以在评论区留言「资料」;也欢迎告诉我你更关注AI大模型里的哪类内容。