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:
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.
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.
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.
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:
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.
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.
AlertEvaluatormembacalastPricesebelum 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. OhlcAggregatormenghasilkan 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.

Memuat komentar...