# Developing Funboost Broker

> 当需要为 funboost 框架添加新的消息中间件时使用。触发场景：实现自定义 Publisher/Consumer 类、使用 register_custom_broker 注册新 broker、使用 override_cls 定制现有 broker 行为。关键词：new broker, AbstractPublisher, AbstractConsumer, register_custom_broker, consumer_override_cls, publisher_override_cls, 扩展中间件, 新增 broker。

- Skill: `ydf0509/developing-funboost-broker` (Agent Skill)
- Install (CLI): `npx skillmds@latest add ydf0509/developing-funboost-broker`
- Raw SKILL.md: https://api.skillmd.com/api/skills/ydf0509/developing-funboost-broker/raw
- Safety review: pending
- Works with: Claude Code, Claude.ai, OpenAI Codex
- Category: Coding & Dev Tools
- Author: ydf0509 (https://skillmd.com/u/ydf0509)
- Updated: 2026-09-17
- Page: https://skillmd.com/skills/ydf0509/developing-funboost-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 行为 | 否 |

## 方式一：静态扩展（框架作者）

需要修改的文件清单：

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（用户空间）

```python
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 混入）

```python
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 拉取消息并交给框架处理。它必须阻塞（不能立即返回）：

```python
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 字典必须包含：

```python
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 访问规范

```python
# 在 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` 中注册默认值：

```python
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` — 编写与运行测试

