Saya pengin mengamati ticker-nya di browser setelah beberapa minggu mengetes lewat log dan CLI tool. Dashboard-nya memakai broker Kafka dan Redis yang sama dari eksperimen kegagalan. Dashboard ini menampilkan harga terkini, alert, consumer assignment, dan lag.
Engine bersama mengumpulkan state broker. Route Server Sent Events mengirim snapshot ke sebuah React component. Konfigurasi Next.js menjaga client broker tetap di server.
Satu engine per proses
Route yang membuat client broker untuk setiap request bisa meninggalkan koneksi duplikat setelah reload waktu development. Tab browser yang terpisah juga bisa membuat consumer yang nggak perlu. Group Kafka yang independen bakal menerima history topic masing-masing. Sebaliknya, consumer di satu group bersama bakal memicu perubahan assignment.
Perbaikannya adalah singleton biasa yang disimpan di globalThis, supaya tetap bertahan melewati HMR maupun banyak request di proses yang sama:
// Process-global singleton (survives Next dev HMR + shared across requests/tabs).
const globalForEngine = globalThis as unknown as { __liveEngine?: LiveEngine };
export function getLiveEngine(): LiveEngine {
if (!globalForEngine.__liveEngine) globalForEngine.__liveEngine = new LiveEngine();
return globalForEngine.__liveEngine;
}
Setiap route memanggil getLiveEngine() untuk mengambil instance bersama itu. Method ensureStarted() menjaga proses startup, yang membuat topic, producer, dan consumer. Request berikutnya memakai ulang state itu:
async ensureStarted(): Promise<void> {
if (this.started) return;
this.started = true;
try {
await this.start();
} catch (e) {
this.lastError = msg(e);
this.started = false; // let a later request retry
}
}
Satu engine, satu producer Kafka, satu koneksi Redis, dua consumer, nggak peduli berapa banyak tab yang lagi menonton.
Apa yang disambungkan oleh engine
Engine mem-produce tick sintetis ke kedua broker. Dua consumer pod mengevaluasi alert dan men-join-nya dengan state harga lokal. Lalu engine mengirim snapshot state ke setiap subscriber.
graph LR
SIM["Market simulator"] --> KP["Kafka producer<br/>market-updates"]
SIM --> RX["Redis XADD<br/>stream:market-updates"]
KP --> KB[("Kafka broker")]
KB --> C1["Pod 1 consumer<br/>co-partition assigner"]
KB --> C2["Pod 2 consumer<br/>co-partition assigner"]
C1 --> ENG["Live engine<br/>snapshot state"]
C2 --> ENG
ENG -->|"broadcast every 400ms"| SSE["SSE route handler<br/>ReadableStream"]
SSE -->|"EventSource"| UI["LiveDashboard"]
Dual write-nya terjadi setiap interval 150ms, tiga tick sekaligus, di-batch jadi satu send Kafka dan satu pipeline Redis:
private async produceTick(): Promise<void> {
if (!this.running || !this.producer) return;
const ticks = [this.sim.next(), this.sim.next(), this.sim.next()];
await this.producer.send({ topic: TOPICS.marketUpdates, messages: ticks.map(tickToKafka) });
this.produced += ticks.length;
const pipe = this.redis.pipeline();
for (const t of ticks) {
pipe.xadd(REDIS_KEYS.marketUpdates, "MAXLEN", "~", "20000", "*", ...tickToRedis(t));
}
await pipe.exec();
this.redisAdded += ticks.length;
}
Di sisi Redis, stream-nya dibatasi sekitar 20.000 entry (MAXLEN ~ 20000) supaya nggak tumbuh tanpa batas selama dashboard dibiarkan jalan. Kafka nggak butuh itu di sini, karena retention-nya berbasis waktu di level topic dan bukan urusan yang harus dikelola loop ini.
Timer setInterval yang terpisah mengukur throughput setiap detik dan melakukan polling lag setiap 1,5 detik. Timer ketiga mengirim snapshot setiap 400ms. Browser butuh state terkini, bukan setiap tick di antaranya. Menggabungkan update jadi snapshot mengurangi pekerjaan rendering.
Dua pod, satu local join
Kedua consumer tergabung di live-evaluators. Keduanya memakai custom co partition assigner dan subscribe ke market-updates dan notifications:
consumer: kafka.consumer({
groupId: "live-evaluators",
partitionAssigners: [coPartitionAssigner],
}),
Assigner itu menaruh partition p dari kedua topic di pod yang sama. Handler eachMessage di pod itu bisa men-join record kalau harga lokal yang dibutuhkan udah tersedia. Sebuah market tick meng-update tabel harga global di dashboard dan map localPrice milik pod:
if (topic === TOPICS.marketUpdates) {
const tick = tickFromKafka(message.value);
this.prices.set(tick.symbol, tick.price);
pod.localPrice.set(tick.symbol, tick.price); // local, co-located state
for (const n of pod.evaluator.process(tick)) {
await this.producer?.send({ topic: TOPICS.notifications, messages: [notifToKafka(n)] });
}
} else {
// Co-partitioned with market-updates, so the symbol's price is already
// sitting on THIS pod - enrich locally, no remote lookup.
const n = notifFromKafka(message.value);
if (pod.localPrice.has(n.symbol)) pod.localJoins++;
else pod.remoteNeeded++;
}
pod.evaluator memakai AlertEvaluator yang stateful untuk menyimpan harga terakhir per simbol dan mendeteksi persilangan threshold. Counter localJoins dan remoteNeeded mencatat apakah sebuah notifikasi menemukan state harga lokal. Keduanya menunjukkan hasil join yang teramati, bukan jaminan bahwa setiap lookup ke depannya bakal berhasil.
Kenapa Server Sent Events, bukan WebSockets
Data snapshot mengalir dari engine ke browser. Pause dan resume memakai endpoint POST yang terpisah. Saya memilih SSE karena feed ini nggak butuh message dari browser lewat koneksi yang sama.
SSE memakai HTTP, dan EventSource di browser menangani reconnection. Handler App Router di Next.js bisa me-return ReadableStream di dalam sebuah Response. TextEncoder mengubah setiap event jadi byte. Server meng-encode setiap snapshot sebagai message SSE. Ini udah cukup untuk feed snapshot yang periodik.
Route handler: response yang nggak pernah selesai
Route-nya sendiri pendek. Route ini mengambil engine singleton, membuka stream, dan menyambungkan subscriber callback milik engine langsung ke stream controller:
export const dynamic = "force-dynamic";
export async function GET(req: Request): Promise<Response> {
const engine = getLiveEngine();
await engine.ensureStarted();
const encoder = new TextEncoder();
let unsubscribe = () => {};
const stream = new ReadableStream({
start(controller) {
const send = (data: unknown) => {
controller.enqueue(encoder.encode(`data: ${JSON.stringify(data)}\n\n`));
};
send(engine.getSnapshot()); // prime immediately
unsubscribe = engine.subscribe(send);
req.signal.addEventListener("abort", () => {
unsubscribe();
controller.close();
});
},
cancel() {
unsubscribe();
},
});
return new Response(stream, {
headers: {
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache, no-transform",
Connection: "keep-alive",
},
});
}
Route ini bergantung pada detail-detail berikut:
export const dynamic = "force-dynamic"mengeluarkan route ini dari static optimization. Stream yang berumur panjang adalah kebalikan dari sesuatu yang bisa kamu cache atau pre render.- Snapshot langsung dikirim begitu terhubung, sebelum subscribe. Tanpa itu, tab yang baru dibuka bakal menampilkan layar kosong sampai 400ms sambil menunggu broadcast terjadwal berikutnya.
req.signalmemicuabortwaktu client terputus, entah karena tab-nya ditutup atau koneksinya putus. Ini memanggilunsubscribe()dan menghapus callback dari setsubscribersmilik engine. Tanpa cleanup ini, callback bakal tetap tertinggal setelah tab-nya ditutup.- Header response (
no-cache, no-transform,keep-alive) menjelaskan perilaku streaming. Setting buffering dan timeout di proxy juga perlu diverifikasi di jalur deployment.
Sequence di bawah menunjukkan snapshot awal, broadcast berikutnya, dan cleanup setelah disconnect:
sequenceDiagram
participant B as Browser
participant H as SSE route handler
participant E as Live engine
B->>H: GET /api/live
H->>E: ensureStarted()
H->>B: send snapshot (priming)
H->>E: subscribe(send)
loop every 400ms
E->>H: broadcast(snapshot)
H->>B: send snapshot
end
B--xH: tab closes, connection aborts
H->>E: unsubscribe()
H->>B: controller.close()
Pause dan resume tanpa menyentuh stream
Kontrol lewat route yang terpisah, bukan lewat message yang dikirim balik melalui koneksi SSE:
export async function POST(req: Request): Promise<Response> {
const { action } = (await req.json().catch(() => ({ action: "" }))) as { action?: string };
const engine = getLiveEngine();
if (action === "pause") {
engine.setRunning(false);
} else if (action === "resume") {
await engine.ensureStarted();
engine.setRunning(true);
}
return Response.json({ running: engine.isRunning() });
}
Pause membalik sebuah boolean yang dicek produceTick sebelum mengerjakan apa pun. Producer dan kedua consumer tetap terhubung, jadi resume langsung jalan: nggak ada reconnect, nggak ada rebalance. Indikator running di dashboard datang langsung dari snapshot broadcast berikutnya, jadi UI nggak pernah perlu melacak optimistic state-nya sendiri.
Menjaga client broker tetap di luar bundle
kafkajs dan ioredis itu Node native: keduanya membuka raw TCP socket dan mengandalkan call require dinamis yang nggak bisa di-resolve secara statis oleh bundler. Kalau dibiarkan, Next.js bakal mencoba mem-bundle keduanya untuk route handler, dan hasilnya build gagal atau sesuatu yang rusak waktu runtime. Config Next mengecualikan keduanya secara eksplisit:
const nextConfig: NextConfig = {
// kafkajs and ioredis are Node-native (net/tls, dynamic requires). Keep them
// out of the bundler so route handlers load them as normal Node modules.
serverExternalPackages: ["kafkajs", "ioredis", "pino", "pino-pretty"],
};
Konfigurasi ini juga meng-externalize pino dan pino-pretty untuk setup logging aplikasi. Client broker butuh runtime Node.js karena keduanya membuka TCP socket. Route-nya nggak bisa memakai export const runtime = "edge" dengan client-client ini.
Apa yang akhirnya tampil di layar
Dashboard-nya adalah client component. Satu useEffect membuka EventSource ke /api/live. Setiap message mengganti state dengan snapshot terbaru lewat setSnap(JSON.parse(e.data)). Component-nya nggak menggabungkan update parsial atau memakai store terpisah.
Yang tampil di layar adalah semua yang dilacak engine:
- Tile throughput: Kafka produced/sec, Kafka consumed/sec, Redis
XADD/sec, dan total berjalan dari alert yang terpicu, langsung dari sampling rate per detik di engine. - Papan ticker: kesepuluh simbol dengan harga terkini dan persentase perubahan terhadap harga dasarnya, diwarnai naik atau turun.
- Feed notifikasi: 16 alert terakhir, masing-masing menampilkan simbol, operator, threshold, dan harga yang memicunya.
- Kepemilikan partition dan join: setiap pod menampilkan partition-nya untuk
market-updatesdannotifications. Sebuah badge menunjukkan apakah daftar partition-nya cocok. Counter menunjukkan join lokal dan remote. - Consumer lag: perhitungan
high - committedyang sama dari post soal consumer group, di-polling setiap 1,5 detik dan diwarnai merah kalau melewati threshold.
Dashboard ini memakai engine broker dari post-post sebelumnya. Setelah pause, producer berhenti dan consumer bisa menghabiskan backlog-nya. Kegagalan broker atau rebalance mengubah assignment yang ditampilkan. Pengamatan ini membantu memeriksa perilaku selagi sistemnya jalan.
Dashboard ini menunjukkan perilaku, tapi nggak membuktikan performa. Artikel berikutnya menjelaskan benchmark harness dan batasannya.

Memuat komentar...