Stream moduli


ULASHISH

Streamlar nima?

Node.js-da streamlar bir vaqtning o‘zida to‘liq mavjud bo‘lmasligi va xotiraga sig‘ishi shart bo‘lmagan ma’lumotlar to‘plamidir.

Ularni ma’lumotlarni bir joydan ikkinchi joyga ko‘chiradigan konveyer bantlari deb tasavvur qiling, bu sizga butun ma’lumotlar to‘plamini kutish o‘rniga har bir bo‘lak bilan ishlashga imkon beradi.

Streams Node.js’ning eng kuchli xususiyatlaridan biri bo‘lib, ular quyidagilarda keng qo‘llaniladi:

  • File system operations (reading/writing files)
  • HTTP so‘rovlari va javoblari
  • Ma’lumotlarni siqish va ochish
  • Ma’lumotlar bazasi operatsiyalari
  • Haqiqiy vaqtda ma’lumotlarni qayta ishlash

Streams bilan ishlashni boshlash

Streamlar Node.js’dagi ma’lumotlar bilan samarali ishlash uchun asosiy tushunchalardan biridir.

Ular sizga hamma narsani bir vaqtning o‘zida xotiraga yuklashdan ko‘ra, mavjud bo‘lganda ma’lumotlarni qismlarga ajratish imkonini beradi.

Asosiy stream misoli

const fs = require('fs');

// Create a readable stream from a file
const readableStream = fs.createReadStream('input.txt', 'utf8');
// Create a writable stream to a file
const writableStream = fs.createWriteStream('output.txt');

// Pipe the data from readable to writable stream
readableStream.pipe(writableStream);

// Handle completion and errors
writableStream.on('finish', () => {
  console.log('File copy completed!');
});

readableStream.on('error', (err) => {
  console.error('Error reading file:', err);
});

writableStream.on('error', (err) => {
  console.error('Error writing file:', err);
});
Misolni ishga tushirish »

Streams nima uchun ishlatiladi?

Oqimlardan foydalanishning bir qancha afzalliklari bor:

  • Xotira samaradorligi: Katta hajmdagi fayllarni xotiraga to‘liq yuklamasdan qayta ishlash
  • Vaqt samaradorligi: Barcha ma’lumotlarni kutish o‘rniga, ma’lumotlarga ega bo‘lishingiz bilanoq uni qayta ishlashni boshlang.
  • Birlashtirish imkoniyati: Oqimlarni ulash orqali kuchli ma’lumotlar payplar (pipes)ini yarating
  • Yaxshiroq foydalanuvchi tajribasi: Ma’lumotlar mavjud bo‘lganda foydalanuvchilarga yetkazing (masalan, video striming)

512 MB operativ xotiraga ega serverda 1 Gb faylni o‘qishni tasavvur qiling:

  • Oqimlarsiz: Butun faylni xotiraga yuklashga urinish jarayoni buziladi
  • Streamlar bilan: Siz faylni kichik bo‘laklarga (masalan, bir vaqtning o‘zida 64KB) qayta ishlaysiz.

Asosiy stream turlari

Node.js to‘rtta asosiy turdagi oqimlarni taqdim etadi, ularning har biri ma’lumotlar bilan ishlashda muayyan maqsadga xizmat qiladi:

Stream turi Tavsif Umumiy misollar
O‘qish mumkin Ma’lumotlarni o‘qish mumkin bo‘lgan oqimlar (ma’lumot manbai) fs.createReadStream(), HTTP responses, process.stdin
Yozilishi mumkin Ma’lumotlarni yozish mumkin bo‘lgan oqimlar (ma’lumot qabul qiluvchisi) fs.createWriteStream(), HTTP requests, process.stdout
Dupleks O‘qilishi va yozilishi mumkin bo‘lgan streamlar TCP soketlari, Zlib oqimlari
O‘zgartirish Ma’lumotlarni yozilishi va o‘qilishi bilan o‘zgartirishi yoki o‘zgartirishi mumkin bo‘lgan ikki tomonlama streamlar Zlib oqimlari, kripto oqimlari

Eslatma: Node.js’dagi barcha streamlar EventEmitter’ning namunalari bo‘lib, ular tinglash va boshqarish mumkin bo‘lgan hodisalarni chiqaradi.



O‘qilishi mumkin bo‘lgan streamlar

