采集AI模型回答时,如果只存储最终解析后的结果,一旦发现数据异常或需要重新计算指标,往往无法还原原始回答。例如,模型返回了格式错误的JSON,或者后续需要提取新的字段,这时原始响应就是唯一可靠的依据。本文面向需要审计或长期维护的AI采集系统,设计一种原始数据存储方案,重点解决可追溯性问题。方案不涉及实时流处理或大规模分布式存储,仅讨论单机或中小规模场景下的实现思路。
存储方案设计
存储格式
原始回答通常以文本形式存储。对于结构化响应(如OpenAI Chat Completion API返回的JSON,包含choices、usage等字段),直接存储完整JSON可保留所有字段,便于后续解析。对于流式响应,需在客户端拼接完整内容后存储为纯文本。以下是一个示例JSON结构:
{"id":"chatcmpl-abc123","object":"chat.completion","created":1722163200,"model":"gpt-4-0613","choices":[{"index":0,"message":{"role":"assistant","content":"原始回答内容"},"finish_reason":"stop"}],"usage":{"prompt_tokens":10,"completion_tokens":20,"total_tokens":30}}对于非结构化文本(如Markdown),存储为纯文本,但建议同时记录原始请求和响应的完整内容。
元数据字段
每条记录需要包含以下元数据:
- 采集时间:精确到毫秒的时间戳
- 平台:如OpenAI、Claude、本地模型等
- 模型名称:如gpt-4-0613、claude-3-opus
- 问题:用户输入的原始问题
- 请求参数:temperature、max_tokens等
- 响应状态:成功、超时、错误码等
- 版本号:用于追踪数据schema变更
数据库表设计
使用关系型数据库(如PostgreSQL)存储元数据,原始回答内容存储在文件系统或对象存储中,数据库只保存路径。表结构如下:
CREATETABLEraw_responses(id BIGSERIALPRIMARYKEY,request_id UUIDNOTNULLUNIQUE,collected_at TIMESTAMPTZNOTNULLDEFAULTNOW(),platformVARCHAR(50)NOTNULL,model_nameVARCHAR(100)NOTNULL,questionTEXTNOTNULL,request_params JSONB,response_statusVARCHAR(20)NOTNULL,raw_content_pathTEXTNOTNULL,content_formatVARCHAR(20)NOTNULLDEFAULT'json',schema_versionINTEGERNOTNULLDEFAULT1,created_at TIMESTAMPTZNOTNULLDEFAULTNOW());CREATEINDEXidx_raw_responses_collected_atONraw_responses(collected_at);CREATEINDEXidx_raw_responses_platformONraw_responses(platform);原始回答文件按日期和request_id组织,例如:/data/raw/2026-07-28/{request_id}.json。
核心实现过程
以下Python代码展示采集并存储原始回答的流程,包含异常处理和重试逻辑:
importjsonimportuuidfromdatetimeimportdatetimefrompathlibimportPathimportpsycopg2frompsycopg2importpoolfromtenacityimportretry,stop_after_attempt,wait_fixedclassRawResponseStorage:def__init__(self,db_config,storage_root,max_retries=3):self.pool=pool.SimpleConnectionPool(1,10,**db_config)self.storage_root=Path(storage_root)self.max_retries=max_retries@retry(stop=stop_after_attempt(3),wait=wait_fixed(2))defstore(self,platform,model,question,params,response,status):request_id=str(uuid.uuid4())collected_at=datetime.utcnow()date_str=collected_at.strftime("%Y-%m-%d")# 保存原始内容到文件file_dir=self.storage_root/date_str file_dir.mkdir(parents=True,exist_ok=True)file_path=file_dir/f"{request_id}.json"raw_data={"request_id":request_id,"collected_at":collected_at.isoformat(),"platform":platform,"model":model,"question":question,"params":params,"response":response,"status":status}try:withopen(file_path,'w',encoding='utf-8')asf:json.dump(raw_data,f,ensure_ascii=False,indent=2)exceptIOErrorase:raiseRuntimeError(f"Failed to write raw content file:{e}")# 插入元数据到数据库conn=self.pool.getconn()try:withconn.cursor()ascur:cur.execute(""" INSERT INTO raw_responses (request_id, collected_at, platform, model_name, question, request_params, response_status, raw_content_path, content_format) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s) """,(request_id,collected_at,platform,model,question,json.dumps(params),status,str(file_path),'json'))conn.commit()exceptExceptionase:conn.rollback()# 清理孤儿文件iffile_path.exists():file_path.unlink()raiseRuntimeError(f"Database insert failed:{e}")finally:self.pool.putconn(conn)returnrequest_id# 调用示例# storage = RawResponseStorage(# db_config={"host": "localhost", "dbname": "ai_collect", "user": "user", "password": "pass"},# storage_root="/data/raw"# )# request_id = storage.store(# platform="openai",# model="gpt-4-0613",# question="什么是可追溯性?",# params={"temperature": 0.7, "max_tokens": 100},# response={"choices": [{"message": {"content": "可追溯性是指..."}}]},# status="success"# )代码说明:
- 使用连接池管理数据库连接,避免频繁创建连接。
- 通过tenacity库实现重试机制,网络超时或临时故障时可自动重试。
- 异常处理:文件写入失败抛出异常;数据库插入失败时回滚并删除已写入的文件,防止孤儿文件。
- 调用示例展示了如何从API获取response并传入store方法。
版本管理
当存储格式或元数据字段发生变化时,通过schema_version字段区分。旧版本数据仍可读取,但解析逻辑需兼容。建议在代码中维护一个版本映射表,根据版本号选择对应的反序列化方法。例如:
VERSION_PARSERS={1:parse_v1,2:parse_v2,}defparse_raw(raw_path,version):parser=VERSION_PARSERS.get(version)ifnotparser:raiseValueError(f"Unsupported schema version:{version}")returnparser(raw_path)验证结果
正常情况下,执行存储后数据库应出现一条记录,且文件系统中存在对应的JSON文件。可通过以下SQL验证:
SELECTid,request_id,collected_at,platform,model_name,response_statusFROMraw_responsesWHEREcollected_at>='2026-07-28'ORDERBYcollected_atDESCLIMIT10;同时检查文件路径是否存在:
ls-la/data/raw/2026-07-28/常见问题与避坑
- 文件系统性能:高并发时文件IO可能成为瓶颈,可考虑使用对象存储(如S3)替代本地文件。
- 数据一致性:写入文件成功但数据库插入失败会导致孤儿文件。上述代码通过异常处理回滚并删除文件,但更可靠的做法是使用两阶段提交或先写数据库再写文件,并增加定期清理任务。
- 存储空间:原始回答可能很大,需要设置保留策略,例如只保留30天,或归档到冷存储。
- 敏感信息:原始回答可能包含用户隐私,存储前需脱敏或加密。
总结
本方案通过分离元数据和内容、记录完整请求上下文、引入版本号,实现了AI回答的可追溯性。在中小规模采集场景下表现良好,但高并发时文件IO可能成为瓶颈,建议使用对象存储。版本管理通过schema_version实现,但需注意旧版本数据的兼容性。实际部署时需根据并发量、存储成本和合规要求调整具体实现。