Kembali ke jurnalCATATAN FAJAR
Software Engineering5 menit baca

Stateful Streams: Trigger Edge vs Level dan Window OHLC

Menyimpan harga sebelumnya untuk alert edge dan mengelompokkan tick ke dalam time window untuk candle OHLC.

BAGIAN 13 DARI 16Kafka vs Redis Streams
  1. 01Kafka vs Redis Streams: Bikin Stock Ticker Dua Kali
  2. 02Kafka Partition dan Key: Menjaga Simbol Tetap Berurutan
  3. 03Dasar-Dasar Redis Streams: XADD, Entry ID, dan MAXLEN
  4. 04Kafka Consumer Group: Lag, Draining, dan Rebalancing
  5. 05Consumer Group Redis Streams: Competing Consumer dan PEL
  6. 06At-Least-Once di Kafka: Retry dan Dead Letter Queue
  7. 07Redis Streams: Message Nyangkut, XAUTOCLAIM, dan Dead Letter
  8. 08Exactly-Once di Kafka: Idempotent Producer dan Transaction
  9. 09Redis Pub/Sub vs Streams: Ephemeral atau Durable
  10. 10Co-Partitioning: Key Sama, Partition Sama, di Kedua Topic
  11. 11Menulis Custom Partition Assigner KafkaJS untuk Local Join
  12. 12Kafka ISR vs Redis Cluster: Replication dan Failover
  13. 13Stateful Streams: Trigger Edge vs Level dan Window OHLCKamu di sini
  14. 14Membangun Live Market Dashboard dengan SSE dan Next.js
  15. 15Benchmark Kafka vs Redis Streams: Throughput dan Latency
  16. 16Kafka atau Redis Streams: Gimana Cara Memilihnya
Di artikel ini 7 bagian

Price alert yang mendeteksi persilangan threshold butuh harga sebelumnya. Tes yang cuma melihat harga saat ini bisa mengirim notifikasi di setiap tick yang berada di atas threshold. State yang disimpan menentukan perilaku mana yang diterima user.

Selama perbandingan ticker, saya awalnya mempelajari transport, consumer group, recovery, dan replication. Saya kira logic alert-nya bakal lebih gampang. Logic itu butuh keputusan tersendiri soal state dan ordering.

Transport versus pemrosesan

Stream processing bisa stateless atau stateful. Handler yang stateless bereaksi terhadap record saat ini. Handler yang stateful juga memakai record sebelumnya atau nilai yang disimpan. Deteksi persilangan threshold butuh state sebelumnya itu.

Library Kafka Streams jalan di JVM. Untuk aplikasi Node.js saya, saya mengimplementasikan alert evaluator dan aggregator untuk window candle. State dan perilaku recovery keduanya tetap jadi tanggung jawab aplikasi.

Alert yang cuma terpicu sekali: edge vs level

Ambil contoh rule seperti "notify Alice when BBCA rises to or above 9500." Rule itu bisa dibaca dengan dua cara:

  • Level: terpicu di setiap tick yang harganya >= 9500. Pengecekan stateless ini mengirim notifikasi lagi untuk setiap tick yang memenuhi kondisinya.
  • Edge: terpicu cuma di tick saat harga naik dari bawah 9500 jadi 9500 atau lebih. Satu notifikasi per persilangan, dan itulah yang diinginkan user dari "notify me when."

Tes level cuma butuh tick saat ini:

typescript
export function matches(op: Operator, price: number, threshold: number): boolean {
  switch (op) {
    case "<": return price < threshold;
    case "<=": return price <= threshold;
    case "=": return price === threshold;
    case ">=": return price >= threshold;
    case ">": return price > threshold;
  }
}

Tes edge butuh satu hal lagi: harga sebelumnya.

typescript
export function crossed(
  op: Operator,
  prevPrice: number,
  price: number,
  threshold: number,
): boolean {
  return !matches(op, prevPrice, threshold) && matches(op, price, threshold);
}

Persilangan berarti kondisinya bernilai false di tick sebelumnya dan true di tick saat ini. Satu tick aja nggak bisa membuktikan perubahan itu. Di ticker saya, AlertEvaluator menyimpan harga sebelumnya untuk setiap simbol.

typescript
export class AlertEvaluator {
  private readonly lastPrice = new Map<StockSymbol, number>();
  private readonly rulesBySymbol = new Map<StockSymbol, AlertRule[]>();