O‘qiladigan streamlar manbadan ma’lumotlarni o‘qish imkonini beradi. Bunga misollar kiradi:

  • Fayldan o‘qish
  • Mijozdagi HTTP javoblari
  • Serverdagi HTTP so‘rovlari
  • process.stdin

O‘qiladigan stream yaratish

const fs = require('fs');

// Create a readable stream from a file
const readableStream = fs.createReadStream('myfile.txt', {
  encoding: 'utf8',
  highWaterMark: 64 * 1024 // 64KB chunks
});

// Events for readable streams
readableStream.on('data', (chunk) => {
  console.log(`Received ${chunk.length} bytes of data.`);
  console.log(chunk);
});

readableStream.on('end', () => {
  console.log('No more data to read.');
});

readableStream.on('error', (err) => {
  console.error('Error reading from stream:', err);
});
Misolni ishga tushirish »

O‘qish rejimlari

O‘qiladigan streamlar ikkita rejimdan birida ishlaydi:

  • Oqimli rejim: Ma’lumotlar manbadan o‘qiladi va voqealar yordamida iloji boricha tezroq ilovangizga taqdim etiladi.
  • To‘xtatilgan rejim: Oqimdagi ma’lumotlarni olish uchun stream.read() ni o‘zingiz aniq chaqirishingiz kerak.
const fs = require('fs');

// Paused mode example
const readableStream = fs.createReadStream('myfile.txt', {
  encoding: 'utf8',
  highWaterMark: 64 * 1024 // 64KB chunks
});

// Manually consume the stream using read()
readableStream.on('readable', () => {
  let chunk;
  while (null !== (chunk = readableStream.read())) {
    console.log(`Read ${chunk.length} bytes of data.`);
    console.log(chunk);
  } });

readableStream.on('end', () => {
  console.log('No more data to read.');
});
Misolni ishga tushirish »

Yoziladigan streamlar

Yoziladigan streamlar ma’lumotlarni belgilangan joyga yozish imkonini beradi. Bunga misollar kiradi:

  • Faylga yozish
  • Mijozdagi HTTP so‘rovlari
  • Serverdagi HTTP javoblari
  • process.stdout

Yoziladigan stream yaratish

const fs = require('fs');

// Create a writable stream to a file
const writableStream = fs.createWriteStream('output.txt');

// Write data to the stream
writableStream.write('Hello, ');
writableStream.write('World!');
writableStream.write('\nWriting to a stream is easy!');

// End the stream
writableStream.end();

// Events for writable streams
writableStream.on('finish', () => {
  console.log('All data has been written to the file.');
});

writableStream.on('error', (err) => {
  console.error('Error writing to stream:', err);
});
Misolni ishga tushirish »

Orqa bosim bilan ishlash

Oqimga yozishda, agar ma’lumotlar qayta ishlanishi mumkin bo‘lganidan tezroq yozilsa, orqa bosim paydo bo‘ladi.

write() usuli yozishni davom ettirish xavfsiz yoki yo‘qligini ko‘rsatuvchi mantiqiy qiymatni qaytaradi.

const fs = require('fs');

const writableStream = fs.createWriteStream('output.txt');

function writeData() {
  let i = 100;
  function write() {
    let ok = true;
    do {
      i--;
      if (i === 0) {
        // Last time, close the stream
        writableStream.write('Last chunk!\n');
        writableStream.end();
      } else {
        // Continue writing data
        const data = `Data chunk ${i}\n`;
        // Write and check if we should continue
        ok = writableStream.write(data);
      }
    }
    while (i > 0 && ok);

    if (i > 0) {
      // We need to wait for the drain event before writing more
      writableStream.once('drain', write);
    }
  }
  write();
}

writeData();
writableStream.on('finish', () => {
  console.log('All data written successfully.');
});
Misolni ishga tushirish »

Piping (oqimlarni ulash - pipe)

pipe() usuli o‘qilishi mumkin bo‘lgan oqimni yoziladigan oqimga ulaydi, ma’lumotlar oqimini avtomatik ravishda boshqaradi va orqa bosimni boshqaradi.

Bu oqimlarni iste’mol qilishning eng oson yo‘li.

const fs = require('fs');

// Create readable and writable streams
const readableStream = fs.createReadStream('source.txt');
const writableStream = fs.createWriteStream('destination.txt');

// Pipe the readable stream to the writable stream
readableStream.pipe(writableStream);

// Handle completion and errors
readableStream.on('error', (err) => {
  console.error('Read error:', err);
});

writableStream.on('error', (err) => {
  console.error('Write error:', err);
});

