Funboost 基础使用
SkillCommunicationUse when writing distributed tasks with funboost. Trigger scenarios: creating @boost decorator task functions, using push/publish to send messages, starting consume consumers, configuring BoosterParams, and setting concurrency and rate limiting. Keywords: boost, BoosterParams, queue_name, push, publ
Available today. Use it from your connected AI after setup.
No other account needed.
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/using-funboost-basics/SKILL.md and read by ahel’s review.
概述
Funboost 用一个 @boost 装饰器把任意 Python 函数变成分布式任务。零侵入设计——func(x, y) 直接本地运行,func.push(x, y) 发送到队列。
核心原则: 你的函数保持为普通函数,不需要任何框架改造。
适用场景
- 创建新的分布式任务函数
- 向任务队列发布消息
- 启动任务消费者(Worker)
- 将现有函数改造为分布式任务
- 配置基础并发和限流
铁律(绝对不可违反)
- 必须使用
BoosterParams对象 — 禁止向@boost传递裸参数 - 禁止使用 Celery 模式 — 不用
self、不用bind=True,获取上下文用fct push只传业务参数;publish用字典传业务参数,并通过task_options=TaskOptions(...)附加框架控制参数
速查表
| 操作 | 代码 |
|---|---|
| 定义任务 | @boost(BoosterParams(queue_name="q1")) |
| 推送消息 | func.push(x, y) |
| 带控制选项发布 | func.publish({"x": 1}, task_options=TaskOptions(countdown=5)) |
| 异步推送 | await func.aio_push(x, y) |
| 启动消费 | func.consume() |
| 多进程消费 | func.multi_process_consume(3) |
| 获取任务上下文 | from funboost import fct; fct.task_id |
| Windows 下 Ctrl+C 退出 | enable_ctrl_c_quit_on_windows()(可选,不加也能运行) |
核心代码模式
以下是零依赖最小示例(使用 MEMORY_QUEUE,无需 Redis/RabbitMQ):
from funboost import boost, BoosterParams, BrokerEnum
@boost(BoosterParams(
queue_name="hello_funboost",
broker_kind=BrokerEnum.MEMORY_QUEUE,
))
def add(a, b):
print(f"计算: {a} + {b} = {a + b}")
return a + b
if __name__ == "__main__":
add.push(1, 2)
add.push(10, 20)
add.consume()
# funboost 消费者永久运行,不会自动退出(和 Celery worker 一样)
# 如果是 Windows 想用 Ctrl+C 停止:
from funboost import enable_ctrl_c_quit_on_windows
enable_ctrl_c_quit_on_windows()
以下示例使用 Redis(需先在 funboost_config.py 中配置 REDIS_HOST 等参数):
from funboost import boost, BoosterParams, BrokerEnum
@boost(BoosterParams(
queue_name="my_task_queue",
broker_kind=BrokerEnum.REDIS_ACK_ABLE,
concurrent_num=30,
qps=10,
max_retry_times=3,
log_level=20,
))
def my_task(url: str, depth: int = 1):
"""你的业务逻辑——保持为普通函数"""
import requests
resp = requests.get(url)
return resp.status_code
if __name__ == "__main__":
# 发布消息
for i in range(100):
my_task.push(f"https://example.com/page/{i}", depth=2)
# 启动消费
my_task.consume()
发布方法详解
push — 只传业务参数
my_task.push("https://example.com", depth=3)
publish — 附带框架控制参数(countdown、task_id 等)
from funboost import TaskOptions
my_task.publish(
{"url": "https://example.com", "depth": 3},
task_options=TaskOptions(
countdown=10, # 延迟 10 秒执行
task_id="custom-id-1", # 自定义 task ID
)
)
异步发布(须在 async def 内)
async def publish_tasks():
await my_task.aio_push("https://example.com", depth=3)
await my_task.aio_publish({"url": "..."}, task_options=TaskOptions(countdown=5))
启动多个消费者
task_a.consume()
task_b.consume()
task_c.consume()
绝对禁止 用 threading.Thread 包装 consume() — 它本身就是非阻塞的。
任务上下文 (fct)
from funboost import fct
@boost(BoosterParams(queue_name="ctx_demo"))
def my_task(x):
print(f"Task ID: {fct.task_id}")
print(f"队列名: {fct.queue_name}")
print(f"执行次数: {fct.function_result_status.run_times}")
print(f"函数参数: {fct.function_params}")
print(f"完整消息: {fct.full_msg}")
fct.logger.info("当前任务 logger")
常用属性:
| 属性 | 说明 |
|---|---|
fct.task_id | 当前任务 ID |
fct.queue_name | 队列名 |
fct.function_result_status.run_times | 运行次数(含重试) |
fct.function_params | 函数入参 |
fct.full_msg | 完整消息体 |
fct.logger | 当前任务 logger |
BoosterParams 核心字段
| 字段 | 类型 | 默认值 | 说明 |
|---|---|---|---|
queue_name | str | — | 必填,队列名 |
broker_kind | str | SQLITE_QUEUE | 消息中间件类型 |
concurrent_num | int | 50 | 并发数量 |
concurrent_mode | str | ConcurrentModeEnum.THREADING | threading/gevent/eventlet/async/single_thread |
qps | float/int/None | None | 每秒执行次数限制 |
max_retry_times | int | 3 | 最大重试次数 |
function_timeout | int/float/None | None | 函数超时秒数 |
log_level | int | 10 (DEBUG) | 日志级别 |
is_using_rpc_mode | bool | False | 是否启用 RPC 获取结果 |
常见错误
| 错误写法 | 正确写法 |
|---|---|
@boost("queue", qps=5) | @boost(BoosterParams(queue_name="queue", qps=5)) |
def task(self, x): 获取上下文 | 使用 fct.task_id |
timeout=30 | function_timeout=30 |
max_retries=5 | max_retry_times=5 |
| 用 threading 启动多个消费者 | 直接顺序调用 func1.consume(); func2.consume() |
func.push(msg_dict) 带控制参数 | 使用 func.publish(msg_dict, task_options=TaskOptions(...)) |
obj.method.push(arg1) 实例方法 push | ClassName.method.push(obj_instance, arg1, arg2),第一个参数传对象实例 |
实例方法 push 语法: 必须写成
ClassName.method.push(obj_instance, arg1, arg2),第一个参数传对象实例。禁止写obj.method.push(arg1)。
消费来自其他系统的消息
当消费非 funboost 发布的消息(如 Java/Go 写入的)时:
@boost(BoosterParams(
queue_name="external_queue",
should_check_publish_func_params=False,
))
def handle_external(**kwargs):
"""使用 **kwargs 接收任意 JSON 结构"""
print(kwargs)
绝对禁止 用 def handle(msg): 单参数接收整个字典——必须用 **kwargs 或 **msg 解包接收。
相关 Skill
understanding-funboost-concepts— 框架概念入门funboost-rpc-mode— 获取任务执行返回值funboost-async-programming— async/await 异步编程funboost-broker-selection— Broker 中间件选型
Signals
- GitHub stars
- 891
- Forks
- 166
- Last commit
- Aug 2026
Advanced
- Catalog kind
- skill
- Gateway key
using-funboost-basics- Source
- github.com/ydf0509/funboost