这里记录一种简单的进程内发布—订阅方式:下游先把处理函数交给通知器,上游在配置发布成功后调用这些函数。

以异步写法为主。同步写法的结构相同,只有函数声明、回调类型和调用方式不同,差异直接标在代码注释中。

一、消息订阅关系图

flowchart LR
subgraph PUB["发布模块"]
direction TB
MAIN["main"]
end
subgraph EVENT["ConfigChangedEvent 事件类(通知数据类型)"]
direction TB
EVENT_INIT["init()"]
end
subgraph NOTIFIER["ConfigChangeNotifier 通知器类"]
direction TB
NOTIFIER_INIT["init()"]
REGISTER["register"]
NOTIFY["notify"]
end
subgraph SUB["下游模块"]
direction TB
RELOAD["reload_tasks"]
end
MAIN -->|"创建通知器"| NOTIFIER_INIT
MAIN -->|"config_id 和 data"| EVENT_INIT
MAIN -->|"调用"| REGISTER
RELOAD -.->|"作为 callback"| REGISTER
REGISTER -.->|"写入 callbacks 列表"| NOTIFY
EVENT_INIT -->|"event,由 main 传入"| NOTIFY
NOTIFY -->|"await callback event"| RELOAD

外层大框表示类或模块,内层小框表示函数。数据和调用参数不单独成块,只标在线条上。图中的名称与下一节示例代码一一对应。

二、基本示例

import asyncio
from collections.abc import Awaitable, Callable
from dataclasses import dataclass
@dataclass(frozen=True)
class ConfigChangedEvent:
"""一次已经发布的配置变更。"""
config_id: str
data: dict[str, object]
# 异步回调返回 Awaitable[None]。
# 同步写法:Callable[[ConfigChangedEvent], None]
ConfigChangedCallback = Callable[[ConfigChangedEvent], Awaitable[None]]
class ConfigChangeNotifier:
def __init__(self) -> None:
# 保存函数本身;注册时不会执行函数。
self._callbacks: list[ConfigChangedCallback] = []
def register(self, callback: ConfigChangedCallback) -> None:
self._callbacks.append(callback)
async def notify(self, event: ConfigChangedEvent) -> None:
# 同步写法:def notify(...)
# 当前按注册顺序逐个执行,前一个完成后才执行下一个。
for callback in self._callbacks:
await callback(event)
# 同步写法:callback(event)
# 每个下游只负责自己的处理逻辑。
async def reload_tasks(event: ConfigChangedEvent) -> None:
# 同步写法:def reload_tasks(...)
print("重新加载配置:", event.config_id, event.data)
async def main() -> None:
# 同步写法:def main() -> None
notifier = ConfigChangeNotifier()
# 通常在应用启动时注册一次。
notifier.register(reload_tasks)
# 上游完成配置发布后发送事件。
event = ConfigChangedEvent("config-1", {"enabled": True})
await notifier.notify(event)
# 同步写法:notifier.notify(event)
asyncio.run(main())
# 同步写法不需要 asyncio,直接调用 main() 即可。

同步和异步的选择只看回调是否需要异步 I/O。只改内存或做少量计算时,同步函数就够了;要异步访问数据库、HTTP 接口等资源时,再使用 async def 和 await。

三、在 TradeOps 中的应用

项目关系图

flowchart LR
subgraph startup[应用启动]
app[应用装配<br/>create_app]
notifier[进程内通知器<br/>InProcessStrategyConfigChangeNotifier]
subscriber[下游订阅者<br/>当前尚未接入]
app -->|创建并持有| notifier
subscriber -.->|注册处理函数| notifier
end
subgraph publish[用户发布配置]
api[发布入口<br/>POST /json-sync]
query[配置汇总<br/>StrategyOverviewQueryService]
service[发布服务<br/>StrategyOverviewJsonSyncService]
snapshot[快照持久化<br/>Repository + Database]
event[变更事件<br/>StrategyConfigChangedEvent]
failed[结束,不通知]
api --> service
query -->|完整 overview JSON| service
service -->|保存快照| snapshot
snapshot -->|commit 成功| service
snapshot -->|commit 失败| failed
service -->|提交后创建| event
end
notifier -->|依赖注入| service
event -->|发布| notifier
notifier -.->|依次调用处理函数| subscriber

实线表示当前项目已有的调用或数据流,虚线表示预留的下游接入点。

业务场景

TradeOps 会把当前全部策略配置整理成一份 overview JSON。用户确认发布后,系统先把 JSON 保存为快照,再通知同一进程中的下游:有一份新配置可以处理了。下游可能据此重载策略任务,但具体处理不属于发布流程。

1. 定义通知内容和函数约定

文件:backend/app/services/configuration/new_strategy_management/config_change_notifier.py

@dataclass(frozen=True, slots=True)
class StrategyConfigChangedEvent:
# 快照 UUID,同时标识本次通知;项目没有额外的 version 字段。
config_id: UUID
# 快照生成并保存的时间。
changed_at: datetime
# 本次发布的完整 overview JSON。
data: Mapping[str, object]
# 下游函数必须接收一个事件,异步处理且不返回业务结果。
StrategyConfigChangedCallback = Callable[
[StrategyConfigChangedEvent],
Awaitable[None],
]
# 如果项目改用同步回调,这里应为 Callable[[StrategyConfigChangedEvent], None]。
class StrategyConfigChangeNotifier(Protocol):
# 发布服务依赖这份协议,不依赖某个具体下游。
def register_config_changed_callback(
self,
callback: StrategyConfigChangedCallback,
) -> None: ...
async def notify_config_changed(
self,
event: StrategyConfigChangedEvent,
) -> None: ...
# 同步写法去掉 async,调用方也不再 await。