  process(tick: Tick): Notification[] {
    const prev = this.lastPrice.get(tick.symbol);
    this.lastPrice.set(tick.symbol, tick.price); // update state AFTER reading prev

    const rules = this.rulesBySymbol.get(tick.symbol);
    if (!rules) return [];

    const out: Notification[] = [];
    for (const rule of rules) {
      const fired =
        rule.trigger === "edge"
          ? prev !== undefined && crossed(rule.operator, prev, tick.price, rule.threshold)
          : matches(rule.operator, tick.price, rule.threshold);
      if (fired) out.push(this.toNotification(rule, tick));
    }
    return out;
  }
}

Evaluator membaca harga sebelumnya, lalu meng-update lastPrice sebelum mengevaluasi rule-nya. Dengan begitu, tick berikutnya dibandingkan dengan tick saat ini. Untuk tick pertama sebuah simbol, prev bernilai undefined. Rule edge nggak bisa terpicu karena nggak ada harga sebelumnya untuk membuktikan persilangan.

Saya me-replay 200 tick sintetis yang sama dengan threshold yang sama di mode edge dan level. Mode edge menghasilkan 9 notifikasi. Mode level menghasilkan 194, kira-kira 21 kali lipatnya. Hasil ini berlaku untuk urutan input itu. Perilaku alert yang dibutuhkan yang menentukan mode mana yang benar.

mermaid
sequenceDiagram
    participant T as Tick Stream
    participant E as AlertEvaluator "state: lastPrice"
    T->>E: price=9480 (below 9500)
    Note over E: lastPrice=9480, no fire
    T->>E: price=9505 (crosses 9500)
    Note over E: EDGE fires (was below, now above)
    Note over E: LEVEL fires too
    T->>E: price=9510 (still above 9500)
    Note over E: EDGE stays silent (no new crossing)
    Note over E: LEVEL fires again

Kepemilikan state di berbagai partition

Setiap proses menyimpan map lastPrice-nya di dalam sebuah instance AlertEvaluator. Di deployment yang dipartisi, setiap pod punya evaluator untuk partition yang di-assign ke pod itu. Map tersebut menyimpan state untuk simbol-simbol di partition itu.

Topic market-updates dan notifications memakai key symbol yang sama dan enam partition, yang menghasilkan nomor partition yang cocok. Assigner yang kompatibel juga harus menaruh partition yang cocok di pod yang sama. Ini mendukung akses state lokal. Tapi ini nggak memulihkan state setelah startup, crash, atau rebalance.

Awalnya saya menganggap co partitioning sebagai cara untuk bikin join lebih murah. Co partitioning juga memberi state pemilik yang jelas. Kalau pod lain butuh state itu, aplikasinya harus mentransfer, menyalin, atau membangunnya ulang.

Mengelompokkan tick ke dalam tumbling window

Alert mengevaluasi tick satu per satu. Candle di chart merangkum sebuah interval dengan empat harga: open, high, low, dan close. Untuk membangun candle dari tick, aplikasi mengelompokkan tick ke dalam time window.

Tumbling window adalah bucket waktu berukuran tetap yang nggak saling tumpang tindih, misalnya setiap 1000 milidetik. Setiap tick yang datang dimasukkan ke bucket tempat timestamp-nya berada. Kalau timestamp sebuah tick masuk ke bucket yang lebih belakang dari bucket yang sedang terbuka, bucket yang terbuka itu selesai. Bucket itu ditutup, di-emit, dan bucket baru dimulai.

OhlcAggregator mengelompokkan tick berdasarkan simbol dan menghitung windowStart dari setiap timestamp:

typescript
export class OhlcAggregator {
  private readonly current = new Map<StockSymbol, Candle>();

  constructor(private readonly windowMs: number) {}

  add(tick: Tick): Candle | null {
    const windowStart = Math.floor(tick.ts / this.windowMs) * this.windowMs;
    const cur = this.current.get(tick.symbol);

    if (!cur || cur.windowStart !== windowStart) {
      const closed = cur && cur.windowStart !== windowStart ? cur : null;
      this.current.set(tick.symbol, {
        symbol: tick.symbol,
        windowStart,
        open: tick.price,
        high: tick.price,
        low: tick.price,
        close: tick.price,
        volume: tick.volume,
        count: 1,
      });
      return closed;
    }

    cur.high = Math.max(cur.high, tick.price);
    cur.low = Math.min(cur.low, tick.price);
    cur.close = tick.price;
    cur.volume += tick.volume;
    cur.count++;
    return null;
  }

