4. 使用框架的各种代码示例

框架极其简单并且自由,只有一个 @boost 装饰器的参数学习。本章所有例子都是调整 @boost 的参数而已。

消息中间件的 ip/端口/密码等配置在首次运行代码时,框架会在项目根目录自动生成 funboost_config.py,按需修改即可。

所有例子的发布和消费都不必写在同一个 py 文件(使用中间件解耦)。


4.索引 快速导航

索引1. 基础用法

章节

内容

4.0

@boost 装饰器入参格式(BoosterParams)

4.1

装饰器方式调度函数(最简示例)

4.2c

动态按队列名生成 booster(BoostersManager.build_booster)

4.2d

BoostersManager 管理(一键启动所有消费)

4.25

push 和 publish 发布消息的区别

索引2. 消费控制

章节

内容

4.3a

多步骤消费(step1 → step2 链式调度)

4.3b

多个消费者共享同一线程池

4.3c

清空队列 / 获取队列消息数量

4.18

暂停 / 恢复消费

4.33

is_auto_start_consuming_message 自动启动

索引3. 并发与性能

章节

内容

4.5

多进程 + 多线程/协程叠加并发

4.7

QPS 精准控频(核心功能)

4.12

asyncio 方式运行协程

4.23

funboost 如何取代手写线程池

4.26

性能调优演示

索引4. 定时与延时

章节

内容

4.4

定时运行(ApsJobAdder)

4.9

延时运行任务

索引5. 高级功能

章节

内容

4.6

RPC 模式(远程调用获取结果)

4.24

设置消费函数重试次数

4.29

任务优先级队列

4.30

远程杀死/取消任务

4.31

fct 上下文获取当前消息和任务状态

4.35

函数入参过滤(去重)

4.36

自定义类型入参(pickle 序列化)

4.38

MemoryFunboostPool / FunboostPool

索引6. 集成与扩展

章节

内容

4.10

Flask / FastAPI / Django 集成

4.11

保存消费状态到 MongoDB

4.13

跨项目发布任务

4.19

自定义消费状态/结果钩子

4.20

broker_exclusive_config 差异化配置

4.21

register_custom_broker 自定义中间件

4.21b

consumer_override_cls / publisher_override_cls

4.28

Celery 框架整体作为 funboost 的 Broker

索引7. 其他

章节

内容

4.14

获取消费进程信息

4.16

文件日志所在位置

4.17

判断函数运行完所有任务

4.32

实例方法/类方法作为消费函数

4.34

PyInstaller 打包说明

4.37

funboost 启动消费的方式大全

4.100

使用控制变量法验证框架功能


4.0 @boost 装饰器入参格式

4.0.1 新推荐写法:BoosterParams

from funboost import boost, BrokerEnum, BoosterParams

@boost(BoosterParams(queue_name='queue_test_f01', qps=0.2, broker_kind=BrokerEnum.REDIS_ACK_ABLE))
def add(a, b):
    print(a + b)

采用 pydantic Model 入参,IDE 可完全自动补全。需在 PyCharm 安装 pydantic 插件(file → settings → Plugins → 搜索 pydantic)。

4.0.2 自定义子类继承 BoosterParams

每次少传重复参数:

class BoosterParamsMy(BoosterParams):
    broker_kind: str = BrokerEnum.RABBITMQ
    max_retry_times: int = 4
    log_level: int = logging.DEBUG

@boost(BoosterParamsMy(queue_name='task_queue_name1', qps=3))
def task_fun(x, y):
    print(f'{x} + {y} = {x + y}')
    time.sleep(3)

4.0.3 老的直接传参方式(仍可用)

@boost('queue_test_f01', qps=0.2, broker_kind=BrokerEnum.REDIS_ACK_ABLE)
def add(a, b):
    print(a + b)

4.1 装饰器方式调度函数

最简单的使用方式:一个 @boost 装饰器即可将普通函数变为分布式可调度函数,支持 push 发布消息和 consume 启动消费:

from funboost import boost, BrokerEnum, BoosterParams

@boost(BoosterParams(queue_name='queue_test_f01', qps=0.2, broker_kind=BrokerEnum.REDIS_ACK_ABLE))
def add(a, b):
    print(a + b)

if __name__ == '__main__':
    for i in range(10, 20):
        add.publish(dict(a=i, b=i * 2))  # publish 发布字典
        add.push(i, b=i * 2)             # push 发布参数
    add.consume()
    # add.multi_process_consume(4)  # 4进程叠加并发

4.2c 动态生成 booster

