在 AI 回答监测系统中,问题库是数据生产的起点。但问题不是越多越好——如何把几百条问题按场景、品牌、平台组合成可管理的批次?如何确保同一批次的任务在合理的时间窗口内完成?又如何把分批采集的结果有序写入样本库?本文围绕问题批次生成、任务调度、失败重试和样本入库四个环节,拆解一条从问题配置到回答样本库的完整数据链路。
问题库不是简单的“一堆问题”的集合。在工程上,它承担三个角色:
因此,问题库的设计需要兼顾可配置性(方便新增和维护)和稳定性(保证跨周期一致性)。
CREATE TABLE question_bank (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
question_id VARCHAR(32) NOT NULL UNIQUE COMMENT '问题唯一标识, 如 Q20260701_001',
scene_type VARCHAR(30) NOT NULL COMMENT '场景分类',
query_text TEXT NOT NULL COMMENT '提问原文',
target_keywords VARCHAR(255) COMMENT '监测关键词, 逗号分隔',
intent_description VARCHAR(255) COMMENT '意图说明, 供运营参考',
is_active TINYINT(1) DEFAULT 1 COMMENT '是否启用',
version INT DEFAULT 1 COMMENT '版本号, 修改后递增',
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
INDEX idx_scene_active (scene_type, is_active),
INDEX idx_question_id (question_id)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;关键设计点:
question_id 使用业务主键而非自增 ID,保证跨环境(开发/测试/生产)一致性。version 字段记录问题修改次数。问题文本变更时版本号递增,旧版本数据仍可通过 question_id + version 关联到历史样本。is_active 控制问题上下线,下线后调度器跳过该问题,但不删除历史样本。-- 问题使用记录表, 跟踪每条问题被调度过的历史
CREATE TABLE question_usage_log (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
question_id VARCHAR(32) NOT NULL,
batch_id VARCHAR(64) NOT NULL COMMENT '批次ID',
platform VARCHAR(30) NOT NULL,
round INT NOT NULL COMMENT '采样轮次',
sub_task_id VARCHAR(64) NOT NULL COMMENT '对应的采集子任务ID',
status VARCHAR(20) DEFAULT 'pending' COMMENT 'pending/success/failed',
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
INDEX idx_batch (batch_id),
INDEX idx_question (question_id)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;question_usage_log 是问题库与采集任务之间的关联表。每次调度都写入一条记录,用于追踪“某问题在某批次中是否被成功采集”。
一次完整的监测任务通常涉及:N 个品牌 × M 个场景 × P 个平台 × R 轮采样。直接组合会产生大量任务,需要按“批次”分批调度。
批次维度 | 说明 | 示例 |
|---|---|---|
按品牌 | 同一品牌的所有问题/平台/轮次组成一个批次 | 品牌A 全部采集任务 = 1 个批次 |
按场景 | 同一场景的所有问题/品牌/平台/轮次组成一个批次 | 推荐决策场景全部任务 = 1 个批次 |
按时间窗口 | 固定时间段内能完成的任务量 | 每 10 分钟调度 50 条任务 |
推荐按品牌分批,原因是:后续分析通常以品牌为单位,同一批次的品牌数据时间窗口接近,对比公平性更好。
CREATE TABLE task_batch (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
batch_id VARCHAR(64) NOT NULL UNIQUE COMMENT '批次ID, 如 BATCH_20260707_001',
brand_name VARCHAR(100) NOT NULL,
scene_types JSON COMMENT '包含的场景列表',
platforms JSON COMMENT '包含的平台列表',
total_questions INT DEFAULT 0 COMMENT '本批次问题总数',
total_sub_tasks INT DEFAULT 0 COMMENT '子任务总数(含多轮)',
completed_sub_tasks INT DEFAULT 0 COMMENT '已完成子任务数',
failed_sub_tasks INT DEFAULT 0 COMMENT '失败子任务数',
batch_status VARCHAR(20) DEFAULT 'pending' COMMENT 'pending/running/completed/partial',
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
completed_at TIMESTAMP NULL,
INDEX idx_brand (brand_name),
INDEX idx_status (batch_status)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;def create_batch(brand_name: str, scene_types: list, platforms: list, rounds: int = 3) -> str:
batch_id = f"BATCH_{datetime.now().strftime('%Y%m%d%H%M%S')}_{brand_name[:4].upper()}"
# 从问题库拉取该品牌、指定场景的启用问题
questions = db.query(
"SELECT question_id, scene_type, query_text FROM question_bank "
"WHERE is_active = 1 AND scene_type IN %s",
scene_types
)
total_sub_tasks = 0
for q in questions:
for platform in platforms:
for r in range(1, rounds + 1):
sub_task = {
"sub_task_id": f"{batch_id}_Q{q['question_id']}_{platform}_R{r}",
"batch_id": batch_id,
"question_id": q["question_id"],
"query_text": q["query_text"],
"scene_type": q["scene_type"],
"platform": platform,
"round": r,
"brand_name": brand_name,
}
send_to_queue(sub_task)
total_sub_tasks += 1
db.insert("task_batch", {
"batch_id": batch_id,
"brand_name": brand_name,
"total_questions": len(questions),
"total_sub_tasks": total_sub_tasks,
"batch_status": "running",
})
return batch_id同一批次内的子任务理论上可以全部并发执行,但实际受限于 AI 平台的 API 限流。调度策略分两级:
控制层级 | 维度 | 机制 |
|---|---|---|
批次间 | 品牌 | 串行或低并发,保证不同品牌的任务时间窗口接近 |
批次内 | 平台 | 令牌桶限流,每平台每秒最多 N 次请求 |
class BatchScheduler:
def __init__(self):
self.platform_limiters = {
"deepseek": TokenBucket(rate=5, capacity=10),
"kimi": TokenBucket(rate=3, capacity=6),
"tongyi": TokenBucket(rate=5, capacity=10),
}
def dispatch(self, sub_task: dict):
limiter = self.platform_limiters.get(sub_task["platform"])
if limiter:
limiter.consume() # 阻塞直到获取令牌
execute_sub_task(sub_task)每个子任务完成后更新批次计数器:
def on_sub_task_complete(sub_task: dict, success: bool):
batch_id = sub_task["batch_id"]
if success:
db.execute(
"UPDATE task_batch SET completed_sub_tasks = completed_sub_tasks + 1 WHERE batch_id = %s",
batch_id
)
else:
db.execute(
"UPDATE task_batch SET failed_sub_tasks = failed_sub_tasks + 1 WHERE batch_id = %s",
batch_id
)
# 检查批次是否全部完成
batch = db.query("SELECT * FROM task_batch WHERE batch_id = %s", batch_id)
if batch["completed_sub_tasks"] + batch["failed_sub_tasks"] == batch["total_sub_tasks"]:
final_status = "completed" if batch["failed_sub_tasks"] == 0 else "partial"
db.execute(
"UPDATE task_batch SET batch_status = %s, completed_at = NOW() WHERE batch_id = %s",
final_status, batch_id
)失败分两类处理:
RETRY_STRATEGY = {
"temp_error": {"max_retries": 3, "backoff": "exponential"}, # 网络超时、限流
"perm_error": {"max_retries": 0}, # API Key 失效、参数错误
}
def handle_failure(sub_task: dict, error_type: str):
strategy = RETRY_STRATEGY.get(error_type, {"max_retries": 0})
retry_count = sub_task.get("retry_count", 0)
if retry_count < strategy["max_retries"]:
sub_task["retry_count"] = retry_count + 1
delay = 60 * (2 ** retry_count) if strategy["backoff"] == "exponential" else 30
requeue_with_delay(sub_task, delay)
else:
# 写入失败记录,标记批次为 partial
db.insert("question_usage_log", {
"question_id": sub_task["question_id"],
"batch_id": sub_task["batch_id"],
"status": "failed",
"error": error_type,
})
on_sub_task_complete(sub_task, success=False)每个子任务完成后,原始回答和提取指标分别入库:
CREATE TABLE answer_sample (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
sub_task_id VARCHAR(64) NOT NULL UNIQUE,
batch_id VARCHAR(64) NOT NULL,
question_id VARCHAR(32) NOT NULL,
brand_name VARCHAR(100) NOT NULL,
platform VARCHAR(30) NOT NULL,
scene_type VARCHAR(30) NOT NULL,
round INT NOT NULL,
raw_response LONGTEXT COMMENT 'AI返回的原始回答',
response_format VARCHAR(20) COMMENT 'markdown/json/text',
response_length INT DEFAULT 0 COMMENT '回答字符数',
response_time_ms INT COMMENT 'API响应耗时',
cos_url VARCHAR(512) COMMENT 'COS存储链接',
is_valid TINYINT(1) DEFAULT 1 COMMENT '是否有效样本',
invalid_reason VARCHAR(50) COMMENT '无效原因: timeout/refuse/error',
collected_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
INDEX idx_batch (batch_id),
INDEX idx_brand_scene (brand_name, scene_type),
INDEX idx_question_round (question_id, round)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;CREATE TABLE answer_metric (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
sub_task_id VARCHAR(64) NOT NULL UNIQUE,
batch_id VARCHAR(64) NOT NULL,
brand_name VARCHAR(100) NOT NULL,
platform VARCHAR(30) NOT NULL,
scene_type VARCHAR(30) NOT NULL,
question_id VARCHAR(32) NOT NULL,
round INT NOT NULL,
is_mentioned TINYINT(1) DEFAULT 0,
is_recommended TINYINT(1) DEFAULT 0,
recommendation_rank INT DEFAULT 0,
has_citation TINYINT(1) DEFAULT 0,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
INDEX idx_brand_scene (brand_name, scene_type)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;一条子任务对应两条数据库写入(原始样本 + 指标),需要保证原子性:
def save_sub_task_result(sub_task: dict, raw_response: str, metrics: dict):
db.begin()
try:
db.insert("answer_sample", {...})
db.insert("answer_metric", {...})
db.insert("question_usage_log", {..., "status": "success"})
db.commit()
except Exception as e:
db.rollback()
raise样本库建成后,常见查询场景:
-- 某问题在多个平台的回答一致性
SELECT platform, is_mentioned, is_recommended, recommendation_rank
FROM answer_metric
WHERE question_id = 'Q20260701_001'
AND brand_name = '品牌A'
ORDER BY platform, round;
-- 某批次整体完成情况
SELECT
brand_name,
COUNT(*) AS total_samples,
SUM(is_valid) AS valid_samples,
SUM(is_mentioned) AS mentioned,
ROUND(SUM(is_mentioned)*100.0/SUM(is_valid), 1) AS mention_rate
FROM answer_metric
WHERE batch_id = 'BATCH_20260707_001_品牌A'
GROUP BY brand_name;1. 问题版本化是趋势分析的前提
问题文本一旦修改,新旧采集结果就不能直接对比。version 字段 + question_usage_log 关联,让每次采集都绑定到具体的问题版本。趋势分析时按 question_id + version 分组,避免跨版本数据混在一起。
2. 批次超时要有兜底
生产环境中,个别子任务可能长时间不返回(AI 平台排队、网络 hang 住)。批次不应无限等待。建议设置批次超时(如 30 分钟),超时后强制标记 partial 并释放等待资源,失败子任务记录在案,后续单独补采。
3. 问题库需要定期巡检
问题库中的问题可能因业务变化而“过时”(如提到的产品已下架)。建议每月巡检一次 question_usage_log,统计每条问题的最近成功采集时间。超过 30 天未成功采集的问题,自动标记为 is_active=0 并通知运营确认。
问题库是 AI 回答监测系统的“生产计划表”。批次生成将大规模任务拆分为可控单元,两级限流兼顾效率与平台友好性,分表存储让原始样本和提取指标各得其所,问题版本化保证跨周期数据的可比性。这四层设计让问题库从“静态配置”升级为“可调度、可追踪、可复盘”的任务编排引擎。
这套方案已在多个消费品牌和企业服务的多平台 AI 可见度监测中实际运行,支持每日数百条问题的自动化调度与样本归集。开发者可在此基础上,扩展问题库的自动生成能力(如基于大模型按场景批量生成提问文本),进一步降低运营维护成本。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。