  flush(): Candle[] {
    const out = [...this.current.values()];
    this.current.clear();
    return out;
  }
}

Ekspresi Math.floor(tick.ts / windowMs) * windowMs menghitung awal window. Di implementasi ini, tick di window yang lebih belakang menutup window saat ini. Kalau nggak ada tick berikutnya yang datang, window-nya tetap terbuka sampai flush() dijalankan.

Desain ini berasumsi timestamp datang secara berurutan untuk setiap simbol. Implementasi untuk production butuh policy yang eksplisit untuk record yang telat dan pemulihan state.

mermaid
graph LR
    subgraph "window 09:00:00"
        A["tick O=9500"] --> B["tick H=9520"] --> C["tick L=9490"]
    end
    subgraph "window 09:00:01"
        D["tick opens new candle"]
    end
    C -->|"tick.ts crosses boundary"| D
    C -.->|closed candle emitted| E["Candle: O/H/L/C, volume, count"]

Seperti apa 200 tick kalau jadi candle

Tesnya memakai 200 tick dari dua simbol, dengan timestamp sintetis yang berjarak 10ms. Window 50ms menghasilkan 80 candle yang udah ditutup. Setiap candle berisi open, high, low, close, dan volume untuk interval-nya. Ini mengurangi jumlah nilai yang harus ditampilkan chart.

Saya mengimplementasikan tumbling window, yang membagi waktu jadi interval berurutan yang nggak saling tumpang tindih. Ini cocok dengan interval candle di chart saya.

Hopping window bergeser kurang dari lebarnya, jadi satu tick bisa masuk ke beberapa window. Sliding window meng-update batasnya seiring event berdatangan. Session window mengelompokkan aktivitas dan ditutup setelah jeda tanpa aktivitas.

State itu kamu yang kelola

Kafka Streams bisa memulihkan state store yang dikonfigurasi dari changelog topic. Implementasi in memory saya nggak punya pemulihan otomatis. Setelah crash, implementasi itu harus membangun ulang state dari source record yang di-retain atau memulihkan snapshot yang durable.

Kafka transaction bisa mengoordinasikan output Kafka dan offset, tapi nggak memulihkan map lokal dengan sendirinya. Desain recovery juga harus menyimpan dan memulihkan state di posisi yang konsisten dengan offset tersebut.

Redis juga bisa mendukung pola-pola ini. Lua script bisa membaca dan meng-update state yang disimpan secara atomic untuk pengecekan edge. Sorted set bisa mendukung aggregate per window. Contoh-contoh Node ini tetap butuh kode aplikasi untuk mengelola window, retention, dan recovery.

State dan window dalam satu gambaran

  • Trigger edge membandingkan nilai saat ini dengan nilai sebelumnya. Di tes 200 tick yang sama, mode edge menghasilkan 9 notifikasi dan mode level menghasilkan 194.
  • AlertEvaluator membaca lastPrice sebelum meng-update nilai itu untuk simbolnya.
  • Di implementasi ini, window yang lebih belakang menutup window saat ini. flush() terakhir meng-emit window yang masih terbuka.
  • OhlcAggregator menghasilkan 80 candle dari 200 tick dengan window 50ms.
  • Partition topic dan consumer assignment yang cocok mendukung akses state lokal. Keduanya nggak menyediakan pemulihan state.
  • Recovery harus memulihkan state dan mengoordinasikannya dengan source offset dan output.

Awalnya saya memverifikasi perilaku ini lewat log. Lalu saya membangun dashboard SSE live untuk memeriksa tick, alert, dan kepemilikan partition di browser.

TOPIK

SELANJUTNYA DI SERI INIMembangun Live Market Dashboard dengan SSE dan Next.js

MAKASIH UDAH BACA

Gimana menurutmu?

Reaksi atau obrolan, dua-duanya selalu ditunggu.

Memuat reaksi…

Bagikan

Memuat komentar...

LANJUT JELAJAH

Satu pikiran bawa ke pikiran lain.

Semua tulisan
Kembali ke semua tulisanSatu catatan, pelan-pelan.