Jawaban singkat untuk pertanyaan yang paling sering diajukan pembaca tentang topik ini.
01Apa itu backpressure dengan bahasa sederhana?
Backpressure adalah consumer yang lambat mendorong balik producer yang cepat agar producer melambat, berhenti sejenak, atau membuang pekerjaan. Tujuannya mencegah buffer di antara keduanya tumbuh tanpa batas. Sinyalnya bisa berupa bounded queue yang penuh, write() yang mengembalikan false, TCP window, atau HTTP 429 atau 503.
02Bagaimana cara menghentikan producer cepat agar tidak merobohkan consumer lambat?
Beri batas pada setiap queue dan tentukan apa yang terjadi saat penuh: block producer, reject pekerjaan baru, atau drop item yang paling tidak bernilai. Biarkan consumer menentukan tempo lewat pull-based consumption atau event drain. Di edge publik, lakukan load shedding dengan 429 atau 503 dan header Retry-After.
03Apa arti write() mengembalikan false di Node.js?
Artinya stream ingin Anda menunggu event drain sebelum menulis data lagi. Dokumentasi Node.js menyatakan jangan menulis chunk lagi setelah false sampai drain dipancarkan. Jika diabaikan, buffer internal tumbuh pada setiap write tambahan, seperti terlihat pada demo di artikel ini.
04Sebaiknya API yang overload mengembalikan 429 atau 503?
Kembalikan 429 Too Many Requests saat satu client melewati batasnya sendiri, dan 503 Service Unavailable saat seluruh service overload. Keduanya bisa membawa header Retry-After berisi jumlah detik atau tanggal HTTP. Client sebaiknya menambah jitter acak agar pemanggil yang ditolak tidak retry di saat yang sama.
05Apakah unbounded queue pernah boleh dipakai?
Hanya jika Anda yakin arrival rate tidak akan lama melebihi service rate, misalnya batch kecil dan terbatas yang Anda buat sendiri. Untuk apa pun yang digerakkan trafik luar, unbounded queue hanya menunda kegagalan dan menambah waktu tunggu. Batasi queue dan pilih kebijakan saat queue penuh.
Backpressure adalah sinyal dari consumer yang lambat agar producer yang cepat memperlambat laju, sehingga queue tetap terbatas dan tidak tumbuh sampai memori habis. Terapkan lewat bounded queue, pull-based consumption, write() yang mengembalikan false lalu menunggu drain di Node.js, TCP window, atau load shedding dengan 429 atau 503 dan Retry-After.
Queue terlihat seperti asuransi gratis: saat consumer lambat, producer tinggal terus menambah dan queue menyerap selisihnya. Itu hanya benar selama consumer rata-rata minimal secepat producer. Begitu tidak, queue berhenti menjadi buffer dan berubah menjadi memory leak yang pelan, lengkap dengan latency di belakangnya.
Saya menjalankan Docker, Redis, dan Postgres di satu VPS, jadi scale out untuk menghindari queue yang membengkak bukan opsi yang bisa saya andalkan. Artikel ini menjawab satu pertanyaan: bagaimana menghentikan producer yang cepat agar tidak merobohkan consumer yang lambat? Kita mulai dari hitungan Little's law, lalu bounded queue, load shedding, worker berbasis pull, stream Node.js dan TCP, dan ditutup dengan aturan autoscaling berdasarkan queue depth. Output demo Node.js di sini asli, dari run yang saya jalankan untuk artikel ini.
Apa itu backpressure di distributed system?
Backpressure adalah mekanisme apa pun yang membuat consumer yang tidak sanggup mengimbangi mendorong hambatan kembali ke producer, sehingga producer melambat, berhenti sejenak, atau membuang pekerjaan alih-alih menumpuknya. Istilahnya berasal dari fluid dynamics, tempat penyempitan di hilir menaikkan tekanan di hulu. Di software, penyempitan itu bisa berupa database yang lambat, API dengan rate limit, atau worker pool yang sibuk.
Ide kuncinya: sinyal berjalan berlawanan arah dengan aliran data. Tanpa sinyal itu, producer baru tahu ada masalah dari timeout atau crash, jauh setelah kerusakan terjadi. Inisiatif Reactive Streams menyatakan tujuannya dengan lugas: mengatur pertukaran data stream melewati batas asynchronous agar penerima tidak dipaksa menampung data tanpa batas.
Apa yang terjadi jika producer lebih cepat dari consumer?
Queue tumbuh tanpa batas, memori naik, dan waktu tunggu ikut naik. Little's law menghubungkan rata-rata jumlah item dalam sistem yang stabil dengan arrival rate dan waktu yang dihabiskan tiap item di sana. Hukum ini butuh kondisi steady state, jadi untuk kasus kelebihan beban hitungannya cukup arrival rate dikurangi service rate. Kedua kasus dihitung di bawah dengan angka asumsi.
Stable system (Little's law, L = lambda x W)
lambda = 100 requests/s arriving, W = 0.05 s inside the service
L = 100 x 0.05 = 5 requests in the system at any moment
Same traffic, downstream now slow (W = 2 s)
L = 100 x 2 = 200 requests in flight
at an assumed 256 KiB of buffered state each: 200 x 256 KiB = 50 MiB
Unstable system (arrivals beat the service rate, so there is no steady state)
lambda = 200 messages/s in, mu = 150 messages/s out
queue grows by 200 - 150 = 50 messages/s
after 1 hour: 50 x 3600 = 180,000 messages queued
at an assumed 2 KiB each: 180,000 x 2048 = 368,640,000 bytes = 351.6 MiB
the next message waits 180,000 / 150 = 1,200 s = 20 minutes
Perhatikan dua biaya yang berbeda. Biaya memori dibatasi oleh mesin, jadi berakhir pada proses yang dibunuh karena out of memory. Biaya latency tidak dibatasi apa pun: pesan yang menunggu 20 menit di queue sering sudah tidak berguna bagi user yang mengirimnya, sehingga kapasitas consumer terbuang untuk pekerjaan yang tidak lagi ditunggu siapa pun.
Queue tanpa batas tidak menghilangkan overload, hanya menyembunyikannya. Kegagalan berpindah dari error cepat yang terlihat di edge menjadi crash lambat di tempat yang tidak Anda pantau, biasanya setelah queue sudah merusak latency semua orang.
Apakah bounded queue solusinya?
Bounded queue adalah bentuk backpressure paling minimal, karena memaksa adanya keputusan saat queue penuh, bukan tidak pernah. Namun membatasi saja belum cukup. Ia mengubah masalah memori tanpa batas menjadi pertanyaan: apa yang harus dialami producer saat queue penuh?
Ada tiga jawaban yang jujur, dan masing-masing cocok untuk beban kerja yang berbeda.
Block producer: pemanggilan push menunggu sampai ada ruang. Cocok jika producer adalah worker internal yang aman untuk berhenti, seperti di stream pipeline.
Reject pekerjaan baru: push gagal cepat dan pemanggil yang memutuskan. Cocok di edge yang menghadap user atau service lain, dan dasar dari load shedding.
Drop item tertua atau yang paling tidak bernilai: pertahankan data terbaru. Cocok untuk metrics, posisi, dan apa pun yang nilainya basi jika terlambat, tapi salah untuk pembayaran atau order.
Bagaimana load shedding bekerja dengan 429, 503, dan Retry-After?
Load shedding menolak pekerjaan di edge begitu batas concurrency tercapai, sehingga pekerjaan yang diterima tetap selesai dengan cepat. HTTP sudah menyediakan kosakatanya. RFC 9110 mendefinisikan 503 Service Unavailable sebagai server yang tidak bisa menangani request karena overload sementara, dan mengizinkan header Retry-After, berupa jumlah detik atau tanggal HTTP, untuk menyarankan berapa lama client menunggu.
import { createServer } from "node:http";
const MAX_IN_FLIGHT = 50; // sized from Little's law: arrival rate x acceptable latency
let inFlight = 0;
createServer(async (req, res) => {
if (inFlight >= MAX_IN_FLIGHT) {
// Refuse in microseconds instead of queueing the work and timing out later.
res.writeHead(503, {
"Retry-After": "5", // delay-seconds form, RFC 9110 section 10.2.3
"Content-Type": "application/json",
});
res.end(JSON.stringify({ error: "overloaded", retryAfterSeconds: 5 }));
return;
}
inFlight++;
try {
await handle(req, res); // your real handler
} finally {
inFlight--; // always release, or the limit leaks shut
}
}).listen(3000);
// Per-client limit instead? Same shape, but answer 429 Too Many Requests.
Gunakan 429 Too Many Requests saat satu client melewati batasnya sendiri, dan 503 saat seluruh service jenuh. RFC 6585 mendefinisikan 429 sebagai user yang mengirim terlalu banyak request dalam waktu tertentu dan mengizinkan Retry-After di sana juga. RFC 9110 juga mencatat bahwa server tidak wajib memakai 503 saat overload; sebagian server langsung menolak koneksi.
Retry-After hanya berguna jika client mematuhinya. Buat client menunggu sejumlah detik yang disebutkan ditambah jitter acak, kalau tidak semua pemanggil yang ditolak akan kembali di detik yang sama dan membuat lonjakan yang baru saja Anda redam.
Bagaimana Node.js menangani backpressure pada stream?
Dokumentasi stream Node.js menyebut writable.write() mengembalikan false jika stream ingin kode pemanggil menunggu event drain sebelum menulis data lagi, dan bahwa setelah false Anda tidak boleh menulis chunk lagi sampai drain dipancarkan. Demo di bawah mengabaikan aturan itu di satu loop dan mematuhinya di loop lain. Consumer-nya sengaja lambat, dan highWaterMark diset ke 16 objek.
import { Writable } from "node:stream";
import { once } from "node:events";
const ITEMS = 100_000;
const PAYLOAD_BYTES = 1024;
const HIGH_WATER_MARK = 16; // objectMode: counts objects, not bytes
// A deliberately slow consumer: each item needs one trip round the event loop.
function makeSlowSink() {
let processed = 0;
const sink = new Writable({
objectMode: true,
highWaterMark: HIGH_WATER_MARK,
write(_chunk, _enc, done) {
processed++;
setImmediate(done);
},
});
return { sink, processed: () => processed };
}
async function withoutBackpressure() {
const { sink, processed } = makeSlowSink();
let falseCount = 0;
for (let i = 0; i < ITEMS; i++) {
// Wrong: the return value is ignored, so every item piles up in the writable buffer.
if (!sink.write(Buffer.alloc(PAYLOAD_BYTES))) falseCount++;
}
console.log("WITHOUT backpressure");
console.log(" write() returned false", falseCount, "times, and the loop ignored it");
console.log(" writableLength right after the loop:", sink.writableLength, "items");
console.log(" processed by consumer so far:", processed());
sink.end();
await once(sink, "finish");
}
async function withBackpressure() {
const { sink, processed } = makeSlowSink();
let maxQueued = 0;
let drains = 0;
for (let i = 0; i < ITEMS; i++) {
const ok = sink.write(Buffer.alloc(PAYLOAD_BYTES));
maxQueued = Math.max(maxQueued, sink.writableLength);
// Right: stop producing until the stream says it has room again.
if (!ok) {
await once(sink, "drain");
drains++;
}
}
console.log("WITH backpressure");
console.log(" highWaterMark:", HIGH_WATER_MARK, "items");
console.log(" max writableLength seen:", maxQueued, "items");
console.log(" drain events awaited:", drains);
sink.end();
await once(sink, "finish");
console.log(" processed by consumer:", processed());
}
// Load shedding: a bounded queue that refuses work instead of growing.
function sheddingDemo() {
const CAPACITY = 100;
const ARRIVALS = 1000;
const queue: number[] = [];
let shed = 0;
for (let i = 0; i < ARRIVALS; i++) {
if (queue.length >= CAPACITY) shed++; // would answer 429 or 503 + Retry-After
else queue.push(i);
}
console.log("LOAD SHEDDING");
console.log(" arrivals:", ARRIVALS, "capacity:", CAPACITY);
console.log(" accepted:", queue.length, "shed:", shed);
}
console.log("node", process.version, "items:", ITEMS, "payload:", PAYLOAD_BYTES, "bytes");
await withoutBackpressure();
await withBackpressure();
sheddingDemo();
Ini output dari menjalankan file tersebut di Node v26.10.0. Saya membuang pembacaan memori dari script karena angka memori proses berbeda di setiap run, jadi hanya jumlah queue yang ditampilkan. Jumlah itu deterministik untuk script ini.
$ node backpressure-demo.ts
node v26.10.0 items: 100000 payload: 1024 bytes
WITHOUT backpressure
write() returned false 99985 times, and the loop ignored it
writableLength right after the loop: 100000 items
processed by consumer so far: 1
WITH backpressure
highWaterMark: 16 items
max writableLength seen: 16 items
drain events awaited: 6250
processed by consumer: 100000
LOAD SHEDDING
arrivals: 1000 capacity: 100
accepted: 100 shed: 900
Bandingkan kedua blok. Mengabaikan return value membuat 100.000 item menumpuk di stream, sementara consumer baru memproses satu. Mematuhinya tidak pernah membiarkan buffer melewati 16 item, dan 6.250 event drain hanyalah 100.000 item dibagi 16 yang muat setiap kali. Blok shedding murni aritmetika: 1.000 kedatangan ke queue 100 slot menerima 100 dan menolak 900. Jika Anda memakai stream.pipeline atau pipe, library yang melakukan penungguan ini; Anda hanya menulis loop sendiri saat memanggil write secara langsung.
Bagaimana perbandingan pull-based worker, TCP window, dan reactive streams?
Semuanya memecahkan masalah yang sama dengan memindahkan siapa yang menentukan laju. TCP adalah contoh aslinya: penerima mengiklankan berapa octet yang mau ia terima, dan pengirim tidak boleh melebihinya, seperti dijelaskan di RFC 9293. Reactive Streams menggeneralisasi ide itu dengan membuat subscriber meminta sejumlah item, sehingga publisher tidak pernah mengirim lebih dari yang diminta.
Mekanisme
Siapa yang menentukan laju
Sinyal
Jika sinyal diabaikan
Node.js write() dan drain
Writable stream
write() mengembalikan false, lalu drain dipancarkan
Item menumpuk di buffer internal
TCP receive window
Penerima
Ukuran window diiklankan di setiap segment
Pengirim membanjiri penerima
Load shedding
Edge server
429 atau 503 dengan Retry-After
Client langsung retry dan memperdalam overload
Pull-based worker
Consumer
Consumer meminta batch berikutnya
Tidak mungkin: tidak ada yang dikirim sebelum diminta
Reactive Streams request(n)
Subscriber
Subscriber meminta sejumlah item
Melanggar spec: publisher tidak boleh mengirim lebih dari yang diminta
Pull paling mudah dipahami karena consumer mengendalikan intake-nya sendiri. Worker berbasis queue yang mengambil batch kecil, menyelesaikannya, baru mengambil lagi tidak akan pernah kebanjiran, secepat apa pun job di-enqueue.
// Pull-based worker: it only asks for work when it has capacity, so the
// producer side can never push more than the workers can hold.
const BATCH = 10;
const IDLE_MS = 500;
async function runWorker(claim: (n: number) => Promise<Job[]>, handle: (j: Job) => Promise<void>) {
while (true) {
const jobs = await claim(BATCH); // e.g. SELECT ... FOR UPDATE SKIP LOCKED LIMIT 10
if (jobs.length === 0) {
await new Promise((r) => setTimeout(r, IDLE_MS)); // empty queue: back off, do not spin
continue;
}
await Promise.all(jobs.map(handle)); // finish the batch BEFORE asking for more
}
}
interface Job { id: number; payload: unknown }
Harga dari pull adalah latency saat idle. Worker yang polling tiap 500 ms menambah delay sebesar itu pada job yang datang ke queue kosong, makanya long-polling atau kanal notifikasi sering dipasangkan dengannya.
Bagaimana autoscaling berdasarkan queue depth?
Queue depth adalah sinyal scaling yang lebih baik daripada CPU untuk worker, karena worker yang menunggu downstream lambat terlihat CPU-nya rendah padahal backlog terus tumbuh. Bagi arrival rate dengan service rate per worker untuk mendapat jumlah worker yang dibutuhkan, dan bagi backlog dengan kapasitas surplus untuk mendapat waktu pengurasan. Angka di bawah adalah asumsi untuk menunjukkan metodenya.
Worker capacity: mu = 50 messages/s per worker (20 ms of work each)
Arrival rate: lambda = 200 messages/s
workers needed to hold the line = ceil(lambda / mu) = ceil(200 / 50) = 4
Backlog already queued: B = 180,000 messages, scale to 6 workers (300/s)
net drain rate = 300 - 200 = 100 messages/s
time to clear = 180,000 / 100 = 1,800 s = 30 minutes
Estimated wait for a new message = depth / (workers x mu)
at 180,000 deep with 4 workers: 180,000 / 200 = 900 s
Kubernetes memungkinkan Horizontal Pod Autoscaler melakukan scaling berdasarkan external metrics, jadi panjang queue yang diekspor pipeline metrics Anda bisa menggerakkan jumlah replica, asalkan ada metrics adapter yang menyajikan external metrics API. Di satu VPS, hitungan yang sama tetap berguna: ia memberi tahu berapa container worker yang perlu dijalankan dan apakah mesin itu sanggup menguras backlog. Pilih umur pesan tertua bila bisa didapat, karena depth saja tidak menunjukkan berapa lama user menunggu.
Strategi backpressure mana yang sebaiknya dipilih?
Telusuri berurutan. Kebanyakan sistem berakhir memakai lebih dari satu, karena edge, queue, dan worker masing-masing butuh perlindungan sendiri.
Beri batas pada setiap queue dan buffer, dan tuliskan apa yang terjadi saat penuh: block, reject, atau drop.
Di edge publik, batasi concurrency dan jawab 503 dengan Retry-After, atau 429 per client, alih-alih mengantre melewati budget latency Anda.
Antar worker internal, pilih pull-based consumption dengan batch kecil agar consumer menentukan temponya sendiri.
Di Node.js, periksa return value write() dan tunggu drain, atau pakai pipeline agar library yang melakukannya.
Pasang alert pada queue depth dan umur pesan, lalu scale worker dari arrival rate dibagi throughput per worker.
Producer cepat hanya merobohkan consumer lambat ketika tidak ada bagian sistem yang bisa berkata tidak. Batasi queue, putuskan arti penuh, dan biarkan consumer menentukan tempo lewat pull, drain, atau window. Lakukan di setiap hop, karena sinyal yang berhenti di satu batas hanya memindahkan tumpukan ke batas berikutnya.