开发 Funboost Broker 中间件

SkillCommunication

Use 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.

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 行为

方式一:静态扩展(框架作者)

需要修改的文件清单:

  1. funboost/constant.py — 在 BrokerEnum 中增加枚举
  2. funboost/funboost_config_deafult.py — 在 BrokerConnConfig 中添加连接配置(如需)
  3. funboost/core/broker_kind__exclusive_config_default_define.py — 注册专属配置默认值
  4. funboost/publishers/ — 创建新的 publisher 文件
  5. funboost/consumers/ — 创建新的 consumer 文件
  6. 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 测试时使用 timeoutos._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