writableStream.on('finish', () => {
  console.log('File copy completed!');
});
Misolni ishga tushirish »

Zanjirli payplar (pipes)

pipe() yordamida bir nechta oqimlarni zanjirlashingiz mumkin.

Bu, ayniqsa, transform oqimlari bilan ishlashda foydalidir.

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

// Create a pipeline to read a file, compress it, and write to a new file
fs.createReadStream('source.txt')
  .pipe(zlib.createGzip()) // Compress the data
  .pipe(fs.createWriteStream('destination.txt.gz'))
  .on('finish', () => {
    console.log('File compressed successfully!');
  });
Misolni ishga tushirish »

Eslatma: pipe() usuli zanjirlanish imkonini beruvchi maqsad oqimini qaytaradi.


Dupleks va Transform oqimlari

Ikki tomonlama streamlar

Dupleks oqimlari ikki tomonlama pipeline kabi o‘qilishi va yozilishi mumkin.

TCP soketsi dupleks oqimning yaxshi namunasidir.

const net = require('net');

// Create a TCP server
const server = net.createServer((socket) => {
  // 'socket' is a duplex stream

  // Handle incoming data (readable side)
  socket.on('data', (data) => {
    console.log('Received:', data.toString());

    // Echo back (writable side)
    socket.write(`Echo: ${data}`);
  });

  socket.on('end', () => {
    console.log('Client disconnected');
  });
});

server.listen(8080, () => {
  console.log('Server listening on port 8080');
});

// To test, you can use a tool like netcat or telnet:
// $ nc localhost 8080
// or create a client:
/*
const client = net.connect({ port: 8080 }, () => {
  console.log('Connected to server');
  client.write('Hello from client!');
});

client.on('data', (data) => {
  console.log('Server says:', data.toString());
  client.end(); // Close the connection
});
*/

Oqimlarni aylantirish

Transformatsiya oqimlari ikki tomonlama streamlar bo‘lib, ular orqali ma’lumotlarni o‘zgartirish mumkin.

Ular payplar (pipes) ma’lumotlarni qayta ishlash uchun ideal.

const { Transform } = require('stream');
const fs = require('fs');

// Create a transform stream that converts text to uppercase
class UppercaseTransform extends Transform {
  _transform(chunk, encoding, callback) {
    // Transform the chunk to uppercase
    const upperChunk = chunk.toString().toUpperCase();
    // Push the transformed data
    this.push(upperChunk);
    // Signal that we're done with this chunk
    callback();
  }
}

// Create an instance of our transform stream
const uppercaseTransform = new UppercaseTransform();

// Create a readable stream from a file
const readableStream = fs.createReadStream('input.txt');

// Create a writable stream to a file
const writableStream = fs.createWriteStream('output-uppercase.txt');

// Pipe the data through our transform stream
readableStream
  .pipe(uppercaseTransform)
  .pipe(writableStream)
  .on('finish', () => {
    console.log('Transformation completed!');
  });
Misolni ishga tushirish »

Oqimli voqealar

Barcha streamlar EventEmitter misollari va bir nechta hodisalarni chiqaradi:

O‘qilishi mumkin bo‘lgan stream hodisalari

  • data : Oqimda o‘qish uchun mavjud ma’lumotlar mavjud bo‘lganda chiqariladi
  • end : Iste’mol qilinadigan ma’lumotlar qolmaganda chiqariladi
  • error : O‘qish paytida xatolik yuzaga kelsa, chiqariladi
  • close : Oqimning asosiy resursi yopilganda chiqariladi
  • readable : O‘qish uchun ma’lumotlar mavjud bo‘lganda chiqariladi

Yoziladigan stream hodisalari

  • drain : write() usuli false qaytarilgandan keyin stream qo‘shimcha ma’lumotlarni qabul qilishga tayyor bo‘lganda chiqariladi
  • finish : Barcha ma’lumotlar asosiy tizimga o‘chirilganda chiqariladi
  • error : Yozish paytida xatolik yuz berganda chiqariladi
  • close : Oqimning asosiy resursi yopilganda chiqariladi
  • pipe : pipe() usuli o‘qilishi mumkin bo‘lgan oqimda chaqirilganda chiqariladi
  • unpipe : unpipe() usuli o‘qilishi mumkin bo‘lgan oqimda chaqirilganda chiqariladi

stream.pipeline() metodi

