问题
在AI回答采集系统中,原始数据是后续所有分析的根基,这句话怎么强调都不过分。如果原始回答存储不规范,后续一旦需要复查某个指标的计算过程、重算某个时间段的统计结果,或者排查数据异常,就会陷入“数据在哪?怎么存的?用了哪个模型?哪个问题?”的困境。这种局面,谁都不想遇到。

本文只聚焦一个核心问题:如何设计原始回答的存储方案,保证每一次采集结果都可追溯、可复查、可重算。
整体流程
关键点在于:原始响应在解析之前就被完整保存。这意味着,任何时候都可以从原始数据重新开始,而不必担心中间环节的丢失或变形。
数据结构设计
下面给出的是示例表结构,具体字段请根据实际业务场景灵活调整。核心思路是分层清晰,每一层各司其职。
1. 采集任务表(task)
这张表记录每次采集任务的基本信息,可以说是整个采集流程的“总控”。
| 字段名 | 类型 | 说明 |
|---|---|---|
| task_id | VARCHAR(64) | 主键,任务唯一ID |
| task_name | VARCHAR(255) | 任务名称 |
| platform | VARCHAR(64) | 采集平台,如openai、claude |
| model | VARCHAR(64) | 模型名称,如gpt-4o、claude-3 |
| question_set_id | VARCHAR(64) | 问题集ID |
| created_at | DATETIME | 任务创建时间 |
| status | TINYINT | 任务状态:0-待执行,1-执行中,2-完成,3-失败 |
2. 原始回答表(raw_response)
这张表是核心中的核心,它保存每次调用的完整响应,绝不提前动手解析。
| 字段名 | 类型 | 说明 |
|---|---|---|
| id | BIGINT AUTO_INCREMENT | 自增主键 |
| task_id | VARCHAR(64) | 关联任务ID |
| question_id | VARCHAR(64) | 问题ID |
| request_params | JSON | 请求参数,包括temperature、max_tokens等 |
| raw_response | TEXT | 模型返回的原始响应(JSON字符串) |
| http_status | INT | HTTP状态码 |
| latency_ms | INT | 调用耗时,单位毫秒 |
| model_name | VARCHAR(64) | 实际使用的模型(可能和任务配置不同) |
| created_at | DATETIME | 记录创建时间 |
| checksum | VARCHAR(64) | 原始响应的哈希值,用于校验完整性 |
| retry_count | TINYINT DEFAULT 0 | 重试次数 |
设计说明:
raw_response字段保存完整的JSON字符串,不提前解析。为什么?因为一旦解析逻辑变了,原始数据还在,可以重新解析,这才是真正的“原始”。checksum字段用于检测数据是否被意外修改,推荐使用SHA256,简单有效。request_params记录实际请求参数,因为同一任务可能因重试而参数略有不同,这一点容易被忽略,但实际复盘时很重要。- 主键使用自增ID,但业务查询通常通过
task_id和question_id进行,所以在这两个字段上建立联合索引是必须的。 retry_count用于标识重试次数,避免重复数据无法区分,这在后面会再提到。
3. 解析结果表(parsed_result)
这张表保存从原始回答中解析出的结构化数据,目的是方便快速查询,但它和原始数据是分离的。
| 字段名 | 类型 | 说明 |
|---|---|---|
| id | BIGINT AUTO_INCREMENT | 自增主键 |
| raw_response_id | BIGINT | 关联原始回答ID |
| content | TEXT | 解析后的回答文本 |
| is_valid | TINYINT | 是否有效回答(0-无效,1-有效) |
| mentioned_entities | JSON | 提及的实体列表 |
| recommended_entities | JSON | 推荐的实体列表 |
| parsed_at | DATETIME | 解析时间 |
| parser_version | VARCHAR(32) | 解析器版本号 |
设计说明:
- 解析结果与原始数据分离,这个设计非常关键。解析逻辑升级后,可以重新解析原始数据,而不影响历史记录。
parser_version字段记录解析器版本,当指标重算时,可以判断是否需要重新解析,避免重复劳动。
核心实现
下面给出关键实现片段,以Python和MySQL为例。代码本身不复杂,但设计思路值得注意。
保存原始响应
import hashlib
import json
from datetime import datetime
import mysql.connector
def sa ve_raw_response(conn, task_id, question_id, request_params, response, http_status, latency_ms, model_name, retry_count=0):
"""
保存原始响应到数据库。
:param conn: 数据库连接
:param task_id: 任务ID
:param question_id: 问题ID
:param request_params: 请求参数字典
:param response: 原始响应字符串
:param http_status: HTTP状态码
:param latency_ms: 调用耗时(毫秒)
:param model_name: 模型名称
:param retry_count: 重试次数
:return: 插入记录的ID
"""
checksum = hashlib.sha256(response.encode('utf-8')).hexdigest()
sql = """
INSERT INTO raw_response
(task_id, question_id, request_params, raw_response, http_status, latency_ms, model_name, created_at, checksum, retry_count)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
"""
params = (
task_id,
question_id,
json.dumps(request_params),
response,
http_status,
latency_ms,
model_name,
datetime.utcnow(),
checksum,
retry_count
)
try:
cursor = conn.cursor()
cursor.execute(sql, params)
conn.commit()
return cursor.lastrowid
except mysql.connector.Error as err:
conn.rollback()
raise RuntimeError(f"Failed to insert raw response: {err}")
finally:
cursor.close()
说明:
- 在调用模型后立即保存,不经过任何业务处理,这是“原始”二字的底线。
- 响应内容作为字符串直接存储,避免序列化问题,简单粗暴但有效。
- 校验和用于后续数据完整性校验,这是事后复查的保障。
- 使用参数化查询防止SQL注入,并包含异常处理和事务回滚,这些在工程实践中必不可少。
数据完整性校验
def verify_raw_response(conn, record_id):
"""
校验原始响应的完整性。
:param conn: 数据库连接
:param record_id: 记录ID
:return: True 如果校验通过,否则抛出异常
"""
cursor = conn.cursor(dictionary=True)
cursor.execute("SELECT raw_response, checksum FROM raw_response WHERE id = %s", (record_id,))
record = cursor.fetchone()
cursor.close()
if not record:
raise ValueError(f"Record {record_id} not found")
expected_checksum = hashlib.sha256(record['raw_response'].encode('utf-8')).hexdigest()
if expected_checksum != record['checksum']:
raise DataIntegrityError(f"Record {record_id} checksum mismatch")
return True
class DataIntegrityError(Exception):
pass
验证方法
1. 数据完整性验证
怎么验证?简单,随机抽取一批原始记录,执行以下SQL查询并计算校验和对比:
SELECT id, raw_response, checksum FROM raw_response ORDER BY RAND() LIMIT 100;
然后对每条记录的 raw_response 计算SHA256,与 checksum 字段对比,应该全部一致。如果发现不一致,那数据就有问题了,需要立刻排查。
2. 可追溯性验证
从最终指标反查原始回答,是检验可追溯性的好方法。比如,假设已知指标来源于 task_id='task_001' 和 question_id='q_002':
SELECT raw_response FROM raw_response WHERE task_id='task_001' AND question_id='q_002';
应该能完整复现当时的响应内容,一字不差。
3. 重算一致性验证
使用同一批原始数据,用新版本解析器重新计算指标:
# 伪代码示例
def recalculate(task_id, new_parser_version):
cursor.execute("SELECT id, raw_response FROM raw_response WHERE task_id = %s", (task_id,))
for row in cursor:
parsed = parse(row['raw_response'], version=new_parser_version)
# 更新 parsed_result 表
结果应与旧版本在相同口径下一致——除非解析逻辑有明确变更,否则不应该出现意外偏差。
常见问题
1. 原始响应过大
某些模型返回的响应可能包含长文本甚至图片,导致单条记录超过MySQL TEXT字段上限(65,535字节)。这怎么办?
解决方法:
- 使用MEDIUMTEXT(最多16MB)或LONGTEXT(最多4GB)。适合中小规模数据,查询方便,但会增加数据库存储压力。
- 将响应内容存储在对象存储(如S3)中,数据库中只保存文件路径。适合大规模数据,但每次读取需要额外网络开销,查询延迟增加。
权衡:如果数据量不大(日均万条以内),直接使用MEDIUMTEXT更简单;如果数据量巨大或需要长期归档,对象存储更经济。没有绝对的对错,看场景。
2. 重试导致重复数据
当请求失败重试时,同一问题可能产生多条原始记录。这很容易造成混乱。
解决方法:
- 在
raw_response表中增加retry_count字段,标识第几次重试。 - 在解析时只使用最后一次成功的结果,或根据业务规则选择(如取最大
retry_count且http_status=200的记录)。 - 也可以使用
task_id + question_id + retry_count作为唯一约束,防止重复插入。
3. 解析器版本升级后,历史数据需要重新解析
这是一个必然会遇到的问题。解析器升级了,历史数据怎么办?
解决方法:
- 在
parsed_result表中记录parser_version。 - 编写脚本,遍历
raw_response表中所有id,对parser_version低于当前版本的记录重新解析并更新。
def reparse_all(conn, new_version):
cursor = conn.cursor(dictionary=True)
cursor.execute("""
SELECT r.id, r.raw_response
FROM raw_response r
LEFT JOIN parsed_result p ON r.id = p.raw_response_id
WHERE p.parser_version IS NULL OR p.parser_version < %s
""", (new_version,))
for row in cursor:
parsed = parse(row['raw_response'], version=new_version)
# 更新或插入 parsed_result
优化建议
- 索引优化:在
raw_response表的task_id和question_id上建立联合索引,加速查询。这是最基础的优化,但常常被忽略。 - 数据归档:对于超过一定时间(如90天)的原始数据,可以迁移到冷存储,但保留元数据以便追溯。这样既能节省成本,又不影响历史查询。
- 分区表:如果数据量极大,可以按时间分区,例如按月分区。这样便于管理和清理,也能提升查询性能。
- 幂等设计:在保存原始响应时,使用
task_id + question_id + retry_count作为唯一约束,防止重复插入。这能从根本上避免数据混乱。
适用边界
- 本文的设计适用于中小规模(日均万条以内)的AI回答采集系统。如果数据量更大,需要考虑分布式方案。
- 如果数据量极大或需要高并发写入,建议使用分布式数据库或消息队列异步写入,避免单点瓶颈。
- 如果对存储成本敏感,可考虑压缩原始响应后再存储,但这会增加读取时的计算开销。
- 本文未涉及数据加密、访问控制等安全措施,生产环境需额外补充,这一点务必注意。
