リアルタイムな UI とは関係のない実験から始めます。パイプに 2 KB のデータが入っています。読み取り可能と通知されたので 1 KB だけ読み、もう一度問い合わせます。カーネルは、まだ読み取り可能だと言うでしょうか。

答えは、登録の仕方で変わります。そしてその答えは、イベントストリームをクライアントがどう消費すべきかに、ほぼそのまま当てはまります。

TL;DR

  • epoll のエッジトリガーは「変化」を1度だけ通知します。データを途中まで読んで待つと、永遠に待つことがあります。レベルトリガーは「状態」(まだ読める)を通知し、状態が解消するまで繰り返します。
  • イベントストリームにも同じ2つの設計があります。差分(「item 7 は X になった」)はエッジです。1つ落とすとクライアントは間違ったままで、それに気づく手段もありません。汚れフラグ(「item 7 は古いかもしれない、見に行って」)はレベルです。重複も順序の入れ替えも畳み込みも問題になりません。クライアントは必ず現在の値を取りに行くからです。
  • ただし、ヒントだけではまだレベルではありません。epoll が堅牢なのは、カーネルが状態を保持し続けるからです。アプリケーションで同じ役割をするのは、クライアントが読み直して比べられるもの(バージョンやカーソル)と、定期的な照合です。損失のあるチャネルのシミュレーションでは、ヒントだけだと300回中179回で収束し、照合を足すと300回すべてで収束しました。
  • どの実装も見落としやすい競合が1つあります。再取得の最中に届いたイベントです。答えは「もう一度再取得する」です。それを行う1行を消すと、3つのテストのうち2つが失敗します。
  • レベルトリガーは取得回数を食います。同じシミュレーションで、きちんと作った差分の消費側は1回の実行あたり約1回の取得で済み、無効化する側は約64回でした。取得のコストより、乱れた配送の下での正しさのほうが重要なとき、あるいは、修復の経路と通常の経路が同じコードであることに価値があるときに、レベルを選びます。

実験1:epoll は部分的な読み取りのあと何を返すか

epoll(7) のマニュアルは、まさにこの場面を説明しています。パイプの読み取り側を登録し、書き手が 2 kB 書き、epoll_wait が通知し、読み手が 1 kB 読み、もう一度 epoll_wait を呼びます。エッジトリガーのフラグ EPOLLET の場合、2回目の呼び出しは、入力バッファにデータが残っているにもかかわらず、「おそらくハングする」(will probably hang)と書かれています。約20行で再現できます(Linux のみ)。

epoll_demo.py
# Edge-triggered vs level-triggered epoll: what happens after a partial read?
# Linux only. Run: python3 epoll_demo.py
import os, select

def run(flags, label):
    r, w = os.pipe()
    os.set_blocking(r, False)
    ep = select.epoll()
    ep.register(r, select.EPOLLIN | flags)
    os.write(w, b"x" * 2048)                      # (2) the writer writes 2 kB
    first = ep.poll(timeout=1)                    # (3) the fd is reported ready
    os.read(r, 1024)                              # (4) the reader consumes only 1 kB
    second = ep.poll(timeout=0.5)                 # (5) wait again
    print(f"{label:15} first={bool(first)} second={bool(second)} (1024 bytes still unread)")
    if not second:
        # The documented fix for edge-triggered use: read until EAGAIN.
        try:
            while os.read(r, 1024):
                pass
        except BlockingIOError:
            print(f"{'':15} drained until EAGAIN")
    ep.close(); os.close(r); os.close(w)

run(0, "level-triggered")
run(select.EPOLLET, "edge-triggered")

Linux 6.12 での出力です。

level-triggered first=True second=True (1024 bytes still unread)
edge-triggered  first=True second=False (1024 bytes still unread)
                drained until EAGAIN

エッジトリガーの登録は、データが届いたときに1度だけパイプを通知しました。そのあと状態は変化していないので通知はなく、1024 バイトが残っているのに待ち続けます。マニュアルが示す対処は、ノンブロッキングのディスクリプターを使い、EAGAIN が返るまで読み切ることです。エッジトリガーの消費側は、変化を丸ごと消費する義務を負います。レベルトリガーの消費側にその義務はありません。少し取り、戻ってきて、もう一度通知を受けられます。

