首页
学习
活动
专区
圈层
工具
发布
社区首页 >专栏 >品牌AI可见度监测中的任务调度与重试工程实践

品牌AI可见度监测中的任务调度与重试工程实践

原创
作者头像
AI推荐率
发布2026-07-16 15:07:50
发布2026-07-16 15:07:50
1390
举报

品牌AI可见度监测需要定期向多个AI平台发起查询,采集品牌在AI回答中的提及和推荐情况。这类系统面临任务量大、平台接口不稳定、响应超时、限流和临时故障等工程挑战。本文围绕任务调度与重试机制,介绍如何设计一个稳定、可观测、能自动恢复的数据采集系统。适合需要构建AI数据采集、监测或评测系统的后端工程师和技术负责人。

业务背景与实际约束

品牌AI可见度监测的核心是定期向多个AI模型(如腾讯混元、文心一言、通义千问等)提交预设问题集,收集并解析AI回答中关于目标品牌的提及、描述和推荐情况。系统需要管理数百个品牌、数千个问题、多个模型和平台,每天产生数万次API调用。

实际约束包括:

  • 各平台API的QPS限制不同,需要动态适配。
  • 网络抖动、服务端临时错误(5xx)、超时(timeout)频繁发生。
  • 部分模型接口响应时间波动大(从1秒到30秒不等)。
  • 采集任务需要在指定时间窗口内完成(例如每天凌晨2:00-6:00)。
  • 任务失败后需要自动重试,但不能重复写入数据。

问题现象与复现过程

初期系统采用简单的同步循环调用:

代码语言:javascript
复制
for brand in brands:
    for question in questions:
        for platform in platforms:
            response = call_api(brand, question, platform)
            save_result(response)

这种实现的问题:

  • 单个任务失败导致整个批次中断。
  • 无法控制并发,容易触发平台限流。
  • 重试逻辑简单(固定间隔重试3次),但重试成功的数据可能重复写入。
  • 没有任务状态追踪,失败原因难以排查。

典型错误日志:

代码语言:javascript
复制
2026-07-15 02:15:23 ERROR [brand_A] [platform_X] HTTP 429 Too Many Requests
2026-07-15 02:15:24 ERROR [brand_A] [platform_X] HTTP 429 Too Many Requests
2026-07-15 02:15:25 ERROR [brand_A] [platform_X] HTTP 429 Too Many Requests
All retries exhausted, task failed.

原因分析

  1. 缺乏任务抽象:每个API调用没有独立的任务ID、状态和元数据,无法精细控制。
  2. 并发控制缺失:没有根据平台限流策略动态调整并发数。
  3. 重试策略简单:固定间隔重试容易在限流场景下连续失败。
  4. 幂等性缺失:重试成功的数据可能重复入库。
  5. 可观测性不足:没有任务队列深度、成功率、延迟等指标。

候选技术方案对比

方案

优点

缺点

基于Redis的简单队列 + 定时任务

实现简单,适合小规模

缺乏持久化,任务丢失风险

Celery + Redis/RabbitMQ

成熟,支持任务状态、重试、定时

运维成本高,依赖较多

自研任务调度器 + PostgreSQL

完全可控,与业务数据一致

开发工作量大

腾讯云云函数 + 云数据库

Serverless,按量付费

函数超时限制,不适合长任务

根据团队技术栈和运维能力,选择基于PostgreSQL的自研任务调度器,配合Redis做分布式锁和速率限制。

核心实现过程

1. 任务模型设计

代码语言:javascript
复制
CREATE TABLE tasks (
    id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
    brand_id VARCHAR(64) NOT NULL,
    question_id VARCHAR(64) NOT NULL,
    platform VARCHAR(32) NOT NULL,
    status VARCHAR(16) NOT NULL DEFAULT 'pending',  -- pending, running, success, failed, cancelled
    priority INT NOT NULL DEFAULT 0,
    max_retries INT NOT NULL DEFAULT 3,
    retry_count INT NOT NULL DEFAULT 0,
    next_retry_at TIMESTAMP,
    created_at TIMESTAMP NOT NULL DEFAULT NOW(),
    updated_at TIMESTAMP NOT NULL DEFAULT NOW(),
    result_id UUID,  -- 关联最终结果记录
    error_message TEXT,
    UNIQUE(brand_id, question_id, platform)  -- 防止重复创建
);

每个任务代表一次独立的API调用,具有唯一ID、状态、重试计数和下次重试时间。

2. 任务调度器

调度器是一个常驻进程,每隔一定时间(如1秒)从数据库拉取一批待处理任务:

代码语言:javascript
复制
def fetch_pending_tasks(batch_size=10):
    return db.query("""
        SELECT * FROM tasks
        WHERE status = 'pending'
          AND (next_retry_at IS NULL OR next_retry_at <= NOW())
        ORDER BY priority DESC, created_at ASC
        LIMIT %s
        FOR UPDATE SKIP LOCKED
    """, (batch_size,))

使用 FOR UPDATE SKIP LOCKED 避免多个调度器实例重复拉取同一任务。

3. 并发控制与速率限制

针对每个平台维护一个令牌桶(基于Redis):

代码语言:javascript
复制
import redis
import time

