Node.js: التدفقات في Node.js

آخر تحديث: 2026-08-26

كُلف تشارلي بمهمة عاجلة: تحليل 2 جيجابايت من سجلات الوصول على الخادم. وما إن قام بتحميلها باستخدام fs.readFile، حتى نفدت الذاكرة على الفور — مما تسبب في تعطل النظام بسبب نفاد الذاكرة (OOM). قال قائد فريقه: «مهما كان حجم الملف، لا تحاول استيعابه دفعة واحدة — استخدم Stream لمعالجته شيئًا فشيئًا!» فانتقل تشارلي إلى createReadStream لمعالجة السجلات سطرًا سطرًا، وظل استهلاك الذاكرة مستقرًا عند 50 ميغابايت. ومنذ ذلك الحين، أصبح يقول لكل من يقابله: «Stream هو السكين السويسري المتعدد الاستخدامات لمعالجة البيانات الضخمة في Node.js.»

1. المفاهيم الأساسية للتيارات

Stream هي واجهة مجردة في Node.js مخصصة لمعالجة البيانات المتدفقة، حيث يتم نقل البيانات على شكل أجزاء، مثل تيار الماء، بدلاً من تحميلها إلى الذاكرة دفعة واحدة.

(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»

JAVASCRIPT
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);
});
▶ جرّب الكود

▶ مثال: الاستماع إلى أحداث تيار قابل للكتابة

JAVASCRIPT
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();
▶ جرّب الكود

3. التسلسل باستخدام pipe()

pipe() هي أقوى طريقة في Stream؛ فهي تربط تلقائيًا مخرج تيار قابل للقراءة بمدخل تيار قابل للكتابة، وتتعامل تلقائيًا مع الضغط العكسي.

▶ مثال:(1) الاستخدام الأساسي لـ "pipe"

JAVASCRIPT
readable.pipe(writable);
▶ جرّب الكود

pipe() تُرجع الدفق المستهدف، بحيث يمكن ربطها بسلسلة من العمليات:

JAVASCRIPT
readable.pipe(transform1).pipe(transform2).pipe(writable);

▶ مثال: نسخ الملفات

JAVASCRIPT
const fs = require('fs');
fs.createReadStream('./source.txt')
  .pipe(fs.createWriteStream('./dest.txt'));
▶ جرّب الكود

▶ مثال: دفق استجابات HTTP

JAVASCRIPT
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) مخطط تدفق البيانات لأنابيب التيار

100%
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


4. تدفقات نظام الملفات وعمليات الملفات

توفر الوحدة النمطية fs تدفقات القراءة والكتابة المتعلقة بنظام الملفات.

(1) createReadStream / createWriteStream

الخيار الوصف القيمة الافتراضية
highWaterMark حجم المخزن المؤقت (بايت) القراءة: 64 كيلوبايت / الكتابة: 16 كيلوبايت
encoding الترميز null (المخزن المؤقت)
start موضع البايت الأول 0
end موضع البايت الأخير (بما في ذلك) اللانهاية
flags علامات فتح الملف قابل للقراءة: r / قابل للكتابة: w

▶ مثال: قراءة جزء من ملف ضمن نطاق محدد

JAVASCRIPT
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()

▶ مثال: معالجة الملفات الكبيرة كتلةً كتلةً

