| name | developing-funboost-broker |
| description | 当需要为 funboost 框架添加新的消息中间件时使用。触发场景:实现自定义 Publisher/Consumer 类、使用 register_custom_broker 注册新 broker、使用 override_cls 定制现有 broker 行为。关键词:new broker, AbstractPublisher, AbstractConsumer, register_custom_broker, consumer_override_cls, publisher_override_cls, 扩展中间件, 新增 broker。 |
开发 Funboost Broker 中间件
概述
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):
:
msg = ._client.receive(.queue_name, timeout=)
msg:
kw = {: msg.body, : msg}
._submit_task(kw)
():
kw[].ack()
():
funboost.core.serialization Serialization
._client.send(.queue_name, Serialization.to_json_str(kw[]))
register_custom_broker(, MyPublisher, MyConsumer)
():
x *
方式三: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):
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,
}
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 访问规范
value = self.consumer_params.broker_exclusive_config["my_key"]
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 — 编写与运行测试