Leaky bucket rate limiting algoritmi: Node.jsda noldan yozamiz
Assalamu Alaykum bugun Leaky Bucketuslubida ishlaydigan Nodejs uchun rate limiter middleware yozib uni tekshirib ko'ramiz !!
Bu maqoladan oldin ushbu rate limiting haqidagi maqolani o'qib chiqing . Rate limiting nima ?

Rate limiterni yozishni boshlaymiz
Bu algoritm uchun bizga har soniyada qancha so'rovni bucketdan olishi hamda bucketning sig'imi kerak bo'ladi. Bu uchun LeakyBucket classini yozamiz :
class LeakyBucket {
constructor(capacity, consumeRate) {
this.capacity = capacity;
this.consumeRate = consumeRate;
}
}Endi so'rov kelganda uni qabul qilish yoki qilmasligni tekshiradigan metodni yozamiz :
class LeakyBucket {
constructor(capacity, consumeRate) {
this.capacity = capacity;
this.consumeRate = consumeRate;
this.bucketQueue = [];
}
fill(r) {
if (this.bucketQueue.length >= this.capacity) {
return false;
}
this.bucketQueue.push(r);
return true;
}
}bucketQueue bu qabul qilingan so'rovlarni ushlab turadi.
Endi buni vaqt bilan so'rovlarni process qiladigan consume metodini yozamiz.
class LeakyBucket {
constructor(capacity, consumeRate) {
this.capacity = capacity;
this.consumeRate = consumeRate;
this.bucketQueue = [];
this.lastConsumeTimestamp = Date.now();
}
...
consume() {
const now = Date.now();
if (this.bucketQueue.length === 0) {
this.lastConsumeTimestamp = now;
return [];
}
const elapsedTime = (now - this.lastConsumeTimestamp) / 1000;
const requestToConsumeCount = Math.floor(elapsedTime * this.consumeRate);
if (requestToConsumeCount === 0) {
return [];
}
this.lastConsumeTimestamp += (requestToConsumeCount / this.consumeRate) * 1000;
return this.bucketQueue.splice(0, requestToConsumeCount);
}
}lastConsumeTimestamp bizga qancha so'rovni server qabul qila olishini bilishga yordam beradi.
bizda so'rovlarni qabul qilib queue ga qo'shadigan va undan olib process qilishga beradigan metodlar tayyor endi shu process qilishni ishga tushuradigan leak metodini yozamiz :
class LeakyBucket {
...
startLeaking(intervalMs = 100) {
if (this.leakTimer) return;
this.leakTimer = setInterval(() => {
if (this.bucketQueue.length > 0) {
const consumed = this.consume();
if (consumed.length > 0) {
this.onLeak?.(consumed);
}
}
}, intervalMs);
}
}Deyarli tugadi endi leakni to'xtatadigan va hozrigi holatni olib beradigan metodlarni yozamiz :
class LeakyBucket {
constructor(capacity, consumeRate) {
this.capacity = capacity;
this.consumeRate = consumeRate;
this.bucketQueue = [];
this.lastConsumeTimestamp = Date.now();
this.leakTimer = null;
}
fill(r) {
if (this.bucketQueue.length >= this.capacity) {
return false;
}
if (this.bucketQueue.length === 0) {
this.lastConsumeTimestamp = Date.now();
}
this.bucketQueue.push(r);
return true;
}
consume() {
const now = Date.now();
if (this.bucketQueue.length === 0) {
this.lastConsumeTimestamp = now;
return [];
}
const elapsedTime = (now - this.lastConsumeTimestamp) / 1000;
const requestToConsumeCount = Math.floor(elapsedTime * this.consumeRate);
if (requestToConsumeCount === 0) {
return [];
}
this.lastConsumeTimestamp +=
(requestToConsumeCount / this.consumeRate) * 1000;
return this.bucketQueue.splice(0, requestToConsumeCount);
}
startLeaking(intervalMs = 100) {
if (this.leakTimer) return;
this.leakTimer = setInterval(() => {
if (this.bucketQueue.length > 0) {
const consumed = this.consume();
if (consumed.length > 0) {
this.onLeak?.(consumed);
}
}
}, intervalMs);
}
stopLeaking() {
clearInterval(this.leakTimer);
this.leakTimer = null;
}
getCurrentState() {
return {
queued: this.bucketQueue.length,
capacity: this.capacity,
consumeRate: this.consumeRate,
isLeaking: this.leakTimer !== null,
};
}
}rate limiter logikasi tugadi endi har bir endpoint yoki user uchun bucket yaratadigan class yozamiz (LeakyBucketRateLimiter ):
class LeakyBucketRateLimiter {
constructor(capacity, consumeRate, intervalMs = 100) {
this.buckets = new Map();
this.capacity = capacity;
this.consumeRate = consumeRate;
this.intervalMs = intervalMs;
}
getBucket(key) {
if (!this.buckets.has(key)) {
const bucket = new LeakyBucket(this.capacity, this.consumeRate);
bucket.onLeak = (items) => {
for (const item of items) item.resolve(true);
};
bucket.startLeaking(this.intervalMs);
this.buckets.set(key, bucket);
}
return this.buckets.get(key);
}
schedule(key) {
const bucket = this.getBucket(key);
return new Promise((resolve) => {
const accepted = bucket.fill({ resolve });
if (!accepted) resolve(false);
});
}
getCurrentState(key) {
return this.getBucket(key).getCurrentState();
}
stopAll() {
for (const bucket of this.buckets.values()) {
bucket.stopLeaking();
}
}
}Endi buni serverga qo'yib tekshirib ko'rishimiz kerak .
serverda /leaky-bucket endpointga rate limiter sifatida ulaymiz.
const http = require("http");
const { LeakyBucketRateLimiter } = require("./leaky-bucket/rate-limiter");
const leakyRateLimiter = new LeakyBucketRateLimiter(5, 1);
function rateLimiterMiddleware(req, res, next) {
const userId = req.headers["x-user-id"] || "default-user";
if (rateLimiter.isAllowed(userId)) {
next();
} else {
res.writeHead(429, { "Content-Type": "application/json" });
res.end(
JSON.stringify({
message: "So'rov rad etildi - limit oshib ketdi",
state: rateLimiter.getCurrentState(userId),
}),
);
}
}
async function handleRequest(req, res) {
const userId = req.headers["x-user-id"] || "default-user";
if (req.url === "/leaky-bucket") {
const accepted = await leakyRateLimiter.schedule(userId);
if (accepted) {
res.writeHead(200, { "Content-Type": "application/json" });
res.end(
JSON.stringify({
message: "So'rov qayta ishlandi (navbatdan oqib chiqdi)",
state: leakyRateLimiter.getCurrentState(userId),
}),
);
} else {
res.writeHead(429, { "Content-Type": "application/json" });
res.end(
JSON.stringify({
message: "So'rov rad etildi - navbat to'lgan",
state: leakyRateLimiter.getCurrentState(userId),
}),
);
}
} else {
res.writeHead(404, { "Content-Type": "application/json" });
res.end(JSON.stringify({ message: "Not found" }));
}
}
const server = http.createServer(handleRequest);
server.listen(3000, () => {
console.log("Server 3000-portda ishga tushdi...");
});endi shu urlga har xil userlar sifatida so'rovlar yuborib ishlayaptimi yoqmi tekshirib ko'ramiz bu uchun foydalanuvchilarni mocklaydigan script yozamiz :
const http = require("http");
const SERVER_URL = "http://localhost:3000/leaky-bucket";
const CAPACITY = 5;
const CONSUME_RATE = 1;
function sendRequest(userId) {
const start = Date.now();
return new Promise((resolve) => {
const options = { headers: { "x-user-id": userId } };
http
.get(SERVER_URL, options, (res) => {
let raw = "";
res.on("data", (chunk) => (raw += chunk));
res.on("end", () => {
resolve({ status: res.statusCode, latency: Date.now() - start });
});
})
.on("error", (err) => {
resolve({ status: null, error: err.message });
});
});
}
async function burstUser(userId, count) {
const results = await Promise.all(
Array.from({ length: count }, () => sendRequest(userId)),
);
const accepted = results.filter((r) => r.status === 200);
const rejected = results.filter((r) => r.status === 429);
return { userId, accepted, rejected, results };
}
(async () => {
const USERS = ["alice", "bob", "charlie"];
const REQUESTS_PER_USER = 8;
console.log(`${"─".repeat(55)}`);
console.log(` Leaky-Bucket Server Test`);
console.log(
` Capacity: ${CAPACITY} queued | Leak: ${CONSUME_RATE} request/sec`,
);
console.log(` Each user fires ${REQUESTS_PER_USER} concurrent requests`);
console.log(`${"─".repeat(55)}\n`);
console.log("Phase 1 — Concurrent burst (all users simultaneously)\n");
const results = await Promise.all(
USERS.map((userId) => burstUser(userId, REQUESTS_PER_USER)),
);
let allGood = true;
for (const log of results) {
const latencies = log.accepted.map((r) => `${r.latency}ms`);
console.log(` User: ${log.userId}`);
console.log(
` Accepted : ${log.accepted.length} (leaked at [${latencies.join(", ")}])`,
);
console.log(` Rejected : ${log.rejected.length} (queue was full)`);
const okAccepted = log.accepted.length === CAPACITY;
const okRejected = log.rejected.length === REQUESTS_PER_USER - CAPACITY;
const paced = log.accepted.some((r) => r.latency >= 500);
if (!okAccepted || !okRejected) {
allGood = false;
console.log(
` ✗ Expected ${CAPACITY} accepted / ${REQUESTS_PER_USER - CAPACITY} rejected`,
);
} else {
console.log(` ✓ Capacity enforced (${CAPACITY} accepted)`);
}
if (!paced) {
allGood = false;
console.log(` ✗ Accepted requests were not paced (leaked instantly)`);
} else {
console.log(` ✓ Accepted requests were paced by leak rate`);
}
console.log();
}
console.log("Phase 2 — After draining, a fresh request is accepted again\n");
await new Promise((resolve) => setTimeout(resolve, 1500));
const followUp = await sendRequest("alice");
console.log(
` alice follow-up → status ${followUp.status} (${followUp.latency}ms)`,
);
if (followUp.status !== 200) {
allGood = false;
console.log(" ✗ Expected 200 after queue drained");
} else {
console.log(" ✓ Bucket recovered — request accepted");
}
console.log(`\n${"─".repeat(55)}`);
console.log(allGood ? " All e2e checks passed." : " Some e2e checks FAILED.");
process.exit(allGood ? 0 : 1);
})();Endi buni tekshirib ko'ramiz :

Hamma kodlarni ushbu repodan olishingiz mumkin. Repo
Xulosa
Leaky bucket algoritmi serverlarga doim bir xil trafik kerak bo'lganda ishlatiladi.
Xulosa qilib aytganda rate limiting implimentatsiya qilayotganda dasturning kattaligi , va sizga keladigan so'rovlar uslubini bilishingiz kerak. Kerakli so'rovlarni bekor qilmaslik uchun. Hozirgi implimentatsiya bu juda oddiy ammo algoritmni to'liq tushunishga yordam beradi.
Bunda tashqari har doim foydalanuvchilar bilan rete limit haqidagi ma'lumotlarni berish , va bu orqali ular o'zlariga mos qayta urinib ko'rishlari mumkin.