在函数内部需要按队列名动态生成 booster 时,使用 BoostersManager.build_booster:

from funboost import Booster, BoostersManager, BoosterParams

def add(a, b):
    print(a + b)

def my_push(queue_name, a, b):
    booster = BoostersManager.build_booster(
        BoosterParams(queue_name=queue_name, qps=0.2, consuming_function=add)
    )  # type: Booster
    booster.push(a, b)

if __name__ == '__main__':
    for i in range(1000000):
        queue_namex = f'queue_{i % 10}'
        my_push(queue_namex, i, i * 2)

    for j in range(10):
        booster = BoostersManager.build_booster(
            BoosterParams(queue_name=f'queue_{j}', qps=0.2, consuming_function=add)
        )
        booster.consume()

build_booster 内部有缓存:同一 queue_name 不会重复创建连接。不要在循环中直接 boost(BoosterParams(...))(add) 创建 booster,这会创建大量重复连接。


4.2d BoostersManager 管理

所有 @boost 或 build_booster 创建的 booster 都会自动登记到 BoostersManager。

4.2d.1 一次性启动所有队列消费

from funboost import BoostersManager

# 方式1:消费组(手动指定哪些函数属于一组)
BoostersManager.consume_group('group1')

# 方式2:自动发现并消费所有
from funboost import BoosterDiscovery
BoosterDiscovery(project_root_path='.', booster_dirs=['./tasks']).auto_discovery()
BoostersManager.consume_all()

4.2d.3 使用 BoostersManager 通过 consume_group 启动一组消费函数

给 @boost 的 BoosterParams 设置 booster_group='group1',然后 BoostersManager.consume_group('group1') 一次性启动该组所有消费函数:

class MyGroup1BoosterParams(BoosterParams):
    booster_group: str = "my_group1"

@boost(MyGroup1BoosterParams(queue_name="queue_g1"))
def f1(x): print(f"f1 {x}")

@boost(MyGroup1BoosterParams(queue_name="queue_g2"))
def f2(x): print(f"f2 {x}")

if __name__ == "__main__":
    BoostersManager.consume_group("my_group1")

注意:如果 @boost 函数分布在多个模块中,需先 import 或用 BoosterDiscovery().auto_discovery() 自动导入。


4.2e funboost 支持实例方法、类方法、静态方法、普通函数 4 种类型

2024 年 6 月新增支持实例方法、类方法作为消费函数。详见 4.32 章节。


4.3a 多步骤消费函数

step1 可以给 step2 发布任务,也可以给自身发布任务(递归调度):

from funboost import boost, BrokerEnum, BoosterParams

@boost(BoosterParams(queue_name='queue_test_step1', qps=0.5, broker_kind=BrokerEnum.LOCAL_PYTHON_QUEUE))
def step1(x):
    print(f'x 的值是 {x}')
    if x == 0:
        for i in range(1, 300):
            step1.push(x + i)
    for j in range(10):
        step2.push(x * 100 + j)

@boost(BoosterParams(queue_name='queue_test_step2', qps=3, broker_kind=BrokerEnum.LOCAL_PYTHON_QUEUE))
def step2(y):
    print(f'y 的值是 {y}')

if __name__ == '__main__':
    step1.push(0)
    step1.consume()
    step2.consume()

4.3.b 共享线程池

多个消费者使用同一个并发池,减少资源浪费:

from funboost import boost, BoosterParams
from funboost.concurrent_pool.flexible_thread_pool import FlexibleThreadPool

pool = FlexibleThreadPool(300)

@boost(BoosterParams(queue_name='test_f1_queue', specify_concurrent_pool=pool, qps=3))
def f1(x):
    print(f'x : {x}')

@boost(BoosterParams(queue_name='test_f2_queue', specify_concurrent_pool=pool, qps=2))
def f2(y):
    print(f'y : {y}')

if __name__ == '__main__':
    for i in range(1000):
        f1.push(i)
        f2.push(i)
    f1.consume()
    f2.consume()

4.3c 清空队列/获取消息数量

booster 对象提供队列管理的便捷方法:

f.clear()                    # 清空队列中所有未消费的消息
count = f.get_message_count()  # 获取队列中的消息数量

4.4 定时运行

4.4.1 核心原理

funboost 的定时任务是定时发布消息到消息队列,而非直接执行函数。add_push_job 本质是每隔 N 秒自动运行 fun.push()。

4.4.2 代码演示

from funboost import boost, BrokerEnum, BoosterParams, ApsJobAdder

@boost(BoosterParams(queue_name='sum_queue5', broker_kind=BrokerEnum.REDIS))
def sum_two_numbers(x, y):
    print(f'The sum of {x} and {y} is {x + y}')

