Funboost 高级重试与容错

SkillDev tools

当需要配置 funboost 的高级重试策略时使用。触发场景:指数退避重试、死信队列、熔断器、任务去重过滤、自定义错误处理。关键词:retry, max_retry_times, dead letter, DLX, exponential backoff, circuit breaker, is_using_advanced_retry, is_push_to_dlx_queue_when_retry_max_times, do_task_filtering, 重试, 死信, 熔断。

Available today. Use it from your connected AI after setup.

Connect ahel once, and every AI you use reads what you have installed.

Then ask your AI: use the Funboost 高级重试与容错 skill

What this skill tells your AI

The instructions your AI receives, as published by ydf0509/funboost in .agents/skills/funboost-advanced-retry/SKILL.md and read by ahel’s review.

概述

Funboost 提供多层错误恢复能力:自动重试、指数退避、死信队列、熔断器和任务去重。

核心原则: 通过 BoosterParams 声明式配置重试行为——不要手动编写 try/except 重试循环。

适用场景

  • 任务可能因瞬态故障失败(网络、API 限流)
  • 需要重试间递增等待时间(指数退避)
  • 希望最终失败的任务路由到死信队列
  • 需要熔断器保护下游服务
  • 需要防止同一任务重复执行

速查表

功能BoosterParams 字段默认值
最大重试次数max_retry_times3(最多重试 N 次 + 首次执行 = 共 N+1 次机会)
指数退避is_using_advanced_retryFalse
死信队列is_push_to_dlx_queue_when_retry_max_timesFalse
任务去重do_task_filteringFalse
函数超时function_timeoutNone

基础重试

from funboost import boost, BoosterParams

@boost(BoosterParams(
    queue_name="fragile_task",
    max_retry_times=5,          # 异常时最多重试 5 次
    function_timeout=30,        # 单次执行超过 30 秒则终止
))
def call_external_api(url: str):
    import requests
    resp = requests.get(url, timeout=10)
    resp.raise_for_status()
    return resp.json()

max_retry_times=N 表示最多重试 N 次,加上首次执行,共 N+1 次执行机会。

is_using_advanced_retry=False(默认)时,失败后会在同一线程内立即重试,无等待间隔。只有启用高级重试后才走指数退避。

指数退避重试

@boost(BoosterParams(
    queue_name="backoff_task",
    max_retry_times=5,
    is_using_advanced_retry=True,  # 启用指数退避
))
def rate_limited_api(endpoint: str):
    """指数退避:间隔 = min(base * 2^n, max_interval),默认 1s,2s,4s,8s...封顶60s"""
    import requests
    resp = requests.get(endpoint)
    if resp.status_code == 429:
        raise Exception("被限流了")
    return resp.json()

advanced_retry_config 参数

通过 BoosterParams.advanced_retry_config 字典自定义退避行为:

参数类型默认值说明
retry_modestr"sleep""sleep" = 当前线程 sleep 等待后重试;"requeue" = 发回队列(附带 countdown),释放线程
retry_base_intervalfloat1.0退避基础间隔(秒)
retry_multiplierfloat2.0退避倍数(每次重试间隔乘以此值)
retry_max_intervalfloat60.0退避最大间隔(秒),封顶值
retry_jitterboolFalse是否加随机抖动,防止多消费者同时重试(惊群)
@boost(BoosterParams(
    queue_name="custom_backoff_task",
    max_retry_times=8,
    is_using_advanced_retry=True,
    advanced_retry_config={
        "retry_mode": "requeue",         # 重试时发回队列,释放线程资源
        "retry_base_interval": 2.0,      # 起始等 2 秒
        "retry_multiplier": 3.0,         # 每次 x3:2s, 6s, 18s, 54s, 60s(封顶)...
        "retry_max_interval": 60.0,      # 最大 60 秒
        "retry_jitter": True,            # 加随机抖动
    },
))
def custom_backoff_task(data: dict):
    process(data)

两种 retry_mode 的区别:

  • "sleep" — 简单场景,线程被占用直到重试完成;适合并发数充足时
  • "requeue" — 消息重新入队(带 countdown 延迟),当前线程立即释放去处理其他消息;适合高并发/长退避场景

死信队列 (DLX)

开启 is_push_to_dlx_queue_when_retry_max_times=True 后,重试耗尽的任务推送到死信队列:

