定义watcher
watcher 是 Oak 后端中的轮询执行器。它和 timer 最大的区别在于:watcher 的调度周期是框架固定的,AppLoader.startWatchers() 会在每一轮执行结束后,用 setTimeout(..., 120000) 再次调度下一轮,也就是每 120 秒执行一次。
因此,watcher 适合处理这类问题:
- 扫描数据库中“待补偿”的数据;
- 重试失败的异步任务;
- 对某些状态做持续收敛;
- 周期性处理“只要符合条件就应被处理”的数据。
watcher 的三种类型
在 oak-domain/src/types/Watcher.ts 中,Oak 定义了三种 watcher:
BBWatcher
最简单的一类 watcher。它不把查询结果交给自定义函数,而是按 filter 直接对目标实体执行固定的 action + actionData。
WBWatcher
最常见的一类 watcher。它会先按 filter + projection 查出数据,再把结果数组交给你写的 fn(context, data) 进行处理。
WBFreeWatcher
它和 WBWatcher 很像,也需要 entity、filter、projection,但执行函数拿到的不是现成 context,而是 builder: () => Promise<context>。这适合某些需要自行控制上下文生命周期的复杂逻辑。
编写位置
一般把 watcher 写在 src/watchers 目录中,并在 src/watchers/index.ts 中统一导出。例如 bm-smart 中:
const watchers = [
...shellMessageWatchers,
] as Watcher<EntityDict, keyof EntityDict, BackendRuntimeContext>[];
export default watchers;
一个典型的 WBWatcher
下面这个例子来自 bm-smart/src/watchers/shellMessage.ts,它会周期性重试发送超时的 Shell 消息:
{
name: '定期重试发送失败的Shell消息',
entity: 'shellMessage',
projection: {
id: 1,
data: 1,
action: 1,
shellId: 1,
error: 1,
},
filter: {
status: 'timeout',
retryOnTimeout: true,
ackAt: {
$exists: false,
}
},
lazy: true,
async fn(context, data) {
for (const msg of data) {
try {
await retryShellMessage(msg);
} catch (err) {
console.error(err);
}
}
return context.opResult;
},
}
这个例子几乎把 watcher 的典型写法都展示出来了:
- 用
filter限定待处理数据; - 用
projection只取处理所需字段; - 用
lazy: true避免应用刚启动时立刻跑一次; - 在
fn里显式处理每一条数据。
这里按行捕获异常,是因为该任务明确选择“一条失败不阻塞同批其他行”的部分进度策略;异常仍应记录并让失败行继续满足下轮 filter。普通 watcher 如果要求整批事务一致,不要 catch 后吞掉异常,应该让异常抛出,由 AppLoader 回滚并统一记录内部错误。
配置项说明
所有 watcher 都有的属性
| 属性 | 是否必填 | 说明 |
|---|---|---|
name | 是 | watcher 的唯一名字 |
entity | 是 | 要扫描的实体 |
filter | 是 | 查询或操作条件;WB 类型可写成异步函数,BB 类型只接受对象或同步函数 |
singleton | 否 | 集群环境下只允许一个实例执行 |
lazy | 否 | 启动后的第一轮跳过执行 |
WBWatcher / WBFreeWatcher 额外属性
| 属性 | 是否必填 | 说明 |
|---|---|---|
projection | 是 | 查询出的字段 |
fn | 是 | 处理函数 |
forUpdate | 否 | 查询时是否加锁,适合需要串行修改的场景 |
exclusive | 否 | 同一实例中跳过仍在处理的同一行;只支持 WB 类型 |
这里的 filter、projection、forUpdate 并不是 watcher 自己发明的新语法,而是直接复用 Oak 查询语法:
filter/projection的写法,参见查询和操作对象;forUpdate对应的就是SelectOption.forUpdate;- 如果 watcher 要筛一批待处理数据,再逐条更新,这里通常就应该考虑是否需要加锁。
BBWatcher 额外属性
| 属性 | 是否必填 | 说明 |
|---|---|---|
action | 是 | 对目标实体执行的动作 |
actionData | 是 | 固定写入的数据,也可以写成异步函数 |
BBWatcher 不支持 exclusive;即使配置,当前 AppLoader 也只会输出警告并忽略。需要按行排他处理时,应改用 WBWatcher 或 WBFreeWatcher。
执行函数长什么样
WBWatcher
fn: async (context, data) => {
// data 是查出来的多行结果
return context.opResult;
}
WBFreeWatcher
fn: async (builder, data) => {
const context = await builder();
// 自行控制 context 的使用
return context.opResult;
}
Oak 在执行 WBWatcher 时,会先用 filter 和 projection 做一次查询,然后把查询结果数组传给 fn。所以 watcher 的思考方式不是“监听事件”,而是“定期扫描待处理行”。
lazy、singleton、exclusive 的区别
这三个参数很容易混淆,但它们解决的是完全不同的问题:
lazy:应用刚启动后的第一轮先跳过;singleton:在多实例部署时,只让一个实例执行;exclusive:在同一个实例中,如果某条数据上一次还没处理完,本次不要并发重复处理。
例如 bm-smart/src/timers/index.ts 中的示例 timer 就使用了 exclusive: true,这类控制同样适用于 watcher。
什么时候该用 watcher
下面这些情况,通常就该优先想到 watcher:
- “数据库里只要有这种状态的数据,就要持续处理”
- “上一次失败了,后面要自动重试”
- “这个逻辑不要求精确到某个秒级 cron 点,只要求不断收敛”
如果你的需求是“每天凌晨 1 点一定执行一次”,那应该用 timer;如果你的需求是“某个事务提交后立即补偿”,那应该优先考虑 commit trigger。
使用 watcher 时的一个原则
watcher 最适合处理幂等、可重复扫描、允许延迟收敛的后台任务。
因为它本质上是轮询。如果你把一个必须“立刻且只执行一次”的动作塞进 watcher,最后往往会把系统设计得更复杂,而不是更简单。