@mnemora/bullmq
v1.2.0
Published
BullMQ で runtime.tick() を駆動する役。outbox は Postgres のまま正本——BullMQ はジョブの中身を持たず、「いま tick して」の合図だけを運ぶ(docs/decisions/0325-bullmq-tick-driver.md)。
Downloads
90
Readme
@mnemora/bullmq
BullMQ で runtime.tick() を駆動する役
(docs/decisions/0325-bullmq-tick-driver.md)。
npm への公開
v1.1.0 から npm に出ている(Issue #205)。初版 1.1.0 は 2026-09-30 にオーナーが手元から
publish した(bootstrap。手順は docs/release-v1.md の @mnemora/bullmq
初回 publish の節)。⚠ この初版には provenance が付いていない——手元からの publish は
OIDC を経由しないため。次の版からは、ほかの @mnemora/* と同じく Release の publish ワークフローが
provenance 付きで上げる。
⚠ Scheduler を実装しない
@mnemora/core の Scheduler interface(enqueue)は実装しない——本番コードのどこからも
呼ばれておらず、実装しても呼び手が無い(ADR 0325「根拠①への応答」)。このパッケージが
することは、BullMQ の Worker が定期的に発火するたびに runtime.tick() を呼ぶ、
それだけである。
outbox は今日どおり Postgres が正本のまま。 BullMQ(Redis)はジョブの中身を一切
持たない——運ぶのは「いま tick して」という合図だけであり、outbox の行と Redis 側の
ジョブが二重に帳簿を持つことはない。同時 tick からの二重処理を防いでいるのは
@mnemora/postgres の PostgresOutboxStore.claimBatch(FOR UPDATE SKIP LOCKED)で
あって、このパッケージや BullMQ 自身ではない(詳しくは ADR 0325「測ったこと」)。
インストール
pnpm add @mnemora/bullmq @mnemora/core
# または
npm i @mnemora/bullmq @mnemora/coreRedis が要る。 BullMQ は Redis(または互換サーバ)への接続を前提にする
——connection オプションにその接続先を渡す(下の例参照)。このパッケージ自身は
Redis サーバを同梱・起動しない。
前提
- Node.js >= 22
- ESM のみ(
"type": "module")。CommonJS からは Node 22.12 以降のrequire(esm)で読み込める(TypeScript はmodule/moduleResolutionをnodenextにし、 TypeScript 5.8 以降を使うこと。5.7 以前のnodenextと、どの版のnode16もTS1479になる) - Redis(または互換サーバ)が要る。
test(pnpm --filter @mnemora/bullmq run test)は 純関数(resolveConcurrency)だけを検査し Redis を要らないが、test:redis(pnpm --filter @mnemora/bullmq run test:redis)は実際に BullMQ のQueue/Workerを 構築するため Redis を要る(.github/workflows/ci.ymlのbullmqjob はredis:7の service container を使う) runtimeはPick<Runtime, "tick">——@mnemora/coreのcreateRuntime()が返すRuntime全体ではなく、tickメソッドさえ満たせば渡せる
動く最小の例(Redis が無いため未実行——型のみ確認)
import { createBullmqTickDriver } from "@mnemora/bullmq";
import type { Runtime } from "@mnemora/core";
declare const runtime: Runtime;
const driver = createBullmqTickDriver({
connection: { host: "127.0.0.1", port: 6379 },
queueName: "mnemora-tick",
runtime,
ctx: { tenantId: "acme" },
tick: { leaseMs: 30 * 60 * 1000, kinds: ["embed"] },
everyMs: 5_000,
});
await driver.start();
// ... プロセスが生きている間、5秒おきに runtime.tick() が呼ばれる ...
await driver.stop();start() を呼ぶまでジョブは処理しない。 createBullmqTickDriver(...) は
Queue/Worker を構築するだけで、Worker は autorun: false で作る——ジョブの処理は
start() が明示的に worker.run() を呼んで初めて始まる。
stop() の後は再開できない。 stop() を呼んだ driver は使い捨てである。その後に
もう一度 start() を呼ぶと Error を投げる(BullMQ の Queue/Worker は close() した後、
同じインスタンスを再利用できないため)。もう一度動かしたいときは
createBullmqTickDriver(...) を新しく呼び直すこと。
複数プロセスで動かすとき
同じ queueName に対して複数プロセスが createBullmqTickDriver(...).start() を
呼んでよい。 BullMQ の Job Scheduler(queue.upsertJobScheduler)を jobSchedulerId
固定値で登録するため、二重登録にはならない——発火した個々の tick ジョブは、その時点で
空いているどのプロセスの Worker が処理してもよい(BullMQ の通常の負荷分散)。
🔴 これは「同じテナントに対して2つの runtime.tick() が同時に走らない」ことを
保証しない。 その重なりから outbox の二重処理を防いでいるのは、上に書いたとおり
@mnemora/postgres 側の行ロックである。
🔴 1台の stop() が、全プロセスの予定を止める。stop() は Worker と Queue を閉じる前に
queue.removeJobScheduler(jobName) を呼ぶ。この scheduler は上のとおり全プロセスで共有している
(同じ jobSchedulerId)ので、1つのプロセスが stop() すると、他のプロセスの Worker は動いたままでも
tick のジョブがもう発火しない。エラーにもならない。動いている driver の start() をもう一度呼んでも
何もしない(冪等)ので、登録はし直されない。新しく作った driver の start() が upsertJobScheduler で登録し直すと、
再び発火する。⟹ rolling deploy や台数の縮小で1台を止めるときは、残りのプロセスのどれかを再起動する
(新しい driver で start() する)こと。
(【実測】redis-server 7.4.7・bullmq 6.3.8。同じ queueName・jobName の driver を2つ start() し、一方を stop() すると、getJobSchedulers() が空になり、
動いたままの他方の Worker は、その後3秒間 tick を1回も呼ばなかった。onTickError も鳴らない。新しい driver の start() で再び発火した。
⚠ rolling deploy で「新しいプロセスを start() してから古いプロセスを stop() する」順だと、古い方の stop() が新しい方の登録を消す——
上の実測と同じ形なので、新しいプロセスが動いていても tick は止まる。次に start() する driver が現れるまで、outbox に積まれた行は処理されないまま溜まる
(データは消えないが、embed・extract が止まったように見える)。止めるプロセスを stop() する代わりに、stop() を呼ばずにプロセスごと終わらせる
(Worker の lock が切れるまでは、その Worker が掴んだ最後のジョブが stalled になりうる)、または stop() の後に残るプロセスのどれかで新しい driver を start() し直すこと。)
🔴 **queueName か jobName は、テナント(ctx)ごとに分けること。**Worker はジョブの中身を見ずに、
自分に渡された ctx で runtime.tick(ctx, tick) を呼ぶ。テナントの違う driver が同じ queueName と
同じ jobName(既定は "mnemora-tick")を使うと、scheduler は1つに上書きされ(everyMs は最後に
start() した driver の値になる)、1回の発火はどれか1つの Worker、つまりどれか1つのテナントの tick に
しかならない。どのテナントが何回 tick されるかは決まらない。上の stop() も、全テナントの予定を止める。
(【実測】redis-server 7.4.7・bullmq 6.3.8。同じ queueName・jobName で everyMs: 100(テナントA)と everyMs: 1000(テナントB、後から start())を動かすと、scheduler は1つ
(every: 1000)になり、6秒間の7回の tick は A に4回・B に3回と、どちらの Worker が拾うかで振り分けられた。A は 100ms ごとには ticks されない。
jobName をテナントごとに分けると、200ms・3秒で A も B も15回ずつ tick され、scheduler は2つになった。)
同じ queueName に、everyMs を変えて start() し直しても同じである。start() は upsertJobScheduler で登録するので、後から start() した
driver の everyMs で共有の scheduler が置き換わり、先に動いていた driver の間隔も変わる(【実測】1000ms で動いていたものが、別の driver の start() で 200ms になり、
さらに 1000ms の driver の start() で 1000ms に戻った)。同じ driver の start() を重ねて呼んでも何も変わらない。everyMs を後から変える口は無いので、
変えたいときは新しい driver を作って start() する(古い driver の stop() は、上のとおり新しい登録を消すので、呼ぶ順に注意)。
詳しい API(CreateBullmqTickDriverOptions の各フィールド)は
src/tick-driver.ts の doc コメントを見ること。
⚠ 完了したジョブ・失敗したジョブは Redis に残り続ける
この driver は、Worker にもジョブにも removeOnComplete・removeOnFail を指定していない。BullMQ(6.3.8)は、
どちらも指定が無いとき、完了したジョブも失敗したジョブも全部残す(redis-queue-backend.js の getKeepJobs が
{ count: -1 } を返す)。⟹ everyMs ごとに1件ずつ、runtime.tick() の戻り値(TickResult)を持った完了ジョブが
Redis に溜まる(everyMs: 5_000 なら1日に 17,280 件)。runtime.tick() が throw した回は、失敗の理由と stack を
持った失敗ジョブとして残る。(【実測】redis-server 7.4.7・bullmq 6.3.8。everyMs: 50 で5秒走らせると、getJobCounts が完了41・失敗13(tick を4回に1回 throw させた)、
その queue のキーが66個、MEMORY USAGE の合計が約83.5KB(1ジョブあたり約1.5KB。戻り値が {} の最小の場合で、実際の TickResult や失敗ジョブの stack では
これより大きい)だった。removeOnComplete: { count: 5 } を付けた素の BullMQ の Worker では、同じ条件で完了は5件で頭打ちになった。
everyMs: 5_000 の1日 17,280 件を 1.5KB と置くと約 26MB。件数は上の読みどおりで、バイト数は戻り値の大きさ次第である。)
driver には保持の設定を渡す口が無い。Queue の側で掃除するには、同じ queueName の Queue を自分で作り、
BullMQ の queue.clean(grace, limit, type) を定期的に呼ぶ(grace ミリ秒より古いジョブを、type ごとに
最大 limit 件消す)。
import { Queue } from "bullmq";
const queue = new Queue("mnemora-tick", { connection: { host: "127.0.0.1", port: 6379 } });
// 1時間より古い完了ジョブと、1日より古い失敗ジョブを、それぞれ最大1000件消す。
await queue.clean(60 * 60 * 1000, 1000, "completed");
await queue.clean(24 * 60 * 60 * 1000, 1000, "failed");
await queue.close();(【実測】redis-server 7.4.7・bullmq 6.3.8。この形の queue.clean は、上の driver が溜めた完了42件・失敗13件を、grace: 0・limit: 0(無制限)で全部消した。limit を付けたときは
その件数までである。)保持の既定値を driver に入れるかどうかは決まっていない。
⚠ everyMs・jobName は構築時に検査する(queueName は BullMQ が検査する)
createBullmqTickDriver(...) は、concurrency(正の整数)に加えて everyMs・jobName を構築時に検査し、不正なら投げる(Queue・Worker は作らず、Redis にも繋がない)。
ADR 0477 が測ったとおり、検査しないと start() が成功したまま tick が黙って止まる入力があったため
(ADR 0498)。
| 入力 | 結果 |
|---|---|
| everyMs が数・有限・1 以上・Number.MAX_SAFE_INTEGER 以下 | 通る。小数(1.5)も通り、BullMQ が切り捨てた間隔(1 ms)で動く |
| everyMs が負・0・1 未満の小数・NaN・Infinity・MAX_SAFE_INTEGER 超(1e21 を含む)・数値の文字列("50")・null・undefined | 構築時に投げる |
| jobName を省略 | 既定 "mnemora-tick" |
| jobName が空でない文字列(: を含む・空白・日本語・300 文字も) | 通る |
| jobName が空文字・文字列でない | 構築時に投げる |
| queueName が空文字・: を含む | BullMQ が createBullmqTickDriver(...) の中で同期的に投げる(driver は検査しない) |
| queueName が空白・日本語・300 文字 | 動く |
(検査を足す前の測定【実測】redis-server 7.4.7・bullmq 6.3.8: 負の everyMs・1 未満の小数・1e21・空文字の jobName は start() が成功し、tick が数回(1e21・空文字は1回)で止まり onTickError も鳴らなかった。0・NaN・null は start() が reject、Infinity は Lua のエラーで reject した。数字は ADR 0477。)
⚠ 利用側は、不正な設定のまま既に動かしていたコードが、更新後は構築時に投げる。 移行は migration-v1 の 🔴 56。
⚠ エラーの通知先(onTickError)
onTickErrorを渡さないと、tick 自体の失敗は誰にも知らされない。runtime.tick()の throw(BullMQ の'failed')も、 Worker の'error'(接続エラーなど)も、driver はopts.onTickError?.(err)へ渡すだけで、 渡していなければ何も出さない(src/tick-driver.ts)。ログにも例外にもならず、 tick が動かないまま見た目は静かである。本番で使うなら渡すこと(ログに出す・メトリクスに積むなど)。import { createBullmqTickDriver } from "@mnemora/bullmq"; import type { CreateBullmqTickDriverOptions } from "@mnemora/bullmq"; declare const base: CreateBullmqTickDriverOptions; // 必須の項目(connection・queueName・runtime など) createBullmqTickDriver({ ...base, onTickError: (error) => console.error("mnemora tick failed", error), });- 🔴
onTickErrorを渡せば「失敗が全部分かる」わけではない。onTickErrorに届くのは、runtime.tick()(とonTickResult)の throw と、Worker・Queue の異常だけである。tick の中の個々のジョブ(outbox の行)が失敗しても、tick が throw しなければonTickErrorは鳴らない——その tick は成功として返り、失敗は戻り値のTickResultに数として載る。- 個々のジョブの失敗は
onTickResultで見る。TickResult.failedは、その tick でoutboxStore.fail()を呼んで リース競合で弾かれなかった件数、TickResult.unsupportedはfailedの内訳のうち「tickがその kind を処理する分岐を 持っていなかった」ジョブの配列である(unsupportedに入ったジョブもfailedに数える。 フィールドの定義はpackages/coreのTickResult)。import { createBullmqTickDriver } from "@mnemora/bullmq"; import type { CreateBullmqTickDriverOptions } from "@mnemora/bullmq"; declare const base: CreateBullmqTickDriverOptions; // 必須の項目(connection・queueName・runtime など) createBullmqTickDriver({ ...base, onTickResult: (result) => { if (result.failed > 0) { console.warn("mnemora tick: jobs failed", { failed: result.failed, unsupported: result.unsupported, }); } }, }); - 失敗した行そのものは、outbox の
last_error列(text)で見る。TickResultは件数とunsupportedの名指ししか持たないので、「どの行が・なぜ」は DB を引くこと。 - ⚠
failedの件数は「行が終端failedになった数」と常に一致するとは限らない(コミット後の接続断など。TickResult.failedの doc、Issue #836)。 onTickErrorの守備範囲は変えていない。failed > 0でonTickErrorを呼ぶ形や、TickResultの形を変える形は採っていない。
- 個々のジョブの失敗は
Queue側のエラーもonTickErrorに届く(onTickErrorを渡しているとき)。driver はWorkerの'error'・'failed'に加え、Queue(繰り返しジョブの登録に使う)の'error'にも listener を付け、 Redis 接続の失敗などをonTickErrorへ渡す。以前はQueueに listener が無く、bullmq がconsole.errorへ 固定で出すだけだった。- **
onTickErrorを渡さないときは、Queueに listener を付けない。**付けると bullmq(6.3.8 のQueueBase.emit: listener の無い'error'は EventEmitter が throw し、bullmq がそれを捕まえてconsole.errorへ出す)の 既定の出力が消え、Queueの異常が完全に黙るため。渡していなければ従来どおり標準エラーに出る (Worker側は従来から、渡していなければ黙る)。 - 同じ障害で、
onTickErrorが複数回呼ばれうる。QueueとWorkerは別々の Redis 接続を持ち、接続ごとに'error'を出す。Redis が落ちると両方が出す——同じ事象の重複ではなく別の接続の事象なので、driver は束ねない。 tick の失敗('failed')は、1回の失敗につき1回。通知のたびに通知先(アラートなど)が鳴るなら、 呼び出し側で間引くこと。 - 【実測】Redis が居ないポートを指した
QueueとWorkerは、それぞれECONNREFUSEDを emit した (bullmq 6.3.8、Redis 無しで確かめた)。Redis が在る状態でのQueueの失敗は測っていない (歯は fake のQueueが emit する形で縛っている)。 - 【実測】redis-server 7.4.7・bullmq 6.3.8。Redis が在る状態で、動いている driver の Redis を止めると、6秒間で
onTickErrorが30回(すべてECONNREFUSED。Queue と Worker の2接続で、再接続のたびに)届いた。 Redis を 永続化あり(appendonly yes)で再起動すると、tick は自動で再開した(5秒で25回)。永続化なしで再起動すると、scheduler が Redis ごと消えるので、tick は再開せず、onTickErrorも鳴らない(driver はstart()のときにしか登録しない)。キャッシュ用途の Redis(永続化なし)を使うなら、再起動のあとにstart()し直す仕組み (新しい driver を作り直す)が要る。 - 【実測】redis-server 7.4.7・bullmq 6.3.8。Redis が落ちている間の
start()は、reject せず、少なくとも15秒 pending のままだった(maxRetriesPerRequest: nullでも、指定しなくても。その間onTickErrorには ECONNREFUSED が届く)。Redis が戻ると resolve した。「失敗したら reject し、もう一度start()すればやり直す」(Issue #963)は、接続が拒まれる間は効かない。start()に自分でタイムアウトを掛けるなら、その後のstop()で後始末すること。 - 【実測】redis-server 7.4.7・bullmq 6.3.8。
connectionに ioredis のインスタンスを渡すとき、maxRetriesPerRequest: nullを指定していないインスタンス(ioredis の既定は 20)だと、createBullmqTickDriver(...)がBullMQ: Your redis options maxRetriesPerRequest must be null.で同期的に throw する(start()ではなく構築の時点)。nullを指定したインスタンスは動き、stop()の後もそのインスタンスは閉じられない(statusはreadyのまま。呼び出し側が閉じる)。オプションのオブジェクトを渡した場合は、BullMQ が警告を出して上書きする。 - 【実測】redis-server 7.4.7・bullmq 6.3.8。Redis が在る状態で、動いている driver の Redis を止めると、6秒間で
onTickErrorが30回(すべてECONNREFUSED。Queue と Worker の2接続で、再接続のたびに)届いた。 Redis を 永続化あり(appendonly yes)で再起動すると、tick は自動で再開した(5秒で25回)。永続化なしで再起動すると、scheduler が Redis ごと消えるので、tick は再開せず、onTickErrorも鳴らない(driver はstart()のときにしか登録しない)。キャッシュ用途の Redis(永続化なし)を使うなら、再起動のあとにstart()し直す仕組み (新しい driver を作り直す)が要る。 - 【実測】redis-server 7.4.7・bullmq 6.3.8。Redis が落ちている間の
start()は、reject せず、少なくとも15秒 pending のままだった(maxRetriesPerRequest: nullでも、指定しなくても。その間onTickErrorには ECONNREFUSED が届く)。Redis が戻ると resolve した。「失敗したら reject し、もう一度start()すればやり直す」(Issue #963)は、接続が拒まれる間は効かない。start()に自分でタイムアウトを掛けるなら、その後のstop()で後始末すること。 - 【実測】redis-server 7.4.7・bullmq 6.3.8。
connectionに ioredis のインスタンスを渡すとき、maxRetriesPerRequest: nullを指定していないインスタンス(ioredis の既定は 20)だと、createBullmqTickDriver(...)がBullMQ: Your redis options maxRetriesPerRequest must be null.で同期的に throw する(start()ではなく構築の時点)。nullを指定したインスタンスは動き、stop()の後もそのインスタンスは閉じられない(statusはreadyのまま。呼び出し側が閉じる)。オプションのオブジェクトを渡した場合は、BullMQ が警告を出して上書きする。
- **
⚠ lock の期限切れ(stalled)で、1回の tick に onTickResult と onTickError の両方が届きうる
🔴 【実測】redis-server 7.4.7・bullmq 6.3.8(ADR 0449)。 ソースの読み(dist/cjs/classes/worker.js の processJob・retryIfFailed、ADR 0440)のとおりだった。
2つの OS プロセスが同じ queueName で動き、一方の runtime.tick() がイベントループを45秒塞ぐ(lock を延長できない)と、もう一方の Worker が約60秒後(lock の期限30秒の後、次の stalled checker の周期)に
同じジョブの2本目の tick を走らせ、塞いでいた1本目は45.8秒で戻ってから onTickResult を呼び、直後に onTickError が2回(Missing lock for job ... moveToFinished、0.1秒以内)届いた。
最終のジョブは完了1件・失敗0件で、attemptsStarted: 2・stalledCounter: 1。(以下の箇条書きの「読み」は、この実測で裏づいた。lockDuration 既定30000ms の値は worker.js の読みのまま。)
- BullMQ は、Worker が処理中のジョブの lock を
lockDuration(既定 30000 ms)で持ち、その半分の間隔で延長する。この driver はlockDurationを設定しない(CreateBullmqTickDriverOptionsに口が無い。Workerには bullmq の既定値が渡る)。runtime.tick()がイベントループを長く塞ぐ・Redis との接続が途切れるなどで lock の延長が間に合わずに期限が切れると、stalled checker(既定stalledInterval30000 ms)がそのジョブを wait へ戻し、別の Worker が2本目の tick を走らせうる(同じジョブの再実行)。 - データは壊れない。 2本の tick が重なっても、outbox の行は
claimBatchの行ロック・リース・CAS(attempts)で二重に処理されない(「複数プロセスで動かすとき」と同じ守り)。driver も同じジョブを二重に数えない。 - ただし通知は素直ではない。 遅れて終わった1本目は、processor の中で
onTickResultを呼んだあと、BullMQ が完了を記録するmoveToCompletedをMissing lockで失敗させ、Worker の'error'経由でonTickErrorに届きうる(2回届く読み)。つまりその tick のonTickResultが届いたあとにonTickErrorが鳴ることがある。onTickErrorを「tick が動かなかった」の意味でだけ扱うと、この場合は誤る。 - 対処は書いていない(口を足す・driver で束ねる、のどちらも採っていない)。
onTickErrorのログには、同じ時刻のonTickResultがあるかを見ること。lockDurationを変える口は、公開 API の追加になるため足していない(ADR 0440)。
確かめていないこと
- BullMQ の Job Scheduler が実運用のワークロードでどの程度「重なる」かは測っていない。
- BullMQ 自身の可用性・再接続・Redis 障害時の挙動は、「エラーの通知先」の実測(Redis を止めて再起動、
start()が pending のまま)の範囲だけ測った。Redis Cluster・Sentinel・フェイルオーバーは測っていない。 - 複数マシン・ネットワーク越しの複数 OS プロセスからの同時 tick は測っていない (同一ホスト上の複数 OS プロセスまでは ADR 0325 の歯が測っている)。
(詳細は ADR 0325 の 「確かめていないこと」「引き受けた負債」を見ること)