@boost(BoosterParams(queue_name='data_queue5', broker_kind=BrokerEnum.REDIS))
def show_msg(data):
    print(f'data: {data}')

if __name__ == '__main__':
    # 每隔5秒发布一次任务
    ApsJobAdder(sum_two_numbers, job_store_kind='redis').add_push_job(
        args=(10, 20), trigger='interval', seconds=5, id='sum_job'
    )
    # 每天凌晨2点发布
    ApsJobAdder(show_msg, job_store_kind='redis').add_push_job(
        kwargs={'data': 'hello'}, trigger='cron', hour=2, id='show_job'
    )
    sum_two_numbers.consume()
    show_msg.consume()

定时语法和入参与 funboost 无关,请学习 apscheduler 3.x 官方文档。

4.4.3 ApsJobAdder 的优势

  1. 用户无需手写 push_msg 包装函数

  2. Redis 作为 job_store 时,每个函数使用独立的 jobs_key

  3. 使用 Redis 分布式锁,多机多进程不重复执行


4.4.4 新增支持 aps_obj.add_job 添加定时任务(2025-08)

除了 add_push_job,现在也支持直接使用 apscheduler 原生的 add_job 方式添加定时任务:

aps_obj = ApsJobAdder(my_task, job_store_kind='redis')
aps_obj.add_job(my_push_func, trigger='interval', seconds=5, id='my_job')

区别在于 add_job 需要用户自己写 push 包装函数,而 add_push_job 会自动帮你生成。


4.5 多进程并发

使用 multi_process_consume 启动多进程叠加并发(多进程 × 多线程/协程),充分利用多核 CPU:

from funboost import boost, BoosterParams, BrokerEnum

@boost(BoosterParams(queue_name='test_multi_process', broker_kind=BrokerEnum.REDIS_ACK_ABLE, qps=100))
def task(x):
    print(x)

if __name__ == '__main__':
    for i in range(10000):
        task.push(i)
    task.multi_process_consume(4)  # 4进程 × 多线程叠加并发
    # 等价于 task.mp_consume(4)

Linux 上建议在脚本顶部加 multiprocessing.set_start_method('spawn', force=True),避免 fork 导致的内存污染问题。


4.6 RPC 模式

客户端调用远程函数并获取结果:

from funboost import boost, BoosterParams, BrokerEnum

@boost(BoosterParams(
    queue_name='rpc_queue',
    broker_kind=BrokerEnum.REDIS_ACK_ABLE,
    is_using_rpc_mode=True,
))
def add(x, y):
    return x + y

if __name__ == '__main__':
    add.consume()

    # 异步获取结果
    async_result = add.push(1, 2)
    print(async_result.result)  # 阻塞等待结果:3

    # 同步 get_future 方式
    future = add.publisher.get_future(x=10, y=20)
    result_status = future.result(timeout=5)
    print(result_status.result)  # 30

4.6.5 asyncio 语法生态下 rpc 获取执行结果

在 asyncio 编程中不能用同步的 async_result.result(会阻塞 event loop),需使用 AioAsyncResult:

import asyncio
from funboost import AioAsyncResult

async def test_get_result(i):
    async_result = add.push(i, i * 2)
    aio_async_result = AioAsyncResult(task_id=async_result.task_id)
    print(await aio_async_result.result)
    print(await aio_async_result.status_and_result)

完整 asyncio 编程示例见 4b.3 章节。


4.7 QPS 控频

4.7.1 核心特性

  • 基于精密计时算法,非 Redis incr 计数

  • 对耗时恒定函数精确度 99.9%,耗时随机波动函数精确度 96%+

  • 支持分布式全局控频:多机器自动均分 QPS 配额

  • 智能自适应:函数耗时变大时自动扩线程,耗时变小时自动缩线程

4.7.2 演示:自适应扩缩容

import time
import threading
from funboost import boost, BrokerEnum, ConcurrentModeEnum, BoosterParams

t_start = time.time()

@boost(BoosterParams(
    queue_name='queue_test2_qps', qps=2,
    broker_kind=BrokerEnum.PERSISTQUEUE,
    concurrent_mode=ConcurrentModeEnum.THREADING,
    concurrent_num=600
))
def f2(a, b):
    result = a + b
    sleep_time = 0.01
    if time.time() - t_start > 60:
        sleep_time = 7
    if time.time() - t_start > 120:
        sleep_time = 30
    if time.time() - t_start > 240:
        sleep_time = 0.8
    print(f'{time.strftime("%H:%M:%S")} 线程数:{threading.active_count()}, {a}+{b}={result}, sleep {sleep_time}s')
    if sleep_time is not None:
        time.sleep(sleep_time)
    return result

