日本語版
最新ニュース
科学&テクノロジー

Postgres テーブルの下で非同期タスクをスイープする方法と理由

私は、各エンドポイントが非常に愚かな DB クエリをラップする、スリムで愚かなサーバーが好きです。 ダムクエリは高速です。高速なクエリにより、Web サイトがスムーズかつ軽快になります。これらのクリック/レンダリング ループは神聖なままにしておきます。 複雑さを一掃する task テーブル: router.post("/signup", async ctx => { const { email, password } = await ctx.request.body().value; const [{ usr_id } = { usr_id: null }] = await sql` with usr_ as (…

Postgres テーブルの下で非同期タスクをスイープする方法と理由

1763752568
2025-11-21 18:28:00

私は、各エンドポイントが非常に愚かな DB クエリをラップする、スリムで愚かなサーバーが好きです。

ダムクエリは高速です。高速なクエリにより、Web サイトがスムーズかつ軽快になります。これらのクリック/レンダリング ループは神聖なままにしておきます。

複雑さを一掃する task テーブル:

router.post("/signup", async ctx => {
  const { email, password } = await ctx.request.body().value;
  const [{ usr_id } = { usr_id: null }] = await sql`
    with usr_ as (
      insert into usr (email, password)
      values (${email}, crypt(${password}, gen_salt('bf')))
      returning *
    ), task_ as (
      insert into task (task_type, params)
      values ('SEND_EMAIL_WELCOME', ${sql({ usr_id })})
    )
    select * from usr_
  `;
  await ctx.cookies.set("usr_id", usr_id);
  ctx.response.status = 204;
});

もちろん使って mailgun.send キューに並べるよりも簡単です task テーブル。間接的な追加によってシステムが作成されることはほとんどありません 少ない 複雑な。しかしどういうわけか、私はまさにそれを主張するためにここにいます。私のマニフェストを無視してもいいし、
私の実装にスキップしてください 最後に。

秘密表面エラー領域

顧客は宇宙線など気にしていません。彼らは何かを望んでいます。さらに重要なことは、彼らが望んでいることは、 即時確認 彼らのこと。彼らは目標に向けた精神的な負担を軽減したいと考えています。

その責任を委任するために、おそらく重要なのは DB だけです。情報がデータベースにコミットされたら、自信を持って「ここから情報を取得します」と言うことができます。

後でメールを送信することもできます。後で支払いを処理することもできます。後からほとんど何でもできます。顧客に、今のような一日を続けてもよいと伝えてください。

明確なフィードバックで顧客を喜ばせましょう。

一度に 1 つの場所に書き込むことで、コンピューターを快適にします。

独自の 2 フェーズ コミットを決してハンドロールしないでください

2 つの場所に「同時に」書き込むことは罪です。

神々が私たちにコンピューターストレージを与えたとき、人々は不幸になりました。彼らは叫びました、「一貫性とは何ですか?私たちの保証はどこにありますか?なぜ私がしなければならないのですか?」 fsync?」そして彼らはコーディング洞窟で何年もの間、袋と灰を着ていました。

神々が Postgres (およびその他の粗悪なデータベース) を石板に走り書きしたとき、人々は大喜びしました。神聖な「データベース トランザクション」により、人類は複数の場所を同時に読み書きできるように見せることができました。

今日に至るまで、データベースは 時々仕事をする

しかし、開発者の中には神の働きを否定する人もいます。彼らは複数のツールを混合するため、複数の場所に書き込むという罪を犯します。

「ああ、行を挿入した後は pubsub メッセージを送信するだけです。」しかし、データは失われます。行を挿入する前のメッセージ?データが失われました。すべての冒涜者は 2 フェーズ コミットを再発明する運命にあります。

物事を行うための 1 つの方法

レゴが好きです。プレイドーが好きです。リンカーン・ログが好きです。しかし、私はそれらを混ぜ合わせるのが好きではありません。

状態が SQS、Redis、PubSub、Celery、Airflow などに分散している場合、システムを調査するのは大変です。プロセスが期待どおりに実行されない理由を調べるために、地元の探偵事務所を開く必要はありません。

最新のプロジェクトのほとんどは SQL を使用します。私はシステムが混在するのが嫌いなので、可能な限り SQL を使用するようにしています。

すべての SQL データベースの中で、Postgres は現在、最新の一流の機能とサードパーティの拡張機能の最適な組み合わせを提供しています。 Postgres は、Kafka の模造品、人工 Airflow、くだらない Clickhouse、厄介な Elasticsearch、貧乏人の PubSub、売り出し中の Celery などになります。

確かに、Postgres には、それぞれの特殊なシステムの優れた機能がすべて備わっているわけではありません。ただし、キュー/パイプライン/非同期データをメイン データベースに同じ場所に配置すると、一連のエラーが排除されます。私の経験では、取引保証は他のすべてに優先します。

TODO主導の開発

while (true) {
  // const rows = await ...
  for (const { task_type, params } of rows)
    if (task_type in tasks) {
      await tasks[task_type](tx, params);
    } else {
      console.error(`Task type not implemented: ${task_type}`);
    }
}

シンプルな再試行システムを備えた非同期デカップリングは、すべての不完全なフローを魔法のように追跡します。

Jira に依存する必要はありません。バグや未実装のタスクはログに記録され、再試行されます。エラー キューから再帰的に作業するのは、本当に素晴らしい経験です。すべてのライブ/緊急 TODO は同じ場所 (開発中と運用中) に出力されます。

このパラダイムでは、スケーラブルなパイプラインに引き寄せられるでしょう。
希望的観測 自然な建築を作ります。

ヒューマンフォールトトレランス

多くのシステムは、無駄な再試行ループを人間に押し付けています。

人間はヒューマンエラーに対するフィードバックを受け取る必要があります。しかし、コンピュータ (およびそのソフトウェア開発者) が処理できる問題については、人間がフィードバックを受け取るべきではありません。

すべての再試行ループはどこかで発生する必要があることに注意してください。顧客や開発者に何を委任するかには注意してください。ビジネスの収益は人間の忍耐力によって左右されます。コンピュータは人間よりも無限に忍耐力があります。

コードを見せて

こちらが task テーブル:

create table task
( task_id bigint primary key not null generated always as identity
, task_type text not null -- consider using enum
, params jsonb not null -- hstore also viable
, created_at timestamptz not null default now()
, unique (task_type, params) -- optional, for pseudo-idempotency
)

タスクワーカーのコードは次のとおりです。

const tasks = {
  SEND_EMAIL_WELCOME: async (tx, params) => {
    const { email } = params;
    if (!email) throw new Error(`Bad params ${JSON.stringify(params)}.`);
    await sendEmail({ email, body: "WELCOME" });
  },
};

(async () => {
  while (true) {
    try {
      while (true) {
        await sql.begin(async (tx: any) => {
          const rows = await tx`
            delete from task
            where task_id in
            ( select task_id
              from task
              order by random() -- use tablesample for better performance
              for update
              skip locked
              limit 1
            )
            returning task_id, task_type, params::jsonb as params
          `;
          for (const { task_type, params } of rows)
            if (task_type in tasks) {
              await tasks[task_type](tx, params);
            } else {
              throw new Error(`Task type not implemented: ${task_type}`);
            }
          if (rows.length 

このスニペットの注目すべき機能をいくつか示します。

  • タスク行は、 ない 場合は削除されます sendEmail 失敗します。 PG トランザクションはロールバックされます。列と sendEmail 再試行されます。
  • PGトランザクション tx タスクに渡されます。これは、行を「処理済み」としてマークする場合などに便利です。
  • トランザクションにより、エラー処理が大幅に改善されます。不可逆的な副作用が発生する前に、必ず可逆的なクエリを編成してください (電子メールを送信する前に DB ステータスをマークするなど)。 DB は最後にコミットすることに注意してください。
  • のため skip locked、これらのワーカーを任意の数で並行して実行できます。彼らはお互いのつま先を踏みつけません。
  • ランダムな順序付けは技術的にはオプションですが、これによりシステムのエラーに対する耐性が高まります。適切なランダム性があれば、単一のタスク タイプがすべてのタスクのキューをブロックすることはできません。
  • 使用 order by (case task_type ... end), random() 簡単に優先順位付きキューを作成できます。
  • 再試行回数を制限するとコードはより複雑になりますが、電子メールなどのユーザーが直面する副作用を考慮すると、間違いなく価値があります。
  • if (rows.length prevents overzealous polling. Your DBA will be
    grateful.

#Postgres #テーブルの下で非同期タスクをスイープする方法と理由

執筆者について: nipponese

Nipponese News編集部は、国内外のニュースを日本語で分かりやすくお届けします。