pipeline() funksiyasi (Node.js v10.0.0 dan beri mavjud) oqimlarni o‘zaro ulashning (pipe) ishonchliroq usuli hisoblanadi, ayniqsa xatolarni qayta ishlashda.

const { pipeline } = require('stream');
const fs = require('fs');
const zlib = require('zlib');

// Create a pipeline that handles errors properly
pipeline(
  fs.createReadStream('source.txt'),
  zlib.createGzip(),
  fs.createWriteStream('destination.txt.gz'),
  (err) => {
    if (err) {
      console.error('Pipeline failed:', err);
    } else {
      console.log('Pipeline succeeded!');
    }
  }
);
Misolni ishga tushirish »

Eslatma: pipeline() barcha oqimlarni to‘g‘ri tozalaydi, agar ulardan birida xatolik yuzaga kelsa, xotira oqishining oldini oladi.


Obyekt rejimi oqimlari

Odatiy bo‘lib, streamlar stringlar va Buffer obyektlari bilan ishlaydi.

Biroq, JavaScript obyektlari bilan ishlash uchun oqimlarni "obyekt rejimi" ga o‘tkazish mumkin.

const { Readable, Writable, Transform } = require('stream');

// Create a readable stream in object mode
const objectReadable = new Readable({
  objectMode: true,
  read() {} // Implementation required but can be no-op
});

// Create a transform stream in object mode
const objectTransform = new Transform({
  objectMode: true,
  transform(chunk, encoding, callback) {
    // Add a property to the object
    chunk.transformed = true;
    chunk.timestamp = new Date();
    this.push(chunk);
    callback();
  } });

// Create a writable stream in object mode
const objectWritable = new Writable({
  objectMode: true,
  write(chunk, encoding, callback) {
    console.log('Received object:', chunk);
    callback();
  } });

// Connect the streams
objectReadable
  .pipe(objectTransform)
  .pipe(objectWritable);

// Push some objects to the stream
objectReadable.push({ name: 'Object 1', value: 10 });
objectReadable.push({ name: 'Object 2', value: 20 });
objectReadable.push({ name: 'Object 3', value: 30 });
objectReadable.push(null); // Signal the end of data
Misolni ishga tushirish »

Kengaytirilgan stream naqshlari

1. pipeline() yordamida xatolarni qayta ishlash

pipeline() usuli stream zanjiridagi xatolarni hal qilishning tavsiya etilgan usulidir:

Misol

const { pipeline } = require('stream');
const fs = require('fs');
const zlib = require('zlib');

pipeline(
  fs.createReadStream('input.txt'),
  zlib.createGzip(),
  fs.createWriteStream('output.txt.gz'),
  (err) => {
   if (err) {
    console.error('Pipeline failed:', err);
   } else {
    console.log('Pipeline succeeded');
   }
  }
);
Misolni ishga tushirish »

2. Obyekt rejimi oqimlari

Streamlar faqat stringlar va bufferlar o‘rniga JavaScript obyektlari bilan ishlashi mumkin:

Misol

const { Readable } = require('stream');

// Create a readable stream in object mode
const objectStream = new Readable({
  objectMode: true,
  read() {}
});
// Push objects to the stream
objectStream.push({ id: 1, name: 'Alice' });
objectStream.push({ id: 2, name: 'Bob' });
objectStream.push(null); // Signal end of stream
// Consume the stream
objectStream.on('data', (obj) => {
  console.log('Received:', obj);
});
Misolni ishga tushirish »

Amaliy misollar

HTTP oqimi

Streamlar HTTP so‘rovlari va javoblarida keng qo‘llaniladi.

const http = require('http');
const fs = require('fs');