if __name__ == '__main__':
    f2.clear()
    for i in range(1000):
        f2.push(i, i * 2)
    f2.consume()

4.7.3 分布式 QPS 控频

@boost(BoosterParams(
    queue_name='distributed_qps_task',
    qps=50000,
    is_using_distributed_frequency_control=True,
    broker_kind=BrokerEnum.REDIS_ACK_ABLE,
))
def heavy_task(x):
    pass

部署 100 个进程(20 台 8 核机器),框架自动统计活跃消费者数量并分配 QPS。增减机器无需改代码。


4.8 QPS 为什么强大?常规并发方式无法完成的需求

常规线程池 ThreadPoolExecutor 只能控制线程数量,无法控制运行速率。QPS 是 funboost 对线程池的降维打击:

  • 设置 qps=5,无论函数耗时 0.01 秒还是 30 秒,框架都能精确保证每秒恰好运行 5 次

  • 线程数量会根据函数耗时自动调节——耗时长就多开线程,耗时短就自动缩减

  • 分布式场景下多机器自动均分配额,全局精准联控

这是 Celery 的 rate_limit 完全无法做到的(Celery 在 rate_limit > 20/s 时精确度只有 60%)。


4.9 延时任务

通过 push 时传递 countdown 参数实现延时执行(单位秒),适用于订单超时检查、延迟通知等场景:

from funboost import boost, BoosterParams, BrokerEnum

@boost(BoosterParams(queue_name='delay_queue', broker_kind=BrokerEnum.REDIS_ACK_ABLE))
def delayed_task(order_id):
    print(f'检查订单 {order_id} 是否已支付')

if __name__ == '__main__':
    # 15分钟后执行
    delayed_task.push(order_id=12345, countdown=900)

专业延时方案:使用 broker_kind=BrokerEnum.REDIS_ZSET_DELAY 可获得原生延时队列支持,延时精度更高、性能更好(无需 APScheduler 二次投递)。


4.10 Web 框架集成

4.10.1 FastAPI 集成

from fastapi import FastAPI
from funboost import boost, BoosterParams, BrokerEnum

app = FastAPI()

@boost(BoosterParams(queue_name='web_task', broker_kind=BrokerEnum.REDIS_ACK_ABLE))
def process_order(order_id: int):
    print(f'处理订单 {order_id}')

@app.on_event("startup")
def startup():
    process_order.consume()

@app.post("/order/{order_id}")
def create_order(order_id: int):
    process_order.push(order_id)
    return {"status": "submitted"}

Flask 和 Django 的集成方式类似:在应用启动时调用 consume(),在视图函数中调用 push()。


4.11 消费状态持久化

配置 FunctionResultStatusPersistanceConfig 后,每次函数执行的状态(成功/失败/耗时/结果/异常)会自动持久化:

from funboost import boost, BoosterParams, BrokerEnum, FunctionResultStatusPersistanceConfig

@boost(BoosterParams(
    queue_name='save_status_task',
    broker_kind=BrokerEnum.REDIS_ACK_ABLE,
    function_result_status_persistance_conf=FunctionResultStatusPersistanceConfig(
        is_save_status=True,
        is_save_result=True,
    ),
))
def my_task(x):
    return x * 2

启用后,每次函数执行的状态(成功/失败/耗时/结果/异常)会保存到 MongoDB,可在 funweb 页面查看。

保存到 MySQL/SQLite/PostgreSQL:通过 user_custom_record_process_info_func 钩子实现,funboost 已内置 save_result_status_to_sqlalchemy:

from funboost.contrib.save_function_result_status.save_result_status_to_sqldb import save_result_status_to_sqlalchemy
@boost(BoosterParams(queue_name='xx', user_custom_record_process_info_func=save_result_status_to_sqlalchemy))

需先建表(见 funboost/contrib/save_function_result_status/save_result_status_to_sqldb.py 中的建表 SQL)并配置 BrokerConnConfig.SQLACHEMY_ENGINE_URL。


4.12 asyncio 并发

使用 ConcurrentModeEnum.ASYNC 启用协程并发模式,支持 async def 消费函数:

from funboost import boost, BoosterParams, BrokerEnum, ConcurrentModeEnum

