Keyboard shortcuts

Press or to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

定义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 很像,也需要 entityfilterprojection,但执行函数拿到的不是现成 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 都有的属性

属性是否必填说明
namewatcher 的唯一名字
entity要扫描的实体
filter查询或操作条件;WB 类型可写成异步函数,BB 类型只接受对象或同步函数
singleton集群环境下只允许一个实例执行
lazy启动后的第一轮跳过执行

WBWatcher / WBFreeWatcher 额外属性

属性是否必填说明
projection查询出的字段
fn处理函数
forUpdate查询时是否加锁,适合需要串行修改的场景
exclusive同一实例中跳过仍在处理的同一行;只支持 WB 类型

这里的 filterprojectionforUpdate 并不是 watcher 自己发明的新语法,而是直接复用 Oak 查询语法:

  • filter / projection 的写法,参见查询和操作对象
  • forUpdate 对应的就是 SelectOption.forUpdate
  • 如果 watcher 要筛一批待处理数据,再逐条更新,这里通常就应该考虑是否需要加锁。

BBWatcher 额外属性

属性是否必填说明
action对目标实体执行的动作
actionData固定写入的数据,也可以写成异步函数

BBWatcher 不支持 exclusive;即使配置,当前 AppLoader 也只会输出警告并忽略。需要按行排他处理时,应改用 WBWatcherWBFreeWatcher

执行函数长什么样

WBWatcher

fn: async (context, data) => {
    // data 是查出来的多行结果
    return context.opResult;
}

WBFreeWatcher

fn: async (builder, data) => {
    const context = await builder();
    // 自行控制 context 的使用
    return context.opResult;
}

Oak 在执行 WBWatcher 时,会先用 filterprojection 做一次查询,然后把查询结果数组传给 fn。所以 watcher 的思考方式不是“监听事件”,而是“定期扫描待处理行”。

lazysingletonexclusive 的区别

这三个参数很容易混淆,但它们解决的是完全不同的问题:

  • lazy:应用刚启动后的第一轮先跳过;
  • singleton:在多实例部署时,只让一个实例执行;
  • exclusive:在同一个实例中,如果某条数据上一次还没处理完,本次不要并发重复处理。

例如 bm-smart/src/timers/index.ts 中的示例 timer 就使用了 exclusive: true,这类控制同样适用于 watcher。

什么时候该用 watcher

下面这些情况,通常就该优先想到 watcher:

  • “数据库里只要有这种状态的数据,就要持续处理”
  • “上一次失败了,后面要自动重试”
  • “这个逻辑不要求精确到某个秒级 cron 点,只要求不断收敛”

如果你的需求是“每天凌晨 1 点一定执行一次”,那应该用 timer;如果你的需求是“某个事务提交后立即补偿”,那应该优先考虑 commit trigger

使用 watcher 时的一个原则

watcher 最适合处理幂等、可重复扫描、允许延迟收敛的后台任务。

因为它本质上是轮询。如果你把一个必须“立刻且只执行一次”的动作塞进 watcher,最后往往会把系统设计得更复杂,而不是更简单。