游乐游手机版
首页/AI热点日报/热点详情

AI回答采集原始数据存储表结构设计与Python实现

类型:热点整理2026-07-23
提出AI回答采集系统中原始数据存储方案,设计任务表、原始回答表和解析结果表,确保每次采集结果可追溯、可复查、可重算。核心在于完整保存原始响应,通过校验和与版本控制保障数据完整性,并支持重算与解析升级。

问题

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

AI回答采集:原始数据存储的表结构设计与Python实现

本文只聚焦一个核心问题:如何设计原始回答的存储方案,保证每一次采集结果都可追溯、可复查、可重算。

整体流程

flowchart TD
    A[采集任务] --> B[构造请求]
    B --> C[调用AI模型]
    C --> D[保存原始响应]
    D --> E[解析与归一化]
    E --> F[指标计算]
    F --> G[结果展示]
    D --> H[原始数据存储]
    H --> I[复查与重算]

关键点在于:原始响应在解析之前就被完整保存。这意味着,任何时候都可以从原始数据重新开始,而不必担心中间环节的丢失或变形。

数据结构设计

下面给出的是示例表结构,具体字段请根据实际业务场景灵活调整。核心思路是分层清晰,每一层各司其职。

1. 采集任务表(task)

这张表记录每次采集任务的基本信息,可以说是整个采集流程的“总控”。

字段名类型说明
task_idVARCHAR(64)主键,任务唯一ID
task_nameVARCHAR(255)任务名称
platformVARCHAR(64)采集平台,如openai、claude
modelVARCHAR(64)模型名称,如gpt-4o、claude-3
question_set_idVARCHAR(64)问题集ID
created_atDATETIME任务创建时间
statusTINYINT任务状态:0-待执行,1-执行中,2-完成,3-失败

2. 原始回答表(raw_response)

这张表是核心中的核心,它保存每次调用的完整响应,绝不提前动手解析。

字段名类型说明
idBIGINT AUTO_INCREMENT自增主键
task_idVARCHAR(64)关联任务ID
question_idVARCHAR(64)问题ID
request_paramsJSON请求参数,包括temperature、max_tokens等
raw_responseTEXT模型返回的原始响应(JSON字符串)
http_statusINTHTTP状态码
latency_msINT调用耗时,单位毫秒
model_nameVARCHAR(64)实际使用的模型(可能和任务配置不同)
created_atDATETIME记录创建时间
checksumVARCHAR(64)原始响应的哈希值,用于校验完整性
retry_countTINYINT DEFAULT 0重试次数

设计说明

  • raw_response 字段保存完整的JSON字符串,不提前解析。为什么?因为一旦解析逻辑变了,原始数据还在,可以重新解析,这才是真正的“原始”。
  • checksum 字段用于检测数据是否被意外修改,推荐使用SHA256,简单有效。
  • request_params 记录实际请求参数,因为同一任务可能因重试而参数略有不同,这一点容易被忽略,但实际复盘时很重要。
  • 主键使用自增ID,但业务查询通常通过 task_idquestion_id 进行,所以在这两个字段上建立联合索引是必须的。
  • retry_count 用于标识重试次数,避免重复数据无法区分,这在后面会再提到。

3. 解析结果表(parsed_result)

这张表保存从原始回答中解析出的结构化数据,目的是方便快速查询,但它和原始数据是分离的。

字段名类型说明
idBIGINT AUTO_INCREMENT自增主键
raw_response_idBIGINT关联原始回答ID
contentTEXT解析后的回答文本
is_validTINYINT是否有效回答(0-无效,1-有效)
mentioned_entitiesJSON提及的实体列表
recommended_entitiesJSON推荐的实体列表
parsed_atDATETIME解析时间
parser_versionVARCHAR(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_counthttp_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_idquestion_id 上建立联合索引,加速查询。这是最基础的优化,但常常被忽略。
  • 数据归档:对于超过一定时间(如90天)的原始数据,可以迁移到冷存储,但保留元数据以便追溯。这样既能节省成本,又不影响历史查询。
  • 分区表:如果数据量极大,可以按时间分区,例如按月分区。这样便于管理和清理,也能提升查询性能。
  • 幂等设计:在保存原始响应时,使用 task_id + question_id + retry_count 作为唯一约束,防止重复插入。这能从根本上避免数据混乱。

适用边界

  • 本文的设计适用于中小规模(日均万条以内)的AI回答采集系统。如果数据量更大,需要考虑分布式方案。
  • 如果数据量极大或需要高并发写入,建议使用分布式数据库或消息队列异步写入,避免单点瓶颈。
  • 如果对存储成本敏感,可考虑压缩原始响应后再存储,但这会增加读取时的计算开销。
  • 本文未涉及数据加密、访问控制等安全措施,生产环境需额外补充,这一点务必注意。
来源:https://segmentfault.com/a/1190000048066131

相关热点

继续查看同栏目近期热点。

延伸阅读

补充最近整理过的热点入口。