@boost(BoosterParams(
    queue_name='async_queue',
    broker_kind=BrokerEnum.REDIS,
    concurrent_mode=ConcurrentModeEnum.ASYNC,
    concurrent_num=100,
))
async def async_task(url):
    import aiohttp
    async with aiohttp.request('GET', url) as resp:
        print((await resp.text())[:50])

注意:ASYNC 模式下函数内不能有阻塞的同步代码(如 requests.get、time.sleep)。部分异步库(如 aiohttp、aiomysql)的连接池绑定了创建时的 loop,在 funboost 子线程的 loop 中使用会报 attached to a different loop,此时需传递 specify_async_loop;而 httpx、sqlalchemy 等库不存在此问题。详见文档 6.26 章节。

THREADING 模式也能运行 async def:funboost 的 FlexibleThreadPool 能自动检测并运行 async def 函数(每个线程临时创建 loop),用户无需手写 asyncio.run() 包装。因此不必为了运行 async 函数就切换到 ASYNC 模式。


4.13 跨项目发布任务

不定义 @boost 消费函数,直接发送消息到指定队列:

from funboost import BoostersManager, BoosterParams, BrokerEnum

booster = BoostersManager.build_booster(BoosterParams(
    queue_name='other_project_queue',
    broker_kind=BrokerEnum.REDIS_ACK_ABLE,
))
booster.publish({'user_id': 123, 'action': 'send_email'})

4.13b 不使用 funboost 消费功能,funboost 作为万能发布者

funboost 可以纯粹作为各种消息队列的统一发布客户端使用,不启动消费:

from funboost import BoostersManager, BoosterParams, BrokerEnum

booster = BoostersManager.build_booster(BoosterParams(
    queue_name='any_queue', broker_kind=BrokerEnum.RABBITMQ
))
for i in range(1000):
    booster.publish({'data': i})

支持近 50 种消息队列,统一 API 发布消息。


4.14 消费进程信息

通过 BoostersManager 查看当前进程中所有已注册的 booster 及其状态:

from funboost import BoostersManager

# 查看所有已注册的 booster
for queue_name, booster in BoostersManager.pid_queue_name__booster_map.items():
    print(queue_name, booster)

4.16 文件日志

funboost 使用 nb_log 包,日志文件默认在项目根目录的 /pythonlogs/ 文件夹下。

可在 nb_log_config.py 中自定义日志路径:LOG_PATH = '/your/custom/path'


4.16.4 funboost 日志由 nb_log 提供

nb_log 是 funboost 作者开发的日志包,提供彩色控制台日志、多进程安全的文件日志、可点击跳转等特性。funboost 的所有日志均由 nb_log 驱动。


4.17 等待全部完成

等待某个队列的任务全部消费完成后再继续执行后续逻辑:

f.consume()
f.wait_for_possible_has_finish_all_tasks(minutes=3)
print("所有任务消费完成")

4.18 暂停消费

运行时动态暂停/恢复某个队列的消费(不需要重启进程):

# 暂停
my_task.pause_consume()

# 恢复
my_task.continue_consume()

4.19 自定义钩子

通过 user_custom_record_process_info_func 参数设置消费后的回调函数,可记录日志、发送告警、写入数据库等:

from funboost import boost, FunctionResultStatus, BoosterParams

def my_save_process_info(function_result_status: FunctionResultStatus):
    """function_result_status 上有丰富的消费状态信息"""
    print(function_result_status.params, function_result_status.result,
          function_result_status.time_cost, function_result_status.success,
          function_result_status.exception)

@boost(BoosterParams(
    queue_name='test_user_custom',
    user_custom_record_process_info_func=my_save_process_info,
))
def my_task(x):
    return x * 2

4.20 broker_exclusive_config

不同中间件的差异化独特配置:

@boost(BoosterParams(
    queue_name='kafka_task',
    broker_kind=BrokerEnum.KAFKA,
    broker_exclusive_config={
        'group_id': 'my_group',
        'bootstrap_servers': '192.168.0.100:9092',
    },
))
def kafka_task(msg):
    print(msg)

如何知道各 broker 支持哪些独有配置? 查看对应 Consumer 类的 BROKER_EXCLUSIVE_CONFIG_DEFAULT 属性,其中的 keys 就是该 broker 支持的配置项。例如 ConsumerKafkaConfluent.BROKER_EXCLUSIVE_CONFIG_DEFAULT 包含 group_id、auto_offset_reset、num_partitions、replication_factor。


4.21 register_custom_broker 完全自定义扩展中间件

完全自定义新的消息中间件(Consumer + Publisher):

from funboost import register_custom_broker