// Create an HTTP server
const server = http.createServer((req, res) => {
  // Handle different routes
  if (req.url === '/') {
    // Send a simple response
    res.writeHead(200, { 'Content-Type': 'text/html' });
    res.end('<h1>Stream Demo</h1><p>Try <a href="/file">streaming a file</a> or <a href="/video">streaming a video</a>.</p>');
  }
  else if (req.url === '/file') {
    // Stream a large text file
    res.writeHead(200, { 'Content-Type': 'text/plain' });
    const fileStream = fs.createReadStream('largefile.txt', 'utf8');

    // Pipe the file to the response (handles backpressure automatically)
    fileStream.pipe(res);

    // Handle errors
    fileStream.on('error', (err) => {
      console.error('File stream error:', err);
      res.statusCode = 500;
      res.end('Server Error');
    });
  }
  else if (req.url === '/video') {
    // Stream a video file with proper headers
    const videoPath = 'video.mp4';
    const stat = fs.statSync(videoPath);
    const fileSize = stat.size;
    const range = req.headers.range;

    if (range) {
      // Handle range requests for video seeking
      const parts = range.replace(/bytes=/, "").split("-");
      const start = parseInt(parts[0], 10);
      const end = parts[1] ? parseInt(parts[1], 10) : fileSize - 1;
      const chunksize = (end - start) + 1;

      const videoStream = fs.createReadStream(videoPath, { start, end });
      res.writeHead(206, {
        'Content-Range': `bytes ${start}-${end}/${fileSize}`,
        'Accept-Ranges': 'bytes',
        'Content-Length': chunksize,
        'Content-Type': 'video/mp4'
      });

      videoStream.pipe(res);
      } else {
        // No range header, send entire video
        res.writeHead(200, {
          'Content-Length': fileSize,
          'Content-Type': 'video/mp4'
        });

        fs.createReadStream(videoPath).pipe(res);
      }
  }&br>   else {
    // 404 Not Found
    res.writeHead(404, { 'Content-Type': 'text/plain' });
    res.end('Not Found');
  }
});

// Start the server
server.listen(8080, () => {
  console.log('Server running at http://localhost:8080/');
});

Katta CSV fayllarni qayta ishlash

const fs = require('fs');
const { Transform } = require('stream');
const csv = require('csv-parser'); // npm install csv-parser

// Create a transform stream to filter and transform CSV data
const filterTransform = new Transform({
  objectMode: true,
  transform(row, encoding, callback) {
    // Only pass through rows that meet our criteria
    if (parseInt(row.age) > 18) {
      // Modify the row
      row.isAdult = 'Yes';
      // Push the transformed row
      this.push(row);
    }
    }
    callback();
  }
});

// Create a writable stream for the results
const results = [];
const writeToArray = new Transform({
  objectMode: true,
  transform(row, encoding, callback) {
    results.push(row);
    callback();
  }
});

// Create the processing pipeline
fs.createReadStream('people.csv')
  .pipe(csv())
  .pipe(filterTransform)
  .pipe(writeToArray)
  .on('finish', () => {
    console.log(`Processed ${results.length} records:`);
    console.log(results);
  }
  })
  .on('error', (err) => {
    console.error('Error processing CSV:', err);
  }
  });
Misolni ishga tushirish »

Eng yaxshi amaliyotlar

  • Xatolarni qayta ishlash: Ilova ishdan chiqishining oldini olish uchun har doim oqimlardagi xato hodisalarini boshqaring.
  • Pipeline(): dan foydalaning Xatolarni yaxshiroq boshqarish va tozalash uchun .pipe() dan stream.pipeline() ni afzal qiling.
  • Tushning orqa bosimi: Xotira bilan bog‘liq muammolarni oldini olish uchun write() qaytariladigan qiymatini hurmat qiling.
  • Oqimlarni tugatish: Ish tugagach, yozish mumkin bo‘lgan oqimlarda har doim end() raqamiga chaqiruv qiling.
  • Sinxron operatsiyalardan saqlaning: Stream ishlov beruvchilari ichidagi sinxron operatsiyalar bilan voqealar siklini bloklamang.
  • Buffer hajmi: HighWaterMark (buffer hajmi) sozlamalariga e’tibor bering.

Ogohlantirish: Oqimlarni noto‘g‘ri boshqarish xotiraning oqishiga va ishlash bilan bog‘liq muammolarga olib kelishi mumkin.

Har doim xatolarni ko‘rib chiqing va oqimlarni to‘g‘ri yakunlang.


Xulosa

Streamlar Node.js-ning asosiy tushunchasi bo‘lib, ma’lumotlarni samarali qayta ishlash imkonini beradi. Ular:

  • Hamma narsani xotiraga yuklamasdan, ma’lumotlarni parcha-parcha qayta ishlang
  • Katta ma’lumotlar to‘plamlari uchun yaxshi xotira samaradorligini ta’minlang
  • Barcha ma’lumotlar mavjud bo‘lgunga qadar qayta ishlashni boshlashga ruxsat bering
  • Kuchli ma’lumotlarni qayta ishlash payplar (pipes)ini yoqing
  • Asosiy Node.js API-larida keng qo‘llaniladi



W3Schools Pathfinder

Yutuqlaringizni kuzating – bu bepul!