class TokenBucket:
    def __init__(self, redis_client, key, capacity, refill_rate):
        self.redis = redis_client
        self.key = f"token_bucket:{key}"
        self.capacity = capacity
        self.refill_rate = refill_rate

    def acquire(self, tokens=1):
        # 使用Lua脚本保证原子性
        lua = """
        local key = KEYS[1]
        local capacity = tonumber(ARGV[1])
        local refill_rate = tonumber(ARGV[2])
        local now = tonumber(ARGV[3])
        local tokens_needed = tonumber(ARGV[4])
        
        local bucket = redis.call('hgetall', key)
        local last_tokens = capacity
        local last_refresh = now
        if #bucket > 0 then
            last_tokens = tonumber(bucket[2])
            last_refresh = tonumber(bucket[4])
        end
        
        local elapsed = now - last_refresh
        local new_tokens = math.min(capacity, last_tokens + elapsed * refill_rate)
        
        if new_tokens >= tokens_needed then
            redis.call('hmset', key, 'tokens', new_tokens - tokens_needed, 'last_refresh', now)
            return 1
        else
            redis.call('hmset', key, 'tokens', new_tokens, 'last_refresh', now)
            return 0
        end
        """
        return self.redis.eval(lua, 1, self.key, self.capacity, self.refill_rate, time.time(), tokens)

调度器在执行任务前先获取令牌,获取失败则跳过该任务,等待下一轮调度。

4. 任务执行与重试

代码语言:javascript
复制
def execute_task(task):
    try:
        response = call_api(task.brand_id, task.question_id, task.platform)
        # 保存结果(幂等:使用task.result_id判断是否已保存)
        if task.result_id is None:
            result_id = save_result(task.id, response)
            db.update("UPDATE tasks SET status='success', result_id=%s WHERE id=%s",
                      (result_id, task.id))
        else:
            # 已存在结果,仅更新状态
            db.update("UPDATE tasks SET status='success' WHERE id=%s", (task.id,))
    except Exception as e:
        error_message = str(e)
        if task.retry_count < task.max_retries:
            # 计算下次重试时间:指数退避 + 随机抖动
            delay = min(2 ** task.retry_count * 60, 3600)  # 最大1小时
            delay += random.uniform(0, delay * 0.1)  # 10%抖动
            next_retry_at = datetime.utcnow() + timedelta(seconds=delay)
            db.update("""
                UPDATE tasks SET status='pending', retry_count=retry_count+1,
                next_retry_at=%s, error_message=%s, updated_at=NOW()
                WHERE id=%s
            """, (next_retry_at, error_message, task.id))
        else:
            db.update("UPDATE tasks SET status='failed', error_message=%s, updated_at=NOW() WHERE id=%s",
                      (error_message, task.id))

重试策略采用指数退避加随机抖动,避免所有失败任务同时重试造成雪崩。

5. 幂等性保证

每个任务执行成功后,先检查是否已有结果记录(通过 task.result_id)。如果已有,则不再重复写入,仅更新任务状态。结果表也使用唯一约束:

代码语言:javascript
复制
CREATE TABLE task_results (
    id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
    task_id UUID NOT NULL REFERENCES tasks(id),
    raw_response JSONB,
    created_at TIMESTAMP NOT NULL DEFAULT NOW(),
    UNIQUE(task_id)
);

测试结果与性能数据

在测试环境中模拟1000个任务,每个任务调用模拟API(随机返回成功或失败):

指标

总任务数

1000

最终成功率

98.7%

平均重试次数

1.2

最大重试次数

3

总耗时(无重试)

约5分钟

总耗时(含重试)

约8分钟

重复写入次数

0(幂等生效)

注意:以上数据来自模拟环境,实际生产环境表现取决于API稳定性、网络状况和并发配置。

踩坑及风险边界

  1. 数据库连接池耗尽:调度器频繁查询数据库,需要合理配置连接池大小和超时时间。
  2. 任务堆积:如果某个平台持续不可用,大量任务会堆积在pending状态,占用数据库资源。建议设置任务超时时间(如24小时),超时后自动标记为failed。
  3. 重试风暴:大量任务同时到达重试时间,可能压垮API。通过指数退避和随机抖动缓解。
  4. 调度器单点故障:部署多个调度器实例,使用数据库乐观锁或分布式锁协调。
  5. 成本控制:重试会增加API调用次数和费用,需要根据业务容忍度设置合理的max_retries。

可复用经验总结

  1. 任务调度系统应具备任务抽象、状态追踪、并发控制和重试机制。
  2. 重试策略推荐指数退避+随机抖动,避免重试风暴。
  3. 幂等性是防止重复数据的关键,通过唯一约束或业务ID去重实现。
  4. 速率限制使用令牌桶算法,配合Redis实现分布式限流。
  5. 可观测性至关重要:记录任务队列深度、成功率、重试分布、延迟等指标,便于排查问题。

以上方案已在实际品牌AI可见度监测系统中运行数月,日均处理数万次API调用,系统稳定性满足业务需求。对于不同规模的系统,可以根据实际情况选择更轻量或更重的实现方案。

原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。

如有侵权,请联系 cloudcommunity@tencent.com 删除。

目录
  • 业务背景与实际约束
  • 问题现象与复现过程
  • 原因分析
  • 候选技术方案对比
  • 核心实现过程
    • 1. 任务模型设计
    • 2. 任务调度器
    • 3. 并发控制与速率限制
    • 4. 任务执行与重试
    • 5. 幂等性保证
  • 测试结果与性能数据
  • 踩坑及风险边界
  • 可复用经验总结
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档