4. 使用框架的各种代码示例
框架极其简单并且自由,只有一个 @boost 装饰器的参数学习。本章所有例子都是调整 @boost 的参数而已。
消息中间件的 ip/端口/密码等配置在首次运行代码时,框架会在项目根目录自动生成 funboost_config.py,按需修改即可。
所有例子的发布和消费都不必写在同一个 py 文件(使用中间件解耦)。
4.索引 快速导航
索引1. 基础用法
章节 |
内容 |
|---|---|
@boost 装饰器入参格式(BoosterParams) |
|
装饰器方式调度函数(最简示例) |
|
动态按队列名生成 booster(BoostersManager.build_booster) |
|
BoostersManager 管理(一键启动所有消费) |
|
push 和 publish 发布消息的区别 |
索引2. 消费控制
章节 |
内容 |
|---|---|
多步骤消费(step1 → step2 链式调度) |
|
多个消费者共享同一线程池 |
|
清空队列 / 获取队列消息数量 |
|
暂停 / 恢复消费 |
|
is_auto_start_consuming_message 自动启动 |
索引3. 并发与性能
章节 |
内容 |
|---|---|
多进程 + 多线程/协程叠加并发 |
|
QPS 精准控频(核心功能) |
|
asyncio 方式运行协程 |
|
funboost 如何取代手写线程池 |
|
性能调优演示 |
索引4. 定时与延时
章节 |
内容 |
|---|---|
定时运行(ApsJobAdder) |
|
延时运行任务 |
索引5. 高级功能
章节 |
内容 |
|---|---|
RPC 模式(远程调用获取结果) |
|
设置消费函数重试次数 |
|
任务优先级队列 |
|
远程杀死/取消任务 |
|
fct 上下文获取当前消息和任务状态 |
|
函数入参过滤(去重) |
|
自定义类型入参(pickle 序列化) |
|
MemoryFunboostPool / FunboostPool |
索引6. 集成与扩展
章节 |
内容 |
|---|---|
Flask / FastAPI / Django 集成 |
|
保存消费状态到 MongoDB |
|
跨项目发布任务 |
|
自定义消费状态/结果钩子 |
|
broker_exclusive_config 差异化配置 |
|
register_custom_broker 自定义中间件 |
|
consumer_override_cls / publisher_override_cls |
|
Celery 框架整体作为 funboost 的 Broker |
索引7. 其他
章节 |
内容 |
|---|---|
获取消费进程信息 |
|
文件日志所在位置 |
|
判断函数运行完所有任务 |
|
实例方法/类方法作为消费函数 |
|
PyInstaller 打包说明 |
|
funboost 启动消费的方式大全 |
|
使用控制变量法验证框架功能 |
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 的优势
用户无需手写
push_msg包装函数Redis 作为 job_store 时,每个函数使用独立的 jobs_key
使用 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(重回队列,不占线程,推荐长间隔场景)。
参数 |
说明 |
|---|---|
|
|
|
基础间隔秒数 |
|
指数退避倍数(1.0 为固定间隔) |
|
最大间隔秒数 |
|
是否随机抖动(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 更符合直觉,publish 更灵活(适合动态构造参数字典)。
4.26 性能调优
提高 concurrent_num:增大并发线程数
使用 multi_process_consume:多进程叠加多线程
选择高性能 broker:REDIS > RABBITMQ > KAFKA(对于小消息)
使用 gevent/eventlet:IO 密集型任务可获得更高并发
减少日志:
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