2. 保存和调用下游函数

同一文件中的进程内实现只有一份回调列表:

class InProcessStrategyConfigChangeNotifier:
def __init__(self) -> None:
self._callbacks: list[StrategyConfigChangedCallback] = []
def register_config_changed_callback(
self,
callback: StrategyConfigChangedCallback,
) -> None:
self._callbacks.append(callback)
async def notify_config_changed(
self,
event: StrategyConfigChangedEvent,
) -> None:
# tuple() 固定本轮遍历对象,避免回调执行期间修改列表而影响本轮通知。
for callback in tuple(self._callbacks):
await callback(event)
# 同步实现改为 callback(event)。

当前行为很直接:回调按注册顺序逐个等待;某个回调抛出异常时,后续回调不会执行,异常继续交给调用方处理。

3. 让所有请求共用同一个通知器

应用创建时把通知器放进 app.state:

backend/app/main.py
def create_app(...) -> FastAPI:
app = FastAPI(...)
# 应用进程内只创建一次,因此启动时注册的函数不会随请求结束而丢失。
app.state.config_change_notifier = InProcessStrategyConfigChangeNotifier()
return app

每次请求创建发布服务时,再把这一个通知器传进去:

backend/app/api/deps.py
def get_strategy_overview_json_sync_service(
request: Request,
session: Session = Depends(get_db_session),
) -> StrategyOverviewJsonSyncService:
return StrategyOverviewJsonSyncService(
session=session,
query_service=get_strategy_overview_query_service(request, session),
snapshot_repository=SqlAlchemyStrategyOverviewJsonSnapshotRepository(session),
# 所有请求取得的是同一个应用级通知器。
notifier=request.app.state.config_change_notifier,
)

4. 快照提交成功后再通知

用户调用 POST /json-sync 确认发布,接口最终进入 StrategyOverviewJsonSyncService.apply():

backend/app/services/configuration/new_strategy_management/
# config_overview_json_sync_service.py
async def apply(...) -> StrategyOverviewJsonSyncApplyResponse:
# 前面已经完成三项检查:
# 1. 相同 sync_id 是否已经处理;
# 2. 预览指纹是否仍与当前数据一致;
# 3. 新 JSON 是否确实发生变化。
snapshot = StrategyOverviewJsonSnapshot(
id=sync_id,
payload_json=current_payload,
updated_at=datetime.now(UTC),
)
try:
self._snapshot_repository.add(snapshot)
self._session.flush()
self._session.commit()
except SQLAlchemyError:
# 提交失败会进入现有错误处理,不会走到 notify_config_changed()。
...
# 数据库已经提交,此时才把快照信息交给所有下游。
await self._notifier.notify_config_changed(
StrategyConfigChangedEvent(
config_id=snapshot.id,
changed_at=snapshot.updated_at,
data=snapshot.payload_json,
)
)
return _build_apply_response(snapshot, idempotent=False)

因此只有“数据有变化且新快照提交成功”会触发通知。预览、无变化、重复提交和数据库提交失败都不会通知。

需要注意,通知发生在提交之后。回调失败可以让接口报错,却不能回滚已经保存的快照。这是当前简单进程内通知的边界,不等同于可靠消息投递。

5. 下游如何接入

当前仓库只提供通知器和发布端,还没有注册实际下游。将来接入任务管理器时,只需在应用启动阶段注册一次:

class StrategyTaskManager:
async def on_strategy_config_changed(
self,
event: StrategyConfigChangedEvent,
) -> None:
# 下游拿到的是完整 JSON,自行决定如何更新任务。
await self.reload_tasks(event.data)
task_manager = StrategyTaskManager(...)
app.state.config_change_notifier.register_config_changed_callback(
task_manager.on_strategy_config_changed,
)

这里注册的是绑定方法本身,不是调用结果,所以不能写成 register_config_changed_callback(task_manager.on_strategy_config_changed(...))。

6. 当前验证范围

# backend/tests/test_config_change_notifier.py 当前验证:
# 注册后的回调可以连续收到事件。
self.assertEqual(received, [event, event])
# 快照提交成功后,事件带有本次 UUID 和完整 JSON。
self.assertEqual(received[0].config_id, preview.sync_id)
self.assertEqual(received[0].data, {"ACCOUNT_A": {"enabled": True}})
# commit 抛出异常时,回调没有执行。
callback.assert_not_called()

四、总结

  • register(callback) 建立订阅,notify(event) 发布消息,callback(event) 处理消息。
  • 事件对象是上下游的数据约定;增加字段时,不必修改通知函数的参数列表。
  • 本项目在数据库提交成功后发布完整策略配置,通知器由应用进程共享。
  • 同步与异步只有函数声明和调用方式的差别;是否需要异步取决于下游有没有异步 I/O。
  • 当前方案只适合同一 Python 进程,不提供注销、并发执行、失败隔离、重试或可靠投递。下游独立成服务后,应改用 HTTP、RPC 或消息队列。