JAVASCRIPT
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`);
});
▶ جرّب الكود

5. آلية الضغط العكسي

يُعد الضغط العكسي عنصراً أساسياً في التحكم في التدفق: فعندما يتعذر على جانب الكتابة مواكبة معدل الدفع من جانب القراءة، يتم تطبيق «الضغط العكسي» على جانب الكتابة لإيقاف جانب القراءة مؤقتاً ومنع تراكم البيانات في الذاكرة.

(1) تعالج الدالة pipe() الضغط العكسي تلقائيًا

عند استخدام pipe()، يتم إدارة الضغط العكسي تلقائيًا داخليًّا بواسطة Node.js، دون الحاجة إلى أي تدخل يدوي.

(2) التحكم اليدوي في الضغط العكسي

عندما لا يتم استخدام pipe()، يجب عليك حساب قيمة الإرجاع لـ write() يدويًّا:

JAVASCRIPT
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();
});

▶ مثال: ملاحظة تأثير الضغط الخلفي

JAVASCRIPT
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();
});
▶ جرّب الكود

6. تيار التحويل

يُعد تيار «Transform» فئة فرعية من «Duplex»؛ حيث يمثل مخرجه النتيجة المحولة للمدخلات، ويُستخدم عادةً لضغط البيانات وتشفيرها وتحويل تنسيقاتها.

(1) الفرق بين «Transform» و«Duplex»

عنصر المقارنة دوبلكس تحويل
العلاقة بين المدخلات والمخرجات مستقلة؛ لا تؤثر إحداهما على الأخرى يتم توليد المخرجات عن طريق تحويل المدخلات
الطرق المطلوبة _read() + _write() _transform()
الاستخدامات الشائعة مآخذ الشبكة الضغط والتشفير وتحويل البيانات
المخزن المؤقت الداخلي واحد للقراءة، وآخر للكتابة الحالة الوسيطة للتحويل

▶ مثال: دفق تحويل مخصص (تحويل حالة الأحرف)

JAVASCRIPT
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

JAVASCRIPT
const fs = require('fs');
const zlib = require('zlib');

fs.createReadStream('./access.log')
  .pipe(zlib.createGzip())
  .pipe(fs.createWriteStream('./access.log.gz'));
▶ جرّب الكود

▶ مثال: فك ضغط ملف zlib

JAVASCRIPT
const fs = require('fs');
const zlib = require('zlib');

fs.createReadStream('./access.log.gz')
  .pipe(zlib.createGunzip())
  .pipe(fs.createWriteStream('./access_restored.log'));
▶ جرّب الكود

7. إيقاف البث مؤقتًا واستئنافه

يكون دفق «Readable» في وضع الإيقاف المؤقت بشكل افتراضي، ويمكن التبديل بين أوضاعه باستخدام pause() وresume().

▶ مثال: عناصر التحكم في الإيقاف المؤقت والاستئناف

JAVASCRIPT
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)

JAVASCRIPT
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.');
});
▶ جرّب الكود

8. مثال شامل: مسار معالجة السجلات

إنشاء مسار معالجة بيانات متكامل: قراءة ملفات السجلات الكبيرة → تحليلها سطراً سطراً → تصفية الأسطر التي تحتوي على أخطاء → ضغط الناتج.

JAVASCRIPT
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؟
ج تربط طريقة pipe مخرج تيار قابل للقراءة بمدخل تيار قابل للكتابة، وتقوم تلقائيًا بإدارة تدفق البيانات والضغط العكسي.
س ما المقصود بالضغط العكسي؟
ج عندما يكون معدل إخراج البيانات في تيار الكتابة أبطأ من معدل إدخال البيانات في تيار القراءة، يرسل تيار الكتابة إشارة إلى تيار القراءة لإيقاف القراءة مؤقتًا؛ وهذه هي آلية الضغط العكسي.

📖 ملخص


📝 تمارين

  1. أكمل جميع أمثلة الأكواد الواردة في هذا الدرس وتأكد من أن كل منها يعمل بشكل صحيح.
  2. قم بتعديل المثال الشامل وأضف الإضافات الخاصة بك
  3. راجع الوثائق الرسمية، وحدد واجهة برمجة تطبيقات (API) واحدة أو اثنتين لم يتم تناولهما في هذا الدرس، واكتب كود اختبار لهما.
  4. التأمل: كيف ستطبق ما تعلمته في هذا الدرس على مشروع في الواقع العملي؟
  5. حاول أن تجمع بين ما تعلمته في هذا الدرس والمواد التي درستها في الدروس السابقة لإنشاء مشروع صغير.
Web-Tutorial.com

فريق Web-Tutorial التقني

منصة دروس برمجية يديرها عدة مطورين. كل درس يتم كتابته ومراجعته بواسطة مطورين متخصصين في المجال. نعمل على ضمان دقة وموثوقية المحتوى — إذا لاحظت أي مشكلة، فيرجى إخبارنا.

100%