この区別は epoll より古く(ハードウェアの信号の「レベル」と「遷移」に由来します)、epoll は手で触って実感しやすい例です。エッジトリガーは、「何も取りこぼさない」という負担を消費側に押し付けます。レベルトリガーは、通知を「現在についての主張」にすることで、その負担をなくします。

イベントストリームでの同じ選択

「パイプ」を「画面に表示している項目の一覧」、「読み取り可能」を「古い」に読み替えます。

エッジ:差分を送るレベル:「汚れている」を送る
イベントの内容item:7 title = "B"item:7 が変わったかもしれない
消費側の動きペイロードを適用するitem:7 を再取得する
イベントの取りこぼしクライアントは間違うが、気づけない次のイベントか照合まで遅れるだけ
イベントの重複冪等でなければならない無料:2回汚すのは1回汚すのと同じ
順序の入れ替わり古い差分が新しい値を上書きしうる無料:再取得は現在の値を返す
100件のバースト100回の適用キーごとに1つの汚れに畳まれる
ペイロードデータを運ぶキーを運ぶ

Kubernetes のコントローラーも、同じ説明をされています。設計提案のアーカイブによれば、コントローラーは「耐障害性を最大にするために」レベルベースで作られ、「途中の状態更新がどれだけ失われても」、望ましい状態と観測された状態さえあれば正しく動作します。一方で、遅延を減らすために、通知型の watch API も併用します。通知は最適化で、真実は状態のほうにあります。PostgreSQL の NOTIFY のドキュメントも同じ方向を向いています。大量の情報を伝えたいなら、テーブルに置いて、レコードのキーを送るのがよい、と書かれています。

実験2:乱れたチャネルで、どの消費側が生き残るか

「差分は壊れやすい」という主張は、試す価値があります。下のシミュレーションでは、サーバーが5つのキーを持ち、200ティック(1ティックは仮想時間の 10 ms)の間、2ティックに約1回の頻度で更新します。5つの消費側が、同じ更新を、それぞれ別の乱れたチャネル経由で受け取ります。各メッセージは確率 p で失われ、確率 0.2 で重複し、0〜7ティック遅れます(つまり追い越しが起きます)。消費側の取得(fetch)は、リクエスト時点の値を返し、1〜8ティック後に届くので、取得の応答どうしも追い越しあいます。更新が止まったあとは、サーバーが 200 ms ごとに、各キーの現在のバージョンを載せた小さな「レベル」メッセージを送ります。このメッセージも同じ乱れたチャネルを通るので、失われることがあります。最後に、すべてのキーがサーバーと一致しているかで判定します。

  • A:差分を届いた順にそのまま適用する。
  • B:差分を、保持しているものより新しいバージョンのときだけ適用する(冪等で、順序に影響されない)。
  • C:ヒントだけ。イベントごとにキーを汚し、畳み込み付きの invalidator が再取得する。
  • D:C に、定期的なレベルメッセージを足す。サーバーのバージョンがクライアントのものより新しいキーを汚す。
  • E:B に、同じ定期メッセージを足す。遅れているキーを取得する。
sim.mjs:決定的。設定ごとにシード付きの300回の試行
// Which consumer design converges when the channel drops, duplicates and reorders messages?
// Deterministic: a seeded PRNG and a virtual clock (1 tick = 10 ms). Run: node sim.mjs
import { createInvalidator } from "./invalidator.mjs";

const KEYS = ["a", "b", "c", "d", "e"];
const rng = (a) => () => { // mulberry32: small, seedable, good enough for a simulation
  a = (a + 0x6d2b79f5) | 0;
  let t = Math.imul(a ^ (a >>> 15), 1 | a);
  t = (t + Math.imul(t ^ (t >>> 7), 61 | t)) ^ t;
  return ((t ^ (t >>> 14)) >>> 0) / 2 ** 32;
};