# 注册自定义 broker
register_custom_broker(
    broker_kind='MY_CUSTOM_BROKER',
    publisher_class=MyCustomPublisher,
    consumer_class=MyCustomConsumer,
)

注册后,即可在 @boost 中使用 broker_kind=BrokerEnum.MY_CUSTOM_BROKER。


4.21b consumer_override_cls 和 publisher_override_cls 自定义消费者/发布者

继承现有消费者/发布者,重写部分方法实现定制:

from funboost import boost, BoosterParams, AbstractConsumer, FunctionResultStatus

class MyConsumer(AbstractConsumer):
    def user_custom_record_process_info_func(self, current_function_result_status: FunctionResultStatus):
        if current_function_result_status.success:
            print(f'成功:{current_function_result_status.result}')
        else:
            print(f'失败:{current_function_result_status.exception}')

@boost(BoosterParams(queue_name='custom_task', consumer_override_cls=MyConsumer))
def my_task(x):
    return x * 2

4.21c 让 AI 帮你扩展 funboost 中间件或定制运行逻辑

将 funboost 源码和教程上传到 AI 大模型(如 Google AI Studio / DeepSeek),描述你的需求,AI 可以帮你:

  • 实现新的自定义 broker(Consumer + Publisher)

  • 继承 override 消费者实现定制逻辑

  • 生成完整的 register_custom_broker 注册代码


4.23 取代线程池

手动线程池写法:

from concurrent.futures import ThreadPoolExecutor

pool = ThreadPoolExecutor(5)
for i in range(100):
    pool.submit(f, i, i * 2)

funboost 等价写法(额外获得 QPS 控频、重试、ACK 等能力):

@boost(BoosterParams(queue_name='test1', broker_kind=BrokerEnum.MEMORY_QUEUE, concurrent_num=5))
def f(x, y):
    print(f'{x} + {y} = {x + y}')
    time.sleep(10)

if __name__ == '__main__':
    f.consume()
    for i in range(100):
        f.push(i, i * 2)

4.24 重试次数

设置 max_retry_times 后,消费函数抛出任何异常都会自动重试(包括 HTTP 200 但内容异常的场景):

@boost(BoosterParams(
    queue_name='retry_task',
    broker_kind=BrokerEnum.REDIS_ACK_ABLE,
    max_retry_times=5,  # 失败重试5次
))
def unstable_task(x):
    import random
    if random.random() < 0.5:
        raise Exception('随机失败')
    print(f'成功: {x}')

特殊异常类型:

  • raise ExceptionForRequeue():主动让消息立即重回消息队列(不计入重试次数)

  • raise ExceptionForPushToDlxqueue():主动将消息发送到死信队列

  • 设置 is_push_to_dlx_queue_when_retry_max_times=True:重试到 max_retry_times 后自动进死信队列

4.24.5 funboost 高级重试:指数退避重试

通过 is_using_advanced_retry=True 和 advanced_retry_config 启用。retry_mode 支持 sleep(当前线程阻塞等待)和 requeue(重回队列,不占线程,推荐长间隔场景)。

参数

说明

retry_mode

requeue(推荐)或 sleep

retry_base_interval

基础间隔秒数

retry_multiplier

指数退避倍数(1.0 为固定间隔)

retry_max_interval

最大间隔秒数

retry_jitter

是否随机抖动(0.5~1.5 倍)

@boost(BoosterParams(
    queue_name='advanced_retry_task',
    max_retry_times=100,
    is_using_advanced_retry=True,
    advanced_retry_config={
        'retry_mode': 'requeue',
        'retry_base_interval': 10.0,
        'retry_multiplier': 2.0,
        'retry_max_interval': 300.0,
        'retry_jitter': True,
    },
))
def api_call(url):
    pass

requeue 模式的优势:间隔越来越大时不会长时间占用工作线程,释放线程去执行其他任务。


4.25 push 和 publish 的区别

方法

入参形式

示例

push

直接传参数(位置参数 + 关键字参数)

f.push(1, y=2)

publish

传一个字典

f.publish({'x': 1, 'y': 2})

两者效果相同,只是调用风格不同。push 更符合直觉,publish 更灵活(适合动态构造参数字典)。


4.26 性能调优

  1. 提高 concurrent_num:增大并发线程数

  2. 使用 multi_process_consume:多进程叠加多线程

  3. 选择高性能 broker:REDIS > RABBITMQ > KAFKA(对于小消息)

  4. 使用 gevent/eventlet:IO 密集型任务可获得更高并发

  5. 减少日志:log_level=logging.WARNING 减少 IO 开销


4.28 Celery 作为 Broker

