这里记录一种简单的进程内发布—订阅方式:下游先把处理函数交给通知器,上游在配置发布成功后调用这些函数。
以异步写法为主。同步写法的结构相同,只有函数声明、回调类型和调用方式不同,差异直接标在代码注释中。
一、消息订阅关系图
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 asynciofrom collections.abc import Awaitable, Callablefrom 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:
def create_app(...) -> FastAPI: app = FastAPI(...)
# 应用进程内只创建一次,因此启动时注册的函数不会随请求结束而丢失。 app.state.config_change_notifier = InProcessStrategyConfigChangeNotifier() return app每次请求创建发布服务时,再把这一个通知器传进去:
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():
# 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 或消息队列。
评论