async function trial(seed, { pDrop, pDup, maxDelay }) {
  const rand = rng(seed), pick = (n) => Math.floor(rand() * n);
  let tick = 0, q = [];
  const at = (d, fn) => { const h = { t: tick + d, fn, dead: false }; q.push(h); return h; };
  const timers = { setTimeout: (fn, ms) => at(Math.max(1, Math.round(ms / 10)), fn), clearTimeout: (h) => (h.dead = true) };

  const server = Object.fromEntries(KEYS.map((k) => [k, 0]));       // key -> version (the value is the version)
  const mk = () => ({ val: Object.fromEntries(KEYS.map((k) => [k, 0])), fetches: 0 });
  const A = mk(), B = mk(), C = mk(), D = mk(), E = mk();            // five consumers, same updates, own channel each

  const refetch = (c) => (k) => new Promise((done) => {
    c.fetches++;
    const answer = server[k];                                       // the server answers with the value NOW
    at(1 + pick(maxDelay), () => { if (answer > c.val[k]) c.val[k] = answer; done(); }); // late answers may arrive out of order
  });
  const invC = createInvalidator({ refetch: refetch(C), timers });
  const invD = createInvalidator({ refetch: refetch(D), timers });
  const fetchE = refetch(E);

  const handlers = {
    A: (m) => m.type === "event" && (A.val[m.k] = m.v),                                    // apply payload in arrival order
    B: (m) => m.type === "event" && m.v > B.val[m.k] && (B.val[m.k] = m.v),                // apply payload if newer
    C: (m) => m.type === "event" && invC.invalidate(m.k),                                  // hint only
    D: (m) => (m.type === "event" ? invD.invalidate(m.k) : KEYS.forEach((k) => m.vers[k] > D.val[k] && invD.invalidate(k))),
    E: (m) => (m.type === "event" ? m.v > E.val[m.k] && (E.val[m.k] = m.v) : KEYS.forEach((k) => m.vers[k] > E.val[k] && fetchE(k))),
  };
  const send = (m) => { for (const h of Object.values(handlers)) {       // an independent lossy channel per consumer
    if (rand() < pDrop) continue;
    for (let i = 0; i < (rand() < pDup ? 2 : 1); i++) at(pick(maxDelay), () => h(m));
  } };

  const END = 200, STOP = 600;
  for (tick = 0; tick <= STOP; tick++) {
    if (tick < END && rand() < 0.5) { const k = KEYS[pick(KEYS.length)]; send({ type: "event", k, v: ++server[k] }); }
    if (tick % 20 === 0 && tick > END) send({ type: "heartbeat", vers: { ...server } }); // the "level": cheap, repeated, idempotent
    for (let moved = true; moved;) {
      moved = false;
      for (const h of q.filter((x) => x.t <= tick && !x.dead)) { h.dead = true; h.fn(); moved = true; }
      q = q.filter((x) => !x.dead);
      await new Promise(setImmediate);                               // let promise continuations run
    }
  }
  const ok = (c) => KEYS.every((k) => c.val[k] === server[k]);
  return { A, B, C, D, E, ok };
}

for (const cfg of [{ pDrop: 0, pDup: 0.2, maxDelay: 8 }, { pDrop: 0.2, pDup: 0.2, maxDelay: 8 }]) {
  const N = 300, tally = Object.fromEntries("ABCDE".split("").map((n) => [n, { ok: 0, fetches: 0 }]));
  for (let s = 1; s <= N; s++) {
    const r = await trial(s, cfg);
    for (const n of "ABCDE") { tally[n].ok += r.ok(r[n]) ? 1 : 0; tally[n].fetches += r[n].fetches; }
  }
  console.log(`\ndrop=${cfg.pDrop} dup=${cfg.pDup} reorder<=${cfg.maxDelay} ticks, ${N} trials`);
  const names = { A: "A delta, apply in arrival order", B: "B delta, apply if newer", C: "C hint only (invalidate+refetch)", D: "D hint + periodic level check", E: "E delta + periodic level check" };
  for (const n of "ABCDE") console.log(`${names[n].padEnd(36)} converged ${String(tally[n].ok).padStart(3)}/${N}  fetches/trial ${(tally[n].fetches / N).toFixed(1)}`);
}

出力です(シード 1〜300、Node.js 20.19。実行は決定的なので、同じ数字になります)。

drop=0 dup=0.2 reorder<=8 ticks, 300 trials
A delta, apply in arrival order      converged 203/300  fetches/trial 0.0
B delta, apply if newer              converged 300/300  fetches/trial 0.0
C hint only (invalidate+refetch)     converged 300/300  fetches/trial 75.2
D hint + periodic level check        converged 300/300  fetches/trial 75.2
E delta + periodic level check       converged 300/300  fetches/trial 0.0

