首页
学习
活动
专区
圈层
工具
发布
社区首页 >专栏 >AI 回答监测问题库如何做任务调度、批次执行和结果入库?

AI 回答监测问题库如何做任务调度、批次执行和结果入库?

原创
作者头像
AI增长技术研究院
发布2026-07-07 11:16:04
发布2026-07-07 11:16:04
1330
举报

在 AI 回答监测系统中,问题库是数据生产的起点。但问题不是越多越好——如何把几百条问题按场景、品牌、平台组合成可管理的批次?如何确保同一批次的任务在合理的时间窗口内完成?又如何把分批采集的结果有序写入样本库?本文围绕问题批次生成、任务调度、失败重试和样本入库四个环节,拆解一条从问题配置到回答样本库的完整数据链路。


一、问题库在监测系统中的定位

问题库不是简单的“一堆问题”的集合。在工程上,它承担三个角色:

  1. 数据生产入口:每条问题对应一组采集任务,问题库的结构决定了后续任务编排的复杂度。
  2. 场景分类载体:每条问题带有场景标签(推荐决策、对比分析、风险判断等),这个标签贯穿采集、诊断、分析全链路。
  3. 采样基线锚点:同一组问题在不同周期的复用,是趋势对比的数据基础。问题一旦上线就不宜随意修改,否则前后数据不可比。

因此,问题库的设计需要兼顾可配置性(方便新增和维护)和稳定性(保证跨周期一致性)。


二、问题库表结构设计

代码语言:javascript
复制
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 控制问题上下线,下线后调度器跳过该问题,但不删除历史样本。

代码语言:javascript
复制
-- 问题使用记录表, 跟踪每条问题被调度过的历史
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 轮采样。直接组合会产生大量任务,需要按“批次”分批调度。

3.1 批次划分策略

批次维度

说明

示例

按品牌

同一品牌的所有问题/平台/轮次组成一个批次

品牌A 全部采集任务 = 1 个批次

按场景

同一场景的所有问题/品牌/平台/轮次组成一个批次

推荐决策场景全部任务 = 1 个批次

按时间窗口

固定时间段内能完成的任务量

每 10 分钟调度 50 条任务

推荐按品牌分批,原因是:后续分析通常以品牌为单位,同一批次的品牌数据时间窗口接近,对比公平性更好。

3.2 批次元数据表

代码语言:javascript
复制
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;
3.3 批次生成逻辑

代码语言:javascript
复制
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

四、任务调度:批次内并发与跨批次隔离
4.1 批次级并发控制

同一批次内的子任务理论上可以全部并发执行,但实际受限于 AI 平台的 API 限流。调度策略分两级:

控制层级

维度

机制

批次间

品牌

串行或低并发,保证不同品牌的任务时间窗口接近

批次内

平台

令牌桶限流,每平台每秒最多 N 次请求

代码语言:javascript
复制
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)
4.2 批次状态追踪

每个子任务完成后更新批次计数器:

代码语言:javascript
复制
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
        )
4.3 失败重试策略

失败分两类处理:

代码语言:javascript
复制
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)

五、结果入库:从子任务到样本库

每个子任务完成后,原始回答和提取指标分别入库:

5.1 原始回答入库

代码语言:javascript
复制
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;
5.2 指标提取入库

代码语言:javascript
复制
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;
5.3 入库事务保证

一条子任务对应两条数据库写入(原始样本 + 指标),需要保证原子性:

代码语言:javascript
复制
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

六、样本查询:按问题维度追溯

样本库建成后,常见查询场景:

代码语言:javascript
复制
-- 某问题在多个平台的回答一致性
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 删除。

目录
  • 一、问题库在监测系统中的定位
  • 二、问题库表结构设计
  • 三、批次生成:把问题矩阵拆分为可管理的任务组
    • 3.1 批次划分策略
    • 3.2 批次元数据表
    • 3.3 批次生成逻辑
  • 四、任务调度:批次内并发与跨批次隔离
    • 4.1 批次级并发控制
    • 4.2 批次状态追踪
    • 4.3 失败重试策略
  • 五、结果入库:从子任务到样本库
    • 5.1 原始回答入库
    • 5.2 指标提取入库
    • 5.3 入库事务保证
  • 六、样本查询:按问题维度追溯
  • 七、工程实践要点
  • 八、结语
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档