Node.js: التدفقات في Node.js
آخر تحديث: 2026-08-26
كُلف تشارلي بمهمة عاجلة: تحليل 2 جيجابايت من سجلات الوصول على الخادم. وما إن قام بتحميلها باستخدام fs.readFile، حتى نفدت الذاكرة على الفور — مما تسبب في تعطل النظام بسبب نفاد الذاكرة (OOM). قال قائد فريقه: «مهما كان حجم الملف، لا تحاول استيعابه دفعة واحدة — استخدم Stream لمعالجته شيئًا فشيئًا!» فانتقل تشارلي إلى createReadStream لمعالجة السجلات سطرًا سطرًا، وظل استهلاك الذاكرة مستقرًا عند 50 ميغابايت. ومنذ ذلك الحين، أصبح يقول لكل من يقابله: «Stream هو السكين السويسري المتعدد الاستخدامات لمعالجة البيانات الضخمة في Node.js.»
1. المفاهيم الأساسية للتيارات
Stream هي واجهة مجردة في Node.js مخصصة لمعالجة البيانات المتدفقة، حيث يتم نقل البيانات على شكل أجزاء، مثل تيار الماء، بدلاً من تحميلها إلى الذاكرة دفعة واحدة.
- المفهوم الأساسي: معالجة البيانات على شكل مجموعات، مما يتيح بدء العمليات دون انتظار توفر جميع البيانات
- علاقة التوريث: جميع التدفقات هي مثيلات لـ
EventEmitterوتستخدم الأحداث للإبلاغ عن تغييرات الحالة. - حالات الاستخدام: قراءة وكتابة الملفات الكبيرة، ونقل البيانات عبر الشبكة، وضغط البيانات وتحويلها
- الوحدات المدمجة:
fs،zlib،crypto،http، وغيرها، توفر جميعها واجهات stream. - الكائنات العامة: توفر الوحدة النمطية
streamالفئات الأساسيةReadableوWritableوDuplexوTransform
(1) أربعة أنواع من التدفقات
| النوع | الوصف | السيناريوهات النموذجية | المدخلات | المخرجات |
|---|---|---|---|---|
| قابل للقراءة | تيار قابل للقراءة، مصدر البيانات | fs.createReadStream، process.stdin |
لا شيء | نعم |
| قابل للكتابة | تيار قابل للكتابة، نقطة نهاية البيانات | fs.createWriteStream، process.stdout |
نعم | لا |
| ثنائي الاتجاه | ثنائي الاتجاه، قراءة/كتابة (مستقل) | net.Socket، tls.Socket |
نعم | نعم |
| تحويل | تيار التحويل؛ الناتج هو تحويل للمدخلات | zlib.createGzip، crypto.createCipheriv |
نعم | نعم |
(2) نوعان من التدفق
يتميز دفق «Readable» بوضعين للتشغيل:
| ميزة | وضع التشغيل | وضع الإيقاف المؤقت |
|---|---|---|
| استرجاع البيانات | الدفع التلقائي؛ الاستهلاك عبر حدث data |
يجب استدعاء read() يدويًّا لإجراء عملية الاسترجاع |
| طريقة التشغيل | إضافة مستمع data / استدعاء pipe() / استدعاء resume() |
الحالة الافتراضية الأولية / استدعاء pause() |
| التعامل مع الضغط العكسي | pipe() معالجة تلقائية؛ يتطلب التشغيل اليدوي مراقبة drain |
التحكم في وتيرة التشغيل من قبل المستهلك، ضغط عكسي طبيعي |
| حالات الاستخدام | مصادر بيانات مستمرة وعالية الإنتاجية | تتطلب تحكماً دقيقاً في توقيت القراءة |
2. شرح مفصل لأحداث البث المباشر
تستند جميع التدفقات إلى EventEmitter وتستخدم آلية قائمة على الأحداث (event) لتمرير البيانات والحالة.
(1) الأحداث العامة
| الحدث | شرط التشغيل | المسار المطبق | معلمات الاستدعاء |
|---|---|---|---|
data |
تم استلام جزء جديد | قابل للقراءة | chunk |
end |
قراءة البيانات | قابلة للقراءة | لا شيء |
error |
خطأ | جميع البث المباشر | Error |
close |
إغلاق الموارد ذات المستوى المنخفض | جميع البث المباشر | لا شيء |
finish |
end() بعد مسح جميع البيانات |
قابل للكتابة | لا شيء |
▶ مثال: الاستماع إلى أحداث تيار «Readable»
const fs = require('fs');
const rs = fs.createReadStream('./access.log', { highWaterMark: 64 * 1024 });
rs.on('data', (chunk) => {
console.log(`Received ${chunk.length} Byte`);
});
rs.on('end', () => {
console.log('Finished reading');
});
rs.on('error', (err) => {
console.error('Output error:', err.message);
});
▶ مثال: الاستماع إلى أحداث تيار قابل للكتابة
const ws = fs.createWriteStream('./output.txt');
ws.on('finish', () => {
console.log('All data has been written');
});
ws.on('error', (err) => {
console.error('Write error:', err.message);
});
ws.write('Hello Stream');
ws.end();
dataيتم تشغيل الأحداث تلقائيًا في وضع التدفق- الفرق بين
endوfinish: يشيرendإلى اكتمال عملية القراءة، بينما يشيرfinishإلى اكتمال عملية الكتابة. errorيجب الاستماع إلى الأحداث؛ وإلا فإن الاستثناءات التي لم يتم التقاطها ستؤدي إلى تعطل العملية.closeلا تؤدي جميع التدفقات إلى تشغيل الحدث؛ فالأمر يعتمد على طريقة التنفيذ الأساسية.
3. التسلسل باستخدام pipe()
pipe() هي أقوى طريقة في Stream؛ فهي تربط تلقائيًا مخرج تيار قابل للقراءة بمدخل تيار قابل للكتابة، وتتعامل تلقائيًا مع الضغط العكسي.
▶ مثال:(1) الاستخدام الأساسي لـ "pipe"
readable.pipe(writable);
pipe() تُرجع الدفق المستهدف، بحيث يمكن ربطها بسلسلة من العمليات:
readable.pipe(transform1).pipe(transform2).pipe(writable);
▶ مثال: نسخ الملفات
const fs = require('fs');
fs.createReadStream('./source.txt')
.pipe(fs.createWriteStream('./dest.txt'));
▶ مثال: دفق استجابات HTTP
const http = require('http');
const fs = require('fs');
http.createServer((req, res) => {
res.writeHead(200, { 'Content-Type': 'text/plain' });
fs.createReadStream('./large.txt').pipe(res);
}).listen(3000);
▶ مثال:(2) مخطط تدفق البيانات لأنابيب التيار
flowchart LR
A[Readable<br/>Data Source] -->|chunk| B[Transform<br/>Data Transformation]
B -->|chunk| C[Writable<br/>Data Endpoint]
C -.->|backpressure| A
B -.->|backpressure| A
style A fill:#4CAF50,color:#fff
style B fill:#FF9800,color:#fff
style C fill:#2196F3,color:#fff
pipe()يدير معدل تدفق البيانات تلقائيًا؛ حيث تتوقف عملية القراءة مؤقتًا عندما تكون عملية الكتابة مشغولةpipe()لا تتم معالجة الأخطاء؛ يجب مراقبة الأحداث بشكل منفصل لكل دفقerror- يتولى استخدام
stream.pipeline()تلقائيًا عمليات التنظيف ونقل الأخطاء، مما يجعله أكثر أمانًا منpipe()
4. تدفقات نظام الملفات وعمليات الملفات
توفر الوحدة النمطية fs تدفقات القراءة والكتابة المتعلقة بنظام الملفات.
(1) createReadStream / createWriteStream
| الخيار | الوصف | القيمة الافتراضية |
|---|---|---|
highWaterMark |
حجم المخزن المؤقت (بايت) | القراءة: 64 كيلوبايت / الكتابة: 16 كيلوبايت |
encoding |
الترميز | null (المخزن المؤقت) |
start |
موضع البايت الأول | 0 |
end |
موضع البايت الأخير (بما في ذلك) | اللانهاية |
flags |
علامات فتح الملف | قابل للقراءة: r / قابل للكتابة: w |
▶ مثال: قراءة جزء من ملف ضمن نطاق محدد
const fs = require('fs');
const rs = fs.createReadStream('./big.bin', {
start: 100,
end: 199,
highWaterMark: 32
});
rs.on('data', (chunk) => {
console.log(chunk.length);
});
(2) مقارنة بين readFile و createReadStream
| المقارنة | readFile | createReadStream |
|---|---|---|
| استخدام الذاكرة | جميع الملفات المحملة في الذاكرة | يستخدم حجم highWaterMark فقط |
| تأخير بدء التشغيل | الانتظار حتى يتم قراءة جميع الملفات قبل استدعاء دالة الاستدعاء | إرجاع الدفق فورًا ومعالجته أثناء القراءة |
| أحجام الملفات المناسبة | الملفات الصغيرة (<10 ميغابايت) | الملفات الكبيرة أو البيانات المتواصلة |
| معالجة الأخطاء | استرداد خطأ في دالة الاستدعاء | الاستماع إلى أحداث error |
| يمكن إيقافه مؤقتًا/استئنافه | غير مدعوم | pause() / resume() |
▶ مثال: معالجة الملفات الكبيرة كتلةً كتلةً
const fs = require('fs');
let totalBytes = 0;
const rs = fs.createReadStream('./2gb.log');
rs.on('data', (chunk) => {
totalBytes += chunk.length;
});
rs.on('end', () => {
console.log(`Total ${totalBytes} Byte`);
});
- يجب معالجة الملفات الكبيرة باستخدام التدفقات لتجنب أخطاء OOM
highWaterMarkكلما كانت الذاكرة أصغر، زادت كفاءتها في استهلاك الطاقة، لكن زادت أيضًا وتيرة استدعاءات النظام.- تتيح الخيارات
start/endقراءة الملفات على شكل أجزاء
5. آلية الضغط العكسي
يُعد الضغط العكسي عنصراً أساسياً في التحكم في التدفق: فعندما يتعذر على جانب الكتابة مواكبة معدل الدفع من جانب القراءة، يتم تطبيق «الضغط العكسي» على جانب الكتابة لإيقاف جانب القراءة مؤقتاً ومنع تراكم البيانات في الذاكرة.
(1) تعالج الدالة pipe() الضغط العكسي تلقائيًا
عند استخدام pipe()، يتم إدارة الضغط العكسي تلقائيًا داخليًّا بواسطة Node.js، دون الحاجة إلى أي تدخل يدوي.
(2) التحكم اليدوي في الضغط العكسي
عندما لا يتم استخدام pipe()، يجب عليك حساب قيمة الإرجاع لـ write() يدويًّا:
const fs = require('fs');
const rs = fs.createReadStream('./source.txt');
const ws = fs.createWriteStream('./dest.txt');
rs.on('data', (chunk) => {
const canContinue = ws.write(chunk);
if (!canContinue) {
rs.pause();
ws.once('drain', () => {
rs.resume();
});
}
});
rs.on('end', () => {
ws.end();
});
▶ مثال: ملاحظة تأثير الضغط الخلفي
const rs = fs.createReadStream('./big.log', { highWaterMark: 1024 });
const ws = fs.createWriteStream('./out.log', { highWaterMark: 512 });
let paused = 0;
rs.on('data', (chunk) => {
const ok = ws.write(chunk);
if (!ok) {
paused++;
rs.pause();
ws.once('drain', () => rs.resume());
}
});
rs.on('end', () => {
console.log(`Backpressure triggered ${paused} times`);
ws.end();
});
write()الإرجاعfalseيشير إلى أن المخزن المؤقت الداخلي ممتلئ ويجب إيقاف القراءة مؤقتًاdrainيشير هذا الحدث إلى أنه تم مسح المخزن المؤقت ويمكن استئناف الكتابة.- سيؤدي تجاهل الضغط العكسي إلى استمرار زيادة استخدام الذاكرة، مما سيؤدي في النهاية إلى حدوث خطأ OOM
- يُوصى باستخدام
pipeline()أكثر منpipe()؛ فهو يتولى تلقائيًا معالجة انتشار الأخطاء وتنظيف الموارد.
6. تيار التحويل
يُعد تيار «Transform» فئة فرعية من «Duplex»؛ حيث يمثل مخرجه النتيجة المحولة للمدخلات، ويُستخدم عادةً لضغط البيانات وتشفيرها وتحويل تنسيقاتها.
(1) الفرق بين «Transform» و«Duplex»
| عنصر المقارنة | دوبلكس | تحويل |
|---|---|---|
| العلاقة بين المدخلات والمخرجات | مستقلة؛ لا تؤثر إحداهما على الأخرى | يتم توليد المخرجات عن طريق تحويل المدخلات |
| الطرق المطلوبة | _read() + _write() |
_transform() |
| الاستخدامات الشائعة | مآخذ الشبكة | الضغط والتشفير وتحويل البيانات |
| المخزن المؤقت الداخلي | واحد للقراءة، وآخر للكتابة | الحالة الوسيطة للتحويل |
▶ مثال: دفق تحويل مخصص (تحويل حالة الأحرف)
const { Transform } = require('stream');
const upper = new Transform({
transform(chunk, encoding, callback) {
callback(null, chunk.toString().toUpperCase());
}
});
process.stdin.pipe(upper).pipe(process.stdout);
▶ مثال: الملفات المضغوطة باستخدام zlib
const fs = require('fs');
const zlib = require('zlib');
fs.createReadStream('./access.log')
.pipe(zlib.createGzip())
.pipe(fs.createWriteStream('./access.log.gz'));
▶ مثال: فك ضغط ملف zlib
const fs = require('fs');
const zlib = require('zlib');
fs.createReadStream('./access.log.gz')
.pipe(zlib.createGunzip())
.pipe(fs.createWriteStream('./access_restored.log'));
- في
_transform(chunk, encoding, callback)، يقومcallback(null, data)بإرسال نتائج التحويل callback(err)خطأ مقبولzlib.createGzip()/createGunzip()هما تيارات «Transform» المدمجة الأكثر استخدامًا- يمكنك أيضًا تطبيق معالجة
_flush(callback)على البيانات المتبقية عند تخصيص دفق «Transform».
7. إيقاف البث مؤقتًا واستئنافه
يكون دفق «Readable» في وضع الإيقاف المؤقت بشكل افتراضي، ويمكن التبديل بين أوضاعه باستخدام pause() وresume().
▶ مثال: عناصر التحكم في الإيقاف المؤقت والاستئناف
const fs = require('fs');
const rs = fs.createReadStream('./big.log');
let count = 0;
rs.on('data', (chunk) => {
count++;
if (count % 10 === 0) {
rs.pause();
console.log(`Processed ${count} chunks, pause 1 second`);
setTimeout(() => rs.resume(), 1000);
}
});
▶ مثال: قراءة ملف سطراً سطراً (readline)
const fs = require('fs');
const readline = require('readline');
const rl = readline.createInterface({
input: fs.createReadStream('./access.log'),
crlfDelay: Infinity
});
rl.on('line', (line) => {
if (line.includes('ERROR')) {
console.log(line);
}
});
rl.on('close', () => {
console.log('The file has been read.');
});
readline.createInterfaceيلف تيار «Readable» كواجهة للقراءة سطراً سطراًcrlfDelay: Infinityيتعامل بشكل صحيح مع فواصل الأسطر في\r\n- بعد
pause()، لا يتم تشغيل الحدثdataحتى يتم استدعاءresume(). - سيقوم
pipe()بتحويل البث إلى الوضع المباشر
8. مثال شامل: مسار معالجة السجلات
إنشاء مسار معالجة بيانات متكامل: قراءة ملفات السجلات الكبيرة → تحليلها سطراً سطراً → تصفية الأسطر التي تحتوي على أخطاء → ضغط الناتج.
const fs = require('fs');
const zlib = require('zlib');
const { Transform } = require('stream');
const filterError = new Transform({
transform(chunk, encoding, callback) {
const lines = chunk.toString().split('\n');
const errors = lines.filter(l => l.includes('ERROR')).join('\n');
callback(null, errors ? errors + '\n' : '');
}
});
const src = fs.createReadStream('./app.log');
const dest = fs.createWriteStream('./errors.log.gz');
src
.pipe(filterError)
.pipe(zlib.createGzip())
.pipe(dest);
dest.on('finish', () => {
console.log('The error log has been written in compressed form.');
});
src.on('error', (err) => {
console.error('Read failed:', err.message);
});
- كل مرحلة من مراحل مسار المعالجة تمثل تيارًا؛ حيث تتدفق البيانات عبره على شكل أجزاء، مما يؤدي إلى انخفاض شديد في استهلاك الذاكرة.
pipe()يتعامل تلقائيًا مع الضغط العكسي بين التدفقات المتجاورة- يجب مراقبة الأخطاء بشكل منفصل في كل دفق، أو معالجتها بشكل موحد باستخدام
pipeline() - تحويل الدفق
filterErrorبحيث يتم الاحتفاظ فقط بالصفوف التي تحتوي علىERROR
❓ أسئلة شائعة
pipe؟pipe مخرج تيار قابل للقراءة بمدخل تيار قابل للكتابة، وتقوم تلقائيًا بإدارة تدفق البيانات والضغط العكسي.- متى ينبغي استخدام Stream؟ استخدم Stream عندما يتجاوز حجم البيانات الذاكرة المتاحة، أو عندما تحتاج إلى معالجة البيانات أثناء قراءتها، أو عند التعامل مع مصادر البيانات في الوقت الفعلي؛ أما بالنسبة للملفات الصغيرة، فمن الأسهل استخدام
readFile. - هل تتعامل
pipe()تلقائيًا مع الضغط العكسي؟ نعم، تقومpipe()داخليًّا بفحص قيمة الإرجاع الواردة من الطرف القابل للكتابةwrite()وتستدعي تلقائيًّاpause()/resume()لإدارة معدل التدفق. - كيف يمكنني قراءة ملف سطراً سطراً؟ استخدم
readline.createInterface({ input: createReadStream(path) })وانتظر حدوث الحدثlineلاسترداد البيانات سطراً سطراً. - ما الفرق بين Transform و Duplex؟ في Duplex، تكون عمليات القراءة والكتابة مستقلة عن بعضها وغير مرتبطة ببعض؛ أما في Transform، فيتم إنشاء المخرجات عن طريق تحويل المدخلات، مما يتطلب فقط تنفيذ طريقة
_transform().** - كيف ينبغي التعامل مع أخطاء البث؟ استمع إلى أحداث
errorفي كل بث، أو استخدمstream.pipeline()لنقل الأخطاء تلقائيًا وتنظيف الموارد، مما يمنع تعطل العمليات الناتج عن الاستثناءات التي لم يتم التقاطها. - ما هي القيمة المناسبة لـ
highWaterMark؟ القيمة الافتراضية البالغة 64 كيلوبايت مناسبة لمعظم الحالات؛ ومعالجة الكتل الكبيرة من البيانات الثنائية، يمكنك زيادتها إلى 256 كيلوبايت، وعندما تكون الذاكرة محدودة، يمكنك تخفيضها إلى 16 كيلوبايت. - أيهما يجب أن تختار:
pipe()أمpipeline()؟ نوصي باستخدامstream.pipeline()، لأنه يقوم تلقائيًا بنقل الأخطاء وتنظيف الموارد؛ أماpipe()فلا يتعامل مع الأخطاء ولا يغلق التدفقات تلقائيًا.
📖 ملخص
- المفاهيم الأساسية واستخدامات أساسيات التدفق
- شرح مفصل للمفاهيم الأساسية لأحداث الدفق وكيفية استخدامها
- المفاهيم الأساسية واستخدامات «Chainable Connections» مع
pipe() - المفاهيم الأساسية واستخدامات تدفقات الملفات وعمليات الملفات في نظام الملفات
- المفاهيم الأساسية لآلية الضغط المضاد وكيفية استخدامها
- المفاهيم الأساسية لـ «Transform Stream» وكيفية استخدامه
- المفاهيم الأساسية واستخدامات إيقاف التدفقات مؤقتًا واستئنافها
- مثال شامل: المفاهيم الأساسية واستخدامات مسار معالجة السجلات
📝 تمارين
- أكمل جميع أمثلة الأكواد الواردة في هذا الدرس وتأكد من أن كل منها يعمل بشكل صحيح.
- قم بتعديل المثال الشامل وأضف الإضافات الخاصة بك
- راجع الوثائق الرسمية، وحدد واجهة برمجة تطبيقات (API) واحدة أو اثنتين لم يتم تناولهما في هذا الدرس، واكتب كود اختبار لهما.
- التأمل: كيف ستطبق ما تعلمته في هذا الدرس على مشروع في الواقع العملي؟
- حاول أن تجمع بين ما تعلمته في هذا الدرس والمواد التي درستها في الدروس السابقة لإنشاء مشروع صغير.