开发 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):
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— 编写与运行测试