Redis से Pub/Sub, स्ट्रीम और दर सीमांकन
इवेंट प्रसारित कीजिए, टिकाऊ Redis Streams बनाइए और टोकन-बकेट दर सीमक परमाणु रूप से लागू कीजिए।
Redis से Pub/Sub, स्ट्रीम और दर सीमांकन, CoddyKit पर Node.js बैकएंड विकास बूटकैंप का एक निःशुल्क पाठ है। यह 4 में से 3वाँ पाठ है। आप नीचे पूरा पाठ निःशुल्क पढ़ सकते हैं—फिर अंतर्निहित कोड संपादक और 24/7 एआई ट्यूटर के साथ ब्राउज़र में इसका व्यावहारिक अभ्यास कर सकते हैं। यह Node.js बैकएंड विकास बूटकैंप सीखने के मार्ग का हिस्सा है और आपकी प्रगति वेब तथा CoddyKit ऐप पर सिंक होती रहती है। Node.js बैकएंड विकास बूटकैंप पाठ्यक्रम में कुल 4 पाठ शामिल हैं।
एक सर्वर में संदेश भेजने की तीन मूल सुविधाएँ
Redis केवल key/value कैश नहीं है। Node.js backend में यह एक हल्के संदेश ब्रोकर और समन्वय इंजन के रूप में भी काम करता है। इस पाठ में आप हर Redis इंस्टॉल के साथ उपलब्ध तीन उत्पादन-स्तरीय पैटर्न सीखेंगे:
- प्रकाशन/सदस्यता — सभी जुड़े हुए श्रोताओं को तुरंत, बिना किसी प्रत्युत्तर की प्रतीक्षा किए प्रसारण।
- Streams — consumer groups और replay वाला, केवल जोड़ने योग्य टिकाऊ लॉग।
- दर सीमित करना — दुरुपयोग करने वाले Client को धीमा करने वाले atomic counters।
इन सबको जोड़ने वाला मुख्य विचार यह है: Redis तक एक ही round-trip ऐसा काम कर सकता है जिसके लिए अन्यथा डेटाबेस, queue server और lock service की आवश्यकता होती। पूरे पाठ में हम ioredis Client का उपयोग करेंगे, जो Node backend के लिए वास्तविक मानक है क्योंकि यह pipelining, Lua scripting और cluster mode का समर्थन करता है।
प्रकाशन/सदस्यता: हर subscriber को प्रसारण
Redis Pub/Sub एक process को किसी channel पर संदेश PUBLISH करने देता है और हर SUBSCRIBE किए गए Client को वह तुरंत प्राप्त होता है। यह cache-invalidation fan-out या WebSocket server को लाइव अपडेट भेजने के लिए उत्तम है।
महत्वपूर्ण सावधानी: subscribe mode वाला connection सामान्य commands नहीं चला सकता। आपको GET/SET के लिए उपयोग किए जाने वाले connection से अलग एक समर्पित subscriber connection बनाना होगा। नीचे sub केवल सुनता है; publish करने के लिए दूसरे Client का उपयोग किया जाएगा।
const Redis = require('ioredis');
const sub = new Redis(); // dedicated subscriber connection
const pub = new Redis(); // separate connection for publishing
sub.subscribe('cache:invalidate', (err, count) => {
if (err) throw err;
console.log('Subscribed to ' + count + ' channel(s)');
});
sub.on('message', (channel, message) => {
console.log('[' + channel + '] ' + message);
});
// Another part of the app publishes an event
pub.publish('cache:invalidate', JSON.stringify({ key: 'user:42' }));पैटर्न सदस्यताएँ और बिना प्रत्युत्तर के प्रसारण की सीमा
सटीक channels के अलावा, Redis glob patterns के साथ PSUBSCRIBE का समर्थन करता है। order:* की सदस्यता लेने पर order:created, order:shipped आदि के संदेश प्राप्त होते हैं। callback pmessage होता है और इसमें मिलान किया गया pattern शामिल रहता है।
ध्यान रखने योग्य मुख्य सीमा यह है: Pub/Sub में कोई persistence नहीं होती। publish करते समय कोई subscriber जुड़ा न हो, तो संदेश हमेशा के लिए चला जाता है — न कोई buffering, न acknowledgement और न replay। यदि कोई subscriber क्रैश होकर फिर जुड़ता है, तो उसके बंद रहने के दौरान हुआ सब कुछ छूट जाता है।
- Pub/Sub का उपयोग ऐसे क्षणिक events के लिए करें जिनका छूट जाना स्वीकार्य हो।
- ऐसे events जिन्हें नहीं खोना चाहिए, उनके लिए Streams आवश्यक हैं (अगला अनुभाग)।
const Redis = require('ioredis');
const sub = new Redis();
sub.psubscribe('order:*', (err, count) => {
console.log('Pattern subscriptions: ' + count);
});
sub.on('pmessage', (pattern, channel, message) => {
console.log('matched ' + pattern + ' on ' + channel + ': ' + message);
});Streams: टिकाऊ और दोबारा पढ़े जा सकने वाला लॉग
Redis Stream केवल जोड़ने योग्य लॉग है। प्रत्येक entry को 1718000000000-0 जैसा लगातार बढ़ता हुआ ID मिलता है (milliseconds-sequence)। Pub/Sub के विपरीत, entries को तब तक संग्रहीत रखा जाता है जब तक आप उन्हें trim न करें, इसलिए देर से जुड़ने वाले या पुनः आरंभ किए गए consumers इतिहास को दोबारा पढ़ सकते हैं।
आप XADD से जोड़ते हैं। विशेष ID * Redis को अगला ID अपने-आप बनाने के लिए कहता है। ID के बाद आप field/value pairs देते हैं, बिल्कुल hash की तरह।
const Redis = require('ioredis');
const redis = new Redis();
async function publishOrder() {
// XADD key * field value field value ...
const id = await redis.xadd(
'stream:orders', '*',
'orderId', '42',
'amount', '99.90',
'status', 'created'
);
console.log('Appended entry with ID ' + id);
// Read the latest 5 entries (newest last)
const entries = await redis.xrange('stream:orders', '-', '+', 'COUNT', 5);
console.log(JSON.stringify(entries, null, 2));
}
publishOrder();Consumer Groups: दोहरे प्रसंस्करण के बिना विस्तार
Streams की वास्तविक शक्ति consumer groups हैं। एक group कई worker processes को एक stream का भार बाँटने देता है: प्रत्येक entry group के ठीक एक consumer को दी जाती है, जिससे क्षैतिज विस्तार संभव होता है।
XGROUP CREATE से group को एक बार बनाएँ। MKSTREAM विकल्प stream के मौजूद न होने पर उसे बना देता है। प्रारंभिक ID $ का अर्थ है 'केवल group बनाए जाने के बाद जोड़ी गई entries दें'; बिल्कुल शुरुआत से उपभोग करने के लिए 0 का उपयोग करें।
const Redis = require('ioredis');
const redis = new Redis();
async function setup() {
try {
await redis.xgroup(
'CREATE', 'stream:orders', 'order-workers', '$', 'MKSTREAM'
);
console.log('Group created');
} catch (e) {
if (e.message.includes('BUSYGROUP')) {
console.log('Group already exists, continuing');
} else {
throw e;
}
}
}
setup();Stream Entries पढ़ना और उनकी पुष्टि करना
प्रत्येक worker XREADGROUP से पढ़ता है और अपने group name तथा एक अद्वितीय consumer name को भेजता है। विशेष ID > का अर्थ है 'मुझे ऐसी entries दें जो इस group के किसी भी consumer को कभी नहीं दी गई हैं'।
जिस entry की पुष्टि नहीं हुई है, वह group की Pending Entries List (PEL) में बनी रहती है। प्रसंस्करण पूरा करने के बाद आप XACK call करते हैं, ताकि Redis जान सके कि उसे सुरक्षित रूप से संभाल लिया गया है। यदि आपका worker XACK से पहले क्रैश हो जाता है, तो entry pending रहती है और उसे फिर से प्राप्त किया जा सकता है — इसी कारण Streams कम-से-कम-एक-बार और टिकाऊ होते हैं।
const Redis = require('ioredis');
const redis = new Redis();
async function consume() {
const res = await redis.xreadgroup(
'GROUP', 'order-workers', 'worker-1',
'COUNT', 10, 'BLOCK', 5000,
'STREAMS', 'stream:orders', '>'
);
if (!res) return; // BLOCK timed out with no new entries
for (const [, entries] of res) {
for (const [id, fields] of entries) {
console.log('processing ' + id, fields);
// ... do real work here ...
await redis.xack('stream:orders', 'order-workers', id);
}
}
}
consume();अटके हुए संदेश फिर प्राप्त करना और छँटाई
यदि worker-1 प्रसंस्करण के बीच में बंद हो जाए तो क्या होगा? उसकी entries फिर से प्राप्त किए जाने तक PEL में हमेशा बनी रहेंगी। किसी स्वस्थ consumer को threshold से अधिक समय से idle entries हस्तांतरित करने के लिए XAUTOCLAIM (Redis 6.2+) का उपयोग करें:
XAUTOCLAIM stream group consumer min-idle-time start—min-idle-timems से अधिक समय से idle entries लेता है।- फिर से प्राप्त करने से पहले लंबित काम की जाँच
XPENDINGसे करें।
Streams हमेशा बढ़ते रहते हैं, इसलिए सीमित XADD से memory की सीमा तय करें: XADD key MAXLEN ~ 10000 * ...। ~ का अर्थ 'लगभग' है, जिससे Redis सटीक गिनती के बजाय पूरे macro-nodes में कुशलता से छँटाई कर सकता है।
const Redis = require('ioredis');
const redis = new Redis();
async function recover() {
// Reclaim entries idle > 30s, hand them to worker-2
const [cursor, claimed] = await redis.xautoclaim(
'stream:orders', 'order-workers', 'worker-2',
30000, '0', 'COUNT', 25
);
console.log('Reclaimed ' + claimed.length + ' entries; next cursor ' + cursor);
// Append with an approximate cap to bound memory
await redis.xadd('stream:orders', 'MAXLEN', '~', 10000, '*', 'orderId', '99');
}
recover();सरल दर-सीमांकन क्यों टूट जाता है
अब throttling पर आते हैं। पहली सहज कोशिश होती है कि GET से counter पढ़ें, Node में उसकी जाँच करें, फिर उसे वापस SET करें। यह एक पारंपरिक race condition है: आपके GET और SET के बीच कोई दूसरा समवर्ती request उसी पुराने मान को पढ़ लेता है और दोनों मानते हैं कि वे limit के भीतर हैं। load के दौरान आप अनुमति से कहीं अधिक requests को आगे जाने देते हैं।
समाधान atomicity है। Redis प्रत्येक command (और प्रत्येक Lua script) को single-threaded और अविभाज्य रूप से चलाता है। एक सरल fixed-window limiter INCR के साथ EXPIRE का उपयोग करता है, जिससे पढ़ना-बदलना-लिखना server-side बिना किसी अंतराल के होता है। नीचे दिया गया pure-JS demo इसे Redis में ले जाने से पहले race की अवधारणा दिखाता है।
// Demonstrates WHY check-then-set races. Two 'requests' interleave.
let counter = 0;
const LIMIT = 3;
function tryRequest(name) {
const current = counter; // read
if (current < LIMIT) {
// imagine an await here: another request runs before we write
counter = current + 1; // write (stale!)
return name + ': allowed (' + counter + ')';
}
return name + ': blocked';
}
// Both read 0 before either writes -> over-admission
const a = tryRequest('reqA');
const b = tryRequest('reqB');
console.log(a);
console.log(b);
console.log('Final counter: ' + counter);INCR + EXPIRE वाला Fixed-Window Limiter
सबसे सरल सही limiter: counter की key Client और time window के आधार पर बनाएँ, उसे INCR करें और window पहली बार खुलने पर TTL सेट करें। क्योंकि INCR नया मान atomic रूप से लौटाता है, इसलिए कोई race नहीं होती।
Fixed windows की कमजोरी सीमा पर अचानक बढ़ा हुआ traffic है: कोई Client 0:59 पर पूरा quota भेज सकता है और फिर 1:00 पर दोबारा, जिससे थोड़े समय के लिए दर दोगुनी हो जाती है। कई APIs के लिए यह स्वीकार्य है, लेकिन token buckets (अगले भाग में) इसे सुचारु कर देते हैं।
const Redis = require('ioredis');
const redis = new Redis();
async function allow(userId, limit = 100, windowSec = 60) {
const key = 'rl:' + userId + ':' + Math.floor(Date.now() / 1000 / windowSec);
const count = await redis.incr(key);
if (count === 1) {
await redis.expire(key, windowSec); // set TTL only on first hit
}
return count <= limit;
}
allow('user:42').then((ok) => {
console.log(ok ? 'request allowed' : '429 Too Many Requests');
});Lua Script वाला Atomic Token Bucket
Token bucket प्रत्येक Client को C क्षमता वाला bucket देता है, जिसमें R tokens/second की दर से tokens फिर भरते हैं। प्रत्येक request में एक token खर्च होता है; bucket खाली होने पर request अस्वीकार कर दी जाती है। यह नियंत्रित bursts की अनुमति देता है और साथ ही स्थिर औसत दर लागू करता है।
फिर से भरने और token खर्च करने को एक atomic चरण होना चाहिए, इसलिए हम उन्हें EVAL के माध्यम से Lua script के रूप में भेजते हैं। Redis पूरी script को अन्य commands के बीच में हस्तक्षेप किए बिना चलाता है। हम hash में दो fields रखते हैं — वर्तमान tokens और अंतिम refill का ts — और पिछली call के बाद से बने tokens की संख्या आवश्यकता पड़ने पर निकालते हैं।
const Redis = require('ioredis');
const redis = new Redis();
const LUA = [
"local cap = tonumber(ARGV[1])",
"local refill = tonumber(ARGV[2])",
"local now = tonumber(ARGV[3])",
"local cost = tonumber(ARGV[4])",
"local b = redis.call('HMGET', KEYS[1], 'tokens', 'ts')",
"local tokens = tonumber(b[1])",
"local ts = tonumber(b[2])",
"if tokens == nil then tokens = cap; ts = now end",
"local delta = math.max(0, now - ts)",
"tokens = math.min(cap, tokens + delta * refill)",
"local allowed = 0",
"if tokens >= cost then allowed = 1; tokens = tokens - cost end",
"redis.call('HMSET', KEYS[1], 'tokens', tokens, 'ts', now)",
"redis.call('EXPIRE', KEYS[1], 3600)",
"return allowed"
].join('\n');
async function take(userId) {
const now = Date.now() / 1000;
// capacity 10, refill 1 token/sec, cost 1
const allowed = await redis.eval(LUA, 1, 'tb:' + userId, 10, 1, now, 1);
return allowed === 1;
}
take('user:42').then((ok) => console.log(ok ? 'allowed' : 'throttled'));Express Middleware में Limiter जोड़ना
वास्तविक Node बैकएंड में सीमाकर्ता मिडलवेयर में रहता है, इसलिए हर रूट सुरक्षित रहता है। अस्वीकृति पर आप HTTP 429 लौटाते हैं और शिष्टाचार के तौर पर Retry-After हेडर भेजते हैं, जो क्लाइंट को बताता है कि दोबारा कब प्रयास करना है।
याद रखने योग्य सर्वोत्तम अभ्यास:
- कुंजी किसी स्थिर पहचान (API कुंजी या प्रमाणित उपयोगकर्ता आईडी) के आधार पर बनाइए, केवल IP के आधार पर नहीं — NAT के पीछे IP साझा किए जाते हैं।
- जानबूझकर खुला या बंद विफलन चुनें: यदि Redis तक पहुँचना संभव न हो, तो तय करें कि ट्रैफ़िक की अनुमति देनी है (उपलब्धता) या उसे रोकना है (सुरक्षा)। इसे स्पष्ट निर्णय बनाइए, दुर्घटनावश होने वाली स्थिति नहीं।
- दर-सीमा वाले हेडर लौटाइए, ताकि सही ढंग से व्यवहार करने वाले क्लाइंट स्वयं अपनी गति सीमित कर सकें।
function rateLimit(redis, take) {
return async (req, res, next) => {
const id = req.user?.id || req.ip;
try {
if (await take(id)) return next();
res.set('Retry-After', '1');
return res.status(429).json({ error: 'Too Many Requests' });
} catch (err) {
// Redis down: fail OPEN here (prioritize availability)
console.error('rate limiter degraded', err.message);
return next();
}
};
}
module.exports = { rateLimit };त्वरित जाँच: सही प्रिमिटिव चुनना
आप ऑर्डर-प्रोसेसिंग पाइपलाइन बना रहे हैं। वर्कर प्रक्रियाएँ क्रैश होकर फिर शुरू हो सकती हैं, और ऑर्डर का कोई भी इवेंट कभी खोना नहीं चाहिए; हर इवेंट को कई वर्करों में से ठीक एक द्वारा संसाधित किया जाना चाहिए और क्रैश हुए काम का अपने-आप दोबारा प्रयास होना चाहिए। Redis की कौन-सी सुविधा उपयुक्त होगी?
पुनरावलोकन: गारंटी के अनुरूप टूल चुनें
अब आपके पास Redis के संदेश-प्रेषण और नियंत्रण के तीन पैटर्न हैं और, इससे भी महत्वपूर्ण, इनके बीच चुनाव करने का विवेक है:
- पब/सब — तुरंत सभी ग्राहकों तक प्रसारण, लेकिन अस्थायी। समर्पित सब्सक्राइबर कनेक्शन का उपयोग करें। कैश अमान्यकरण और ऐसे लाइव नोटिफ़िकेशन के लिए बढ़िया, जहाँ संदेश छूट जाने से कोई नुकसान न हो।
- स्ट्रीम्स — टिकाऊ और दोबारा चलाए जा सकने वाला लॉग। उपभोक्ता समूह ठीक-एक-में-से-N डिलीवरी देते हैं;
XACKके साथ Pending Entries List औरXAUTOCLAIMक्रैश से उबरते हुए कम-से-कम-एक-बार प्रोसेसिंग देते हैं।MAXLEN ~से वृद्धि सीमित करें। - दर सीमित करना — ऐप कोड में पहले-जाँचें-फिर-सेट न करें; इससे रेस की स्थिति बनती है। निश्चित विंडो के लिए परमाण्विक
INCR+EXPIREया सुचारु उछाल के लिएEVALLua टोकन बकेट का उपयोग करें। इसे मिडलवेयर में लपेटें,429के साथRetry-Afterलौटाएँ और सोच-समझकर तय करें कि विफल होने पर अनुमति देनी है या रोकना है।
एकीकृत सिद्धांत यह है: एक ही राउंड-ट्रिप में परमाण्विक काम Redis से करवाइए और अपने उपयोग-प्रसंग के लिए आवश्यक टिकाऊपन की गारंटी के आधार पर प्रिमिटिव चुनिए।
एआई शिक्षक के साथ JavaScript सीखें — निःशुल्क
अपने ब्राउज़र में वास्तविक कोड लिखें और चलाएँ, चौबीसों घंटे एआई शिक्षक से तुरंत सहायता पाएँ, और वेब या ऐप पर वहीं से शुरू करें जहाँ आपने छोड़ा था।
- पाठ्यक्रम
- 22
- पाठ
- 92
अक्सर पूछे जाने वाले प्रश्न
क्या “Redis से Pub/Sub, स्ट्रीम और दर सीमांकन” पाठ निःशुल्क है?
हाँ—“Redis से Pub/Sub, स्ट्रीम और दर सीमांकन” का पूरा पाठ यहाँ वेब पर निःशुल्क पढ़ा जा सकता है। इंटरैक्टिव अभ्यास (अंतर्निहित कोड संपादक और 24/7 एआई ट्यूटर) करने और Node.js बैकएंड विकास बूटकैंप पाठ्यक्रम का बाकी हिस्सा अनलॉक करने के लिए CoddyKit PRO लें। Node.js बैकएंड विकास बूटकैंप पाठ्यक्रम में कुल 4 पाठ शामिल हैं।
“Redis से Pub/Sub, स्ट्रीम और दर सीमांकन” में मैं क्या सीखूँगा?
इवेंट प्रसारित कीजिए, टिकाऊ Redis Streams बनाइए और टोकन-बकेट दर सीमक परमाणु रूप से लागू कीजिए। आप ब्राउज़र में सीधे चलाए जाने वाले व्यावहारिक कोड के साथ Node.js बैकएंड विकास बूटकैंप का अभ्यास करते हैं, और पाठ पूरा करते समय 24/7 एआई ट्यूटर आपके प्रश्नों के उत्तर देता है।
क्या Node.js बैकएंड विकास बूटकैंप शुरू करने के लिए मुझे किसी अनुभव की आवश्यकता है?
पहले के अनुभव की आवश्यकता नहीं है। CoddyKit पर Node.js बैकएंड विकास बूटकैंप शुरुआती से लेकर उन्नत शिक्षार्थियों तक सभी के लिए व्यवस्थित किया गया है, इसलिए आप यहीं से या शुरुआत से सीखना शुरू कर सकते हैं और अपनी गति से आगे बढ़ सकते हैं। यह 4 में से 3वाँ पाठ है।
“Redis से Pub/Sub, स्ट्रीम और दर सीमांकन” पाठ पूरा करने में कितना समय लगता है?
CoddyKit का अधिकांश पाठ लगभग 5–10 मिनट में पूरा हो जाता है। हर पाठ छोटा और संवादात्मक है, इसलिए आप लगातार प्रगति करते हैं और वेब या ऐप पर वहीं से सीखना जारी रख सकते हैं जहाँ आपने छोड़ा था।
क्या मैं इस Node.js बैकएंड विकास बूटकैंप पाठ में कोड लिख और चला सकता हूँ?
हाँ। हर Node.js बैकएंड विकास बूटकैंप पाठ में एक अंतर्निर्मित कोड संपादक शामिल है, जिससे आप सीधे अपने ब्राउज़र में वास्तविक कोड लिख और चला सकते हैं और तुरंत एआई प्रतिक्रिया पा सकते हैं—स्थानीय सेटअप की आवश्यकता नहीं है।
इस पाठ्यक्रम के सभी पाठ
- कैश-आसाइड, राइट-थ्रू और TTL रणनीतियाँ
- वितरित लॉक और Redlock एल्गोरिदम
- Redis से Pub/Sub, स्ट्रीम और दर सीमांकन
- कैश स्टैम्पीड और थंडरिंग हर्ड रोकना