funboost 能将 Celery 框架整体作为自己的一个 Broker 使用:

@boost(BoosterParams(
    queue_name='celery_as_broker_task',
    broker_kind=BrokerEnum.CELERY,
))
def my_task(x):
    print(x)

这是降维打击——将 celery 作为 funboost 的一个组件。


4.29 优先级队列

使用 BrokerEnum.REDIS_PRIORITY 作为 broker,支持消息优先级调度(数字越大越优先消费):

@boost(BoosterParams(
    queue_name='priority_task',
    broker_kind=BrokerEnum.REDIS_PRIORITY,
))
def task_with_priority(data):
    print(data)

if __name__ == '__main__':
    task_with_priority.push('低优先级', priority=1)
    task_with_priority.push('高优先级', priority=10)  # 数字越大优先级越高
    task_with_priority.consume()

4.30 远程杀死任务

消费端(须设置 is_support_remote_kill_task=True):

from funboost import boost, BoosterParams

@boost(BoosterParams(queue_name='test_kill_fun_queue', is_support_remote_kill_task=True))
def long_task(x, y):
    import time
    print(f'start {x} + {y}')
    time.sleep(120)
    print(f'over {x} + {y} = {x + y}')

if __name__ == '__main__':
    long_task.consume()

发布端(发送远程杀死命令):

from funboost import RemoteTaskKiller

async_result = long_task.push(3, 4)
RemoteTaskKiller(long_task.queue_name, async_result.task_id).send_kill_remote_task_comd()

安全警告:function_timeout 和远程杀死功能通过杀死线程实现,如果函数内持有不可重入锁,可能导致死锁。推荐使用 expire_lock(pip install expire_lock 或 from funboost.utils import expire_lock),它规定了锁的最大占用时间,到期自动释放,避免永久锁死。详见 expire_lock 文档。


4.31 fct 上下文

在消费函数内部获取当前消息的元数据:

from funboost import boost, BoosterParams, fct

@boost(BoosterParams(queue_name='fct_demo'))
def my_task(x):
    print(f'task_id: {fct.task_id}')
    print(f'queue_name: {fct.queue_name}')
    print(f'function_params: {fct.function_params}')
    print(f'publish_time: {fct.function_result_status.publish_time}')
    print(f'run_times: {fct.function_result_status.run_times}')
    return x

4.32 实例方法/类方法

2024.06 月新增,支持实例方法和类方法作为消费函数。

关键注意点:

  • 实例方法 push 必须写成 类.方法.push(实例对象, 其他入参),不能写成 实例.方法.push(其他入参)

  • 类方法 push 必须写成 类.方法.push(类, 其他入参),不能省略 cls

  • 实例方法的类必须在 __init__ 中定义 obj_init_params 属性保存初始化入参,用于消费时还原对象

import copy
from funboost import boost, BoosterParams
from funboost.utils.class_utils import ClsHelper

class MyService:
    m = 1

    def __init__(self, prefix):
        self.obj_init_params: dict = ClsHelper.get_obj_init_params_for_funboost(copy.copy(locals()))
        self.prefix = prefix

    @boost(BoosterParams(queue_name='instance_method_queue'))
    def process(self, name):
        print(f'{self.prefix} {name}')

    @classmethod
    @BoosterParams(queue_name='class_method_queue')
    def class_process(cls, y):
        print(cls.m + y)

if __name__ == '__main__':
    # 实例方法:必须 类.方法.push(实例, 入参)
    MyService.process.push(MyService("Hello"), 'World')
    MyService.process.consume()

    # 类方法:必须 类.方法.push(类, 入参)
    MyService.class_process.push(MyService, 2)
    MyService.class_process.consume()

4.33 自动启动消费

设置 is_auto_start_consuming_message=True 后,booster 定义完成时自动开始消费,无需手动调用 .consume():

@boost(BoosterParams(
    queue_name='auto_start_task',
    is_auto_start_consuming_message=True,  # 定义后自动启动消费
))
def auto_task(x):
    print(x)

4.34 PyInstaller 打包

打包 funboost 项目为 exe 时,需在 .spec 文件中添加 hidden imports(具体取决于你使用的 broker_kind)。


4.35 任务过滤

基于函数参数的去重:

@boost(BoosterParams(
    queue_name='filter_task',
    do_task_filtering=True,
    task_filtering_expire_seconds=3600,  # 1小时内相同参数不重复执行
))
def crawl_page(url):
    print(f'抓取: {url}')

if __name__ == '__main__':
    crawl_page.push('https://example.com')
    crawl_page.push('https://example.com')  # 会被过滤
    crawl_page.consume()