drop=0.2 dup=0.2 reorder<=8 ticks, 300 trials
A delta, apply in arrival order      converged  68/300  fetches/trial 0.0
B delta, apply if newer              converged 101/300  fetches/trial 0.0
C hint only (invalidate+refetch)     converged 179/300  fetches/trial 64.2
D hint + periodic level check        converged 300/300  fetches/trial 64.6
E delta + periodic level check       converged 300/300  fetches/trial 1.0

目を引く点が4つあり、最後の1つは、ヒントを送る設計にとっていちばん都合が悪いものです。

  1. 重複と順序の入れ替えだけで、素朴な差分の消費側は壊れます(A:300回中203回)。バージョンの確認を入れれば完全に直ります(B)。
  2. 損失は、「そのキーの最後のイベント」に頼るものをすべて壊します。 損失20%のとき、B は101回、C は179回しか収束しませんでした。そのキーについての最後のメッセージなら、失われたヒントも、失われた差分と同じく取り返せません。
  3. 損失を直すのは、ヒントではなくレベルです。 D も E も毎回収束しました。C と D の差は、定期的な照合だけです。epoll でこの役をするのはカーネルの準備完了フラグですが、アプリケーションは自分で作る必要があります。クライアントがいつでもサーバーと比べられるバージョン、カーソル、ダイジェストです。
  4. 修復の経路を持つ、きちんと作った差分の消費側(E)のほうが安い。 1回の実行あたり約1回の取得で、無効化する側の約64回と比べて大きく少ないです。

つまり、レベルトリガーの設計を支持する理由は、「差分では動かない」ではありません。D では、修復の経路と通常の経路が同じコードであることです。E には2つの経路があります。常に動く適用の経路と、何かがおかしくなったあとにしか動かない取得の経路です。後者はいちばん実行されないコードで、初めて必要になったときに壊れている可能性が最も高いコードでもあります。取得が安いなら D を、高いなら E を選び、修復の経路を適用の経路と同じくらい念入りにテストしてください。

invalidator を作る

シミュレーションの消費側は、下のコードです。4つのことをしていて、それぞれが具体的な間違いに対応しています。

invalidator.mjs
// A level-triggered consumer: events only say "this key may be stale".
// The consumer owns the question "what is the current value?" and always answers it by refetching.
export function createInvalidator({
  refetch,                       // async (key) => void; must read the CURRENT value and store it
  windowMs = 50,                 // coalescing window: fixed, so latency is bounded
  maxBackoffMs = 5000,
  timers = { setTimeout, clearTimeout },
}) {
  const dirty = new Set();       // keys that may be stale and are not being fetched right now
  const fetching = new Set();    // keys with a refetch in flight
  const failures = new Map();    // key -> consecutive failures
  let timer = null;

  function schedule(ms = windowMs) {
    if (timer === null) timer = timers.setTimeout(flush, ms);
  }

  function flush() {
    timer = null;
    for (const key of [...dirty]) {
      if (fetching.has(key)) continue;           // stays dirty; re-armed when that fetch settles
      dirty.delete(key);
      fetching.add(key);
      Promise.resolve(refetch(key)).then(
        () => { failures.delete(key); settle(key, windowMs); },
        () => {
          const n = (failures.get(key) ?? 0) + 1;
          failures.set(key, n);
          dirty.add(key);                        // still stale
          // full jitter backoff so that many clients do not retry in lockstep
          settle(key, Math.random() * Math.min(maxBackoffMs, windowMs * 2 ** n));
        },
      );
    }
  }

  function settle(key, nextMs) {
    fetching.delete(key);
    if (dirty.has(key)) schedule(nextMs);        // an event arrived during the fetch, or the fetch failed
  }

  return {
    invalidate(key) { dirty.add(key); schedule(); },                          // any order, any number of times
    invalidateAll(keys) { for (const k of keys) dirty.add(k); schedule(); },  // after (re)connect
    pending: () => dirty.size + fetching.size,
  };
}