@boost(BoosterParams(
    queue_name="important_task",
    max_retry_times=3,
    is_push_to_dlx_queue_when_retry_max_times=True,  # -> important_task_dlx
))
def process_payment(order_id: str, amount: float):
    """3 次重试全部失败后,消息进入 'important_task_dlx' 队列"""
    charge(order_id, amount)

主动控制重试/入队的异常类

funboost 提供 3 个内置异常类,可在消费函数中主动 raise 来控制任务流向:

异常类行为是否计入 max_retry_times
ExceptionForRequeue立即将消息重新放回当前队列(不等待、不计重试次数)
ExceptionForPushToDlxqueue立即将消息推入死信队列 {queue_name}_dlx
ExceptionForRetry语义标记(任意异常本就触发重试,此类仅增强语义清晰度)
from funboost import ExceptionForRequeue, ExceptionForPushToDlxqueue

@boost(BoosterParams(queue_name="controlled_retry", broker_kind=BrokerEnum.REDIS_ACK_ABLE))
def process_order(order_id):
    status = check_order_status(order_id)
    if status == "not_ready":
        raise ExceptionForRequeue()  # 还没准备好,重回队列稍后再试
    if status == "invalid":
        raise ExceptionForPushToDlxqueue()  # 无效订单,直接进死信
    # 正常处理...

注意: ExceptionForRequeue 不受 max_retry_times 限制,消息会无限重回队列直到不再抛此异常。使用时需确保有退出条件,避免无限循环。

熔断器 (Mixin)

连续失败达阈值时熔断(OPEN 状态下消费端 _submit_task 阻塞等待恢复,消息暂留队列;不影响发布端 push/publish),防止级联故障:

from funboost import boost, BoosterParams
from funboost.contrib.override_publisher_consumer_cls.circuit_breaker_mixin import CircuitBreakerConsumerMixin

@boost(BoosterParams(
    queue_name="protected_task",
    consumer_override_cls=CircuitBreakerConsumerMixin,
    user_options={
        "circuit_breaker_options": {
            "failure_threshold": 5,     # 连续 5 次失败后熔断(OPEN)
            "recovery_timeout": 60,     # 60 秒后进入半开(HALF_OPEN)试探
        },
    },
))
def call_fragile_service(data: dict):
    return external_service.process(data)

任务去重(过滤)

消费端过滤:相同入参的任务完成一次消费周期后,再次发布会被跳过(需 Redis):

@boost(BoosterParams(
    queue_name="dedup_task",
    do_task_filtering=True,  # 需要 Redis 配置
))
def send_notification(user_id: int, message: str):
    """相同 (user_id, message) 组合完成消费后不会重复处理"""
    notify(user_id, message)

获取重试上下文

from funboost import fct

@boost(BoosterParams(queue_name="ctx_retry", max_retry_times=3))
def task_with_context(url: str):
    run_times = fct.function_result_status.run_times
    if run_times > 1:
        print(f"第 {run_times} 次执行(第 {run_times - 1} 次重试)")
    fetch(url)

组合多种策略

@boost(BoosterParams(
    queue_name="resilient_pipeline",
    max_retry_times=5,
    is_using_advanced_retry=True,
    is_push_to_dlx_queue_when_retry_max_times=True,
    function_timeout=60,
    do_task_filtering=True,
))
def resilient_task(job_id: str, payload: dict):
    """
    - 指数退避重试 5 次
    - 每次超时 60 秒
    - 相同入参完成消费后自动跳过
    - 重试耗尽后进入死信队列
    """
    process(job_id, payload)

常见错误

错误修正
max_retries=5正确字段:max_retry_times=5
timeout=30正确字段:function_timeout=30
手写重试循环让 funboost 通过 BoosterParams 处理重试
do_task_filtering=True 但没配 Redis任务过滤依赖 Redis
搞混 DLX 队列名自动命名为 {原队列名}_dlx
try/except 吞掉所有异常让异常正常抛出,funboost 才能执行重试

相关 Skill

  • developing-funboost-mixin — Consumer/Publisher Mixin 扩展
  • funboost-troubleshooting — 排错与 FAQ

Signals

GitHub stars
891
Forks
166
Last commit
Aug 2026
Advanced
Catalog kind
skill
Gateway key
funboost-advanced-retry
Source
github.com/ydf0509/funboost