警告:funboost 的 RPC 功能和函数入参过滤(do_task_filtering)不要同时使用——因为过滤会导致相同参数的消息不入队,RPC 调用方将永远收不到结果。


4.35c 使用 nb_cache 作为缓存装饰器

nb_cache 是 funboost 作者开发的缓存装饰器包,可以为函数结果添加缓存(内存/Redis/文件),避免重复计算:

from nb_cache import cache

@cache(expire=3600, cache_type='redis')
def expensive_compute(x):
    import time
    time.sleep(10)
    return x * x

也可以将 cache 装饰器传给 BoosterParams 的 consuming_function_decorator 参数,好处是不需要设置 should_check_publish_func_params=False:

from nb_cache import Cache

dual_cache = Cache().setup("dual://localhost:6379/0?memory_size=1000&local_ttl=30", prefix="myapp")

@boost(BoosterParams(
    queue_name='queue_test', concurrent_num=10,
    broker_kind=BrokerEnum.REDIS_ACK_ABLE,
    consuming_function_decorator=dual_cache.cache(ttl=100, key="user:{user_id}"),
))
def get_user(user_id):
    pass

4.36 自定义类型入参

2025-07 新增支持不可 JSON 序列化的入参类型(自动使用 pickle):

from pydantic import BaseModel
from funboost import boost, BoosterParams, BrokerEnum

class Order(BaseModel):
    order_id: int
    items: list

@boost(BoosterParams(queue_name='custom_type_task', broker_kind=BrokerEnum.REDIS_ACK_ABLE))
def process_order(order: Order):
    print(f'处理订单 {order.order_id}, 商品数: {len(order.items)}')

if __name__ == '__main__':
    process_order.push(Order(order_id=1, items=['apple', 'banana']))
    process_order.consume()

4.37 启动消费的方式大全

funboost 提供多种消费启动方式,按场景选择:

# 方式1:单进程消费(最常用)
task_fun.consume()

# 方式2:多进程叠加并发
task_fun.multi_process_consume(4)  # 简写: task_fun.mp_consume(4)

# 方式3:消费组启动
BoostersManager.consume_group('my_group')

# 方式4:一键启动所有
BoostersManager.consume_all()

# 方式5:自动启动(装饰时)
@boost(BoosterParams(queue_name='x', is_auto_start_consuming_message=True))
def auto_task(x): pass

4.38 FunboostPool

funboost 提供两种线程池包装,API 兼容 ThreadPoolExecutor:

4.38.1 MemoryFunboostPool(纯内存,快速替代线程池)

from funboost import MemoryFunboostPool

pool = MemoryFunboostPool(queue_name='my_pool', concurrent_num=10, qps=5)
future = pool.submit(my_func, arg1, arg2)
result = future.result(timeout=10)

4.38.2 FunboostPool(带消息队列持久化)

from funboost import FunboostPool, BrokerEnum

pool = FunboostPool(queue_name='persistent_pool', broker_kind=BrokerEnum.REDIS_ACK_ABLE, qps=10)
pool.submit(my_func, arg1, arg2)

特性

MemoryFunboostPool

FunboostPool

持久化

无(内存)

支持(消息队列)

分布式

不支持

支持

配置灵活度

低(仅并发数、QPS)

高(所有 BoosterParams 参数)

API 兼容

ThreadPoolExecutor

ThreadPoolExecutor


4.100 控制变量法验证

验证原则:对 funboost 任何功能的疑问,都可以抽象精简为一个 time.sleep() + print('hello') 的 demo 来验证。

from funboost import boost, BoosterParams

@boost(BoosterParams(queue_name='test_queue',
    ...  # 修改各种参数测试效果
))
def f(x):
    time.sleep(10)     # 修改 sleep 大小测试耗时影响
    print(f'hello: {x}')
    return x

示例:验证超时杀死

import random
import time
from funboost import boost, BoosterParams

@boost(BoosterParams(queue_name='test_timeout', concurrent_num=5, function_timeout=20, max_retry_times=4))
def add(x, y):
    t_sleep = random.randint(10, 30)
    print(f'计算 {x} + {y},需要 {t_sleep} 秒')
    time.sleep(t_sleep)
    print(f'{x} + {y} = {x + y}')
    return x + y

if __name__ == '__main__':
    for i in range(100):
        add.push(i, i * 2)
    add.consume()

function_timeout=20 会自动杀死超过 20 秒的函数执行(不是杀死进程/脚本),然后触发重试。


4.200 分布式函数调度框架 QQ 群

QQ 群:189603256