开发 Funboost Broker 中间件
SkillCommunicationUse when adding a new message broker to the funboost framework. Trigger scenarios: implementing custom Publisher/Consumer classes, using register_custom_broker to register a new broker, using override_cls to customize existing broker behavior. Keywords: new broker, AbstractPublisher, AbstractConsume
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 Broker 中间件 skill
What this skill tells your AI
The instructions your AI receives, as published by ydf0509/funboost in .agents/skills/developing-funboost-broker/SKILL.md and read by ahel’s review.
概述
Funboost 支持 3 种方式添加新 broker:静态扩展(框架作者)、register_custom_broker(新中间件)、override_cls(mixin 定制)。
核心原则: 实现抽象方法,正确注册。非抽象方法一般应调用 super()(完全替换父类逻辑时除外)。
适用场景
- 为 funboost 添加全新的消息队列后端
- 使用 override_cls 定制现有 broker 行为
- 用
register_custom_broker在用户空间注册 broker - 修改特定场景的发布/消费流程
三种扩展方式
| 方式 | 适用场景 | 是否修改源码 |
|---|---|---|
| 静态扩展 | 框架作者添加核心 broker | 是 |
register_custom_broker | 用户添加全新 broker | 否 |
consumer_override_cls / publisher_override_cls | 用户定制现有 broker 行为 | 否 |
方式一:静态扩展(框架作者)
需要修改的文件清单:
funboost/constant.py— 在BrokerEnum中增加枚举funboost/funboost_config_deafult.py— 在BrokerConnConfig中添加连接配置(如需)funboost/core/broker_kind__exclusive_config_default_define.py— 注册专属配置默认值funboost/publishers/— 创建新的 publisher 文件funboost/consumers/— 创建新的 consumer 文件funboost/factories/broker_kind__publsiher_consumer_type_map.py— 注册映射
方式二:register_custom_broker(用户空间)
from funboost import register_custom_broker, boost, BoosterParams
from funboost.publishers.base_publisher import AbstractPublisher
from funboost.consumers.base_consumer import AbstractConsumer
class MyPublisher(AbstractPublisher):
def custom_init(self):
super().custom_init()
self._client = connect_to_my_mq()
def _publish_impl(self, msg: str):
self._client.send(self.queue_name, msg)
def clear(self):
self._client.purge(self.queue_name)
def get_message_count(self):
return self._client.queue_length(self.queue_name)
def close(self):
self._client.close()
class MyConsumer(AbstractConsumer):
def custom_init(self):
super().custom_init()
self._client = connect_to_my_mq()
def _dispatch_task(self):
while True:
msg = self._client.receive(self.queue_name, timeout=5)
if msg:
kw = {"body": msg.body, "raw_msg": msg}
self._submit_task(kw)
def _confirm_consume(self, kw):
kw["raw_msg"].ack()
def _requeue(self, kw):
# 注意:_requeue 被调用时 kw["body"] 已是 dict,需序列化后再入队
from funboost.core.serialization import Serialization
self._client.send(self.queue_name, Serialization.to_json_str(kw["body"]))
# 注册
register_custom_broker("MY_BROKER", MyPublisher, MyConsumer)
# 使用
@boost(BoosterParams(queue_name="test", broker_kind="MY_BROKER"))
def my_task(x):
return x * 2
方式三:override_cls(Mixin 混入)
from funboost import boost, BoosterParams, BrokerEnum
class MyConsumerMixin:
"""Mixin 混入到任意 broker 的 consumer 中"""
def _submit_task(self, kw):
self._pre_check(kw)
super()._submit_task(kw)
def _pre_check(self, kw):
pass
@boost(BoosterParams(
queue_name="custom_task",
broker_kind=BrokerEnum.REDIS_ACK_ABLE,
consumer_override_cls=MyConsumerMixin,
))
def my_task(x):
return x
Publisher 必须实现的方法
| 方法 | 是否抽象 | 说明 |
|---|---|---|
_publish_impl(msg) | 是 | 核心发布逻辑——必须实现。普通 MQ broker 收到 JSON 字符串;MEMORY_QUEUE/FASTEST_MEM_QUEUE 收到 dict |
clear() | 是 | 清空队列所有消息 |
get_message_count() | 是 | 返回队列深度 |
close() | 是 | 关闭连接(可以写 pass) |
custom_init() | 否 | 可选的初始化钩子 |
注意: 自定义 broker 的
_publish_impl接收的msg通常是 JSON 字符串(框架已完成序列化),直接写入中间件即可。内存队列例外,可能收到 dict。
Consumer 必须实现的方法
| 方法 | 是否抽象 | 说明 |
|---|---|---|
_dispatch_task() | 是 | 主循环:取消息,调用 self._submit_task(kw) |
_confirm_consume(kw) | 是 | 确认消费(ACK) |
_requeue(kw) | 是 | 消息重入队 |
custom_init() | 否 | 可选的初始化钩子 |
_dispatch_task 实现
_dispatch_task 负责从 MQ 拉取消息并交给框架处理。它必须阻塞(不能立即返回):
def _dispatch_task(self):
# pull 模式:自写循环
while True:
msg = self._client.receive(timeout=5)
if msg:
self._submit_task({"body": msg.body, "raw_msg": msg})
如果 MQ 客户端本身提供阻塞消费方法(如 RabbitMQ 的 channel.start_consuming(callback=...)),也可以直接调用它。
_dispatch_task 内部不需要捕获网络异常。框架通过 keep_circulating 包裹它——如果因网络断开等异常退出,框架会自动重新调用实现重连。
kw 字典结构
传给 _submit_task 的 kw 字典必须包含:
kw = {
"body": message_body_string, # JSON 字符串(与 _publish_impl 收到的 msg 相同)
# broker 特有字段,用于 ack/requeue:
# "receipt_handle": ..., # SQS 用
# "message": ..., # AMQP 用
# "channel": ..., # RabbitMQ 用
}
kw["body"]可以是 JSON 字符串,也可以是 dict。框架在_submit_task内部通过_convert_msg_before_run统一转为 dict(Serialization.to_dict(msg)兼容两种输入)。
super() 调用规则
| 场景 | 必须调 super()? | 调用位置 |
|---|---|---|
custom_init() | 是 | 先调 super(),再执行自己的初始化 |
_submit_task() 重写 | 是 | 先执行前置逻辑,再调 super() |
_run() / _async_run() | 是 | 在自己的上下文中包裹 super() |
_publish_impl() | 否 | 抽象方法——直接实现 |
_dispatch_task() | 否 | 抽象方法——直接实现 |
_confirm_consume() | 否 | 抽象方法——直接实现 |
broker_exclusive_config 访问规范
# 在 consumer 中
value = self.consumer_params.broker_exclusive_config["my_key"]
# 在 publisher 中
value = self.publisher_params.broker_exclusive_config["my_key"]
用 [] 方括号访问——如果拼错了 key,直接 KeyError 报错。不要用 .get(key),否则拼写错误会默默降级为默认值而不自知。
在 broker_kind__exclusive_config_default_define.py 中注册默认值:
register_broker_exclusive_config_default("MY_BROKER", {
"my_key": "default_value",
})
参考代码位置
- 基础 publisher:
funboost/publishers/base_publisher.py - 基础 consumer:
funboost/consumers/base_consumer.py - 动态扩展示例:
funboost/contrib/register_custom_broker_contrib/ - Mixin 示例:
funboost/contrib/override_publisher_consumer_cls/ - 工厂注册:
funboost/factories/broker_kind__publsiher_consumer_type_map.py
常见错误
| 错误 | 修正 |
|---|---|
非抽象方法忘记调 super() | 除抽象方法外始终调用 super() |
exclusive_config 用 .get() | 推荐用 [] 方括号访问已注册的 key |
_dispatch_task 中没调 self._submit_task(kw) | 每条消息必须调用此方法 |
kw 字典缺少 body 键 | _submit_task 必须有 kw["body"] |
| 过度防御性编程(到处 try) | funboost 偏好简洁代码,让异常正常抛出 |
| 静态扩展后没注册到工厂映射 | 必须在 broker_kind__publsiher_consumer_type_map.py 中注册 |
测试新 Broker
实现后编写测试验证功能:
- AI 写测试放在
tests/ai_codes/regression_testing/或tests/ai_codes/ai_demos/{子文件夹}/ - 发布消息 -> 启动消费 -> 验证消费结果
- 运行约 30 秒后 kill(funboost 不会自动停止)
- AI 测试时使用
timeout或os._exit自动终止
相关 Skill
developing-funboost-mixin— Consumer/Publisher Mixin 扩展developing-funboost-testing— 编写与运行测试
Signals
- GitHub stars
- 891
- Forks
- 166
- Last commit
- Aug 2026
Advanced
- Catalog kind
- skill
- Gateway key
developing-funboost-broker- Source
- github.com/ydf0509/funboost