イベントをキーの集合に畳み込む。 invalidate(key) は dirty に加えて、タイマーを1つ設定します。2つのキーに対する100件のイベントは、2つのエントリーになります。

デバウンスではなく、固定の窓を使う。 タイマーは最初のイベントで1度設定し、後続のイベントでリセットしません。イベントのたびに再開する本物のデバウンスは、イベントが続く限り発火しません。流れ続けるストリームでは UI が飢餓状態になります。固定の窓なら、遅延の上限は windowMs です。

取得の最中に届いたイベントを扱う。 これが競合です。k の再取得が走っていて、サーバーをすでに読んだとします。そこへ k のイベントが届きます。そのイベントが告げているのは、取得が返す値より新しいものです。消費側が「いま取得中だから何もしない」とこれを捨てると、無関係なイベントが来るまで、画面は古い値のままです。そこで、取得中も k を dirty に残し、取得が終わったときに、もう一度 flush を予約します。それが if (dirty.has(key)) schedule(nextMs) の行です。消すとテストが落ちます。

invalidator.test.mjs
import assert from "node:assert/strict";
import { test } from "node:test";
import { createInvalidator } from "./invalidator.mjs";

const sleep = (ms) => new Promise((r) => setTimeout(r, ms));

test("100 events on 2 keys become 2 fetches", async () => {
  const fetched = [];
  const inv = createInvalidator({ refetch: async (k) => void fetched.push(k), windowMs: 20 });
  for (let i = 0; i < 100; i++) inv.invalidate(i % 2 ? "list" : "item:1");
  await sleep(80);
  assert.deepEqual(fetched.sort(), ["item:1", "list"]);
});

test("an event that arrives during a fetch causes a second fetch", async () => {
  let version = 1, shown = 0, calls = 0;
  const inv = createInvalidator({
    windowMs: 10,
    refetch: async () => {
      calls++;
      const v = version;            // the server answers with the value at request time...
      await sleep(40);              // ...and the response is slow
      shown = v;
    },
  });
  inv.invalidate("k");              // version 1 is announced
  await sleep(20);                  // fetch #1 is now in flight, holding version 1
  version = 2; inv.invalidate("k"); // version 2 is announced while fetch #1 is running
  await sleep(200);
  assert.equal(calls, 2);
  assert.equal(shown, 2);           // without the re-arm, the UI would stay at 1
});

test("a failed fetch is retried with backoff instead of being forgotten", async () => {
  let calls = 0;
  const inv = createInvalidator({
    windowMs: 5,
    maxBackoffMs: 20,
    refetch: async () => { if (++calls < 3) throw new Error("boom"); },
  });
  inv.invalidate("k");
  await sleep(300);
  assert.equal(calls, 3);            // two failures, then success, then it stops
  assert.equal(inv.pending(), 0);
});
ok 1 - 100 events on 2 keys become 2 fetches
ok 2 - an event that arrives during a fetch causes a second fetch
ok 3 - a failed fetch is retried with backoff instead of being forgotten

この行をコメントアウトすると、2番目と3番目のテストが失敗します(確認しました)。通るのは1番目だけです。

失敗はジッターつきのバックオフで再試行し、キーごとに同時に取得するのは1つだけにする。 取得に失敗したキーは dirty に戻ります。待ち時間はランダムで(フルジッター。バックオフとジッターの解説を参照)、同時に失敗した多数のクライアントが同時に再試行しないようにします。また、同じキーを同時に2回取得しないので、同じキーの応答が追い越しあうことはありません。

この40行が、あえてやっていないことが2つあります。

  • 再接続。 再接続したあとは、何を取りこぼしたか分からないので、表示しているキーについて invalidateAll(keys) を呼びます。
  • 定期的なレベルの照合。 シミュレーションでは、サーバーからのメッセージでした。実際のシステムでは、レスポンスヘッダーのバージョン、ETag、カーソル、軽いダイジェスト用のエンドポイントなどになるでしょう。仕組みより形が大切です。つまり、クライアントが比べられるもので、イベントが失われても残り続けるものです。

