大数据转大模型:Demo能跑只是开始,权限日志才是真正的分水岭

📅 2026/8/4 12:34:51
大数据转大模型: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_keyTrue) content Column(Text, nullableFalse) embedding Column(Vector(1536)) # OpenAI ada-002 维度 doc_id Column(String, nullableFalse) source_system Column(String, nullableFalse) access_level Column(String, nullableFalse) def __repr__(self): return fDocumentChunk(id{self.id}, doc_id{self.doc_id}) engine create_engine(postgresqlpsycopg2://user:passhost/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[i1] if i1 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( queryquery, k5, filterfilter_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_code403, 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(endpointrequest.url.path).observe(duration) REQUEST_COUNT.labels( endpointrequest.url.path, statusresponse.status_code ).inc() if response.status_code 500: ERROR_COUNT.labels(error_typeserver_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/pgvectorPrometheus官方文档https://prometheus.io/docs/资料展示下面是我整理的AI大模型学习资料和工具包预览适合收藏后按主题逐步学习。如果你想看完整资料目录可以在评论区留言「资料」也欢迎告诉我你更关注AI大模型里的哪类内容。