トレードオフ

  • 取得の増幅。 クライアントがイベントから新しい値を計算できる場合でも、汚れたキーは取得を1回消費します。シミュレーションでは、修復の経路を持つ差分の消費側が1回の実行で約1回取得したのに対し、無効化する消費側は約64〜75回取得しました。
  • サンダリングハード(一斉殺到)。 サーバーが1つのヒントを多数のクライアントに配ると、全員が同時に再取得します。再試行だけでなく最初の窓にもジッターを入れるか、サーバー側でヒントを散らします。
  • バージョンを運ぶヒント。 中間の安い道があります。{key, version} を送り、クライアントがそのバージョン以上をすでに持っているなら取得を省くのです。自分で作った重複を減らせます。
  • ベストエフォートなチャネル上のヒント。 プッシュ通知のようにメッセージが落ちたり畳まれたりしうるチャネルは、ヒントを運ぶのには向いています。ただし、レベルを運ぶ唯一の手段にはしないでください。

壊れ方

失敗症状防御
再取得中に届いたイベントを無視する無関係なイベントが来るまで画面が古いままキーを汚れたままにして、もう一度再取得する(再設定の行)。
イベントのたびに再開するデバウンス流れ続けるストリームの間、UI が更新されない最初のイベントからの固定の窓。
あるキーの最後のイベントが失われる次の変更まで古いままバージョンの定期的な比較、またはフォーカス時・再接続時の更新。
1つのキーの応答が順序を入れ替わる古い値が新しい値を置き換えるキーごとに同時に1回だけ取得する。または書き込み時にバージョンで守る。
失敗した取得を忘れる永続的に古いままキーを dirty に戻し、ジッターつきのバックオフで再試行する。
再同期なしの再接続すべてが古い可能性がある表示中のキーに invalidateAll。
1回のブロードキャストで多数のクライアントが再取得オリジンの負荷スパイク最初の窓にジッター。バージョンつきのヒント。

使う場面と使わない場面

使うのは、クライアントが表示するものについて、「現在の値を読む」操作が安いとき、配送が乱れやすいとき(再接続、不安定な回線、中継を挟んだファンアウト)、そして修復の経路を通常の経路と同じにしたいときです。

次のような場合は向きません。

  • イベントそのものがデータであるストリーム:監査ログ、台帳、各ステップを順番どおりに見る必要があるもの。これらに必要なのは無効化ではなく、カーソルつきのログです(再開できるストリームを参照)。
  • 大きなオブジェクトの小さな変更:イベントのたびに値全体を再取得するのは無駄です。バージョンつきの差分を使い、フォールバックにヒント方式の再同期を置く形(消費側 E)を検討します。
  • 通知そのものが内容であるもの:何かが起きたと知らせるトーストのようなもの。ヒントを読み直しても同じ効果にはなりません。

実験を動かす

# 実験1(Linux のみ)
python3 epoll_demo.py

# 実験2とテスト(Node.js 18 以降。実行したのは Node.js 20)
node --test invalidator.test.mjs
node sim.mjs

そして、わざと壊してみてください。if (dirty.has(key)) schedule(nextMs) の行を消してテストをやり直します。sim.mjs の pDrop を 0.5 にして、どの消費側が生き残るかを見ます。シミュレーションは、更新が止まったあとにしかレベルメッセージを送りません(tick > END)。この条件を外して更新中にも送り、取得回数がどう変わるかを見ます。

事実か、ヒントか

イベントストリームは、世界についての事実を運ぶこともできますし、世界がいまは違うというヒントを運ぶこともできます。事実のほうは、すべてが1回ずつ順番に届くなら安上がりです。ヒントのほうは、配送が乱れても動き続けます。ヒントは間違いようがなく、遅れるだけだからです。どちらの場合でも、クライアントが自分の状態とサーバーの状態を比べる手段を持たせてください。ストリームは、真実のすべてではないからです。

確認したことと、モデルにすぎないこと

実験1は Linux 6.12 と Python 3.13 で実行しました。実験2とテストは Node.js 20.19 で実行しました。シミュレーションはモデルです。チャネル(独立した損失・重複・遅延)、ティックの長さ、パラメーターは僕が選んだもので、どのネットワークの計測値でもなく、取得回数はそれらに依存します。実際のブラウザや実際のネットワークでは試しておらず、invalidator のベンチマークも取っていません。epoll の挙動はドキュメントに書かれているもので、この1つのカーネルで再現しました。