استریم‌ها - راهنمای قطعی

یاد بگیرید که چگونه از استریم‌های قابل خواندن، قابل نوشتن و تبدیل با API استریم‌ها استفاده کنید.

API استریمز به شما این امکان را می‌دهد که به صورت برنامه‌نویسی به استریم‌های داده‌ای که از طریق شبکه دریافت می‌شوند یا به هر وسیله‌ای به صورت محلی ایجاد می‌شوند، دسترسی پیدا کنید و آنها را با جاوا اسکریپت پردازش کنید. استریمینگ شامل تجزیه منبعی است که می‌خواهید دریافت، ارسال یا تبدیل کنید و سپس این قطعات را بیت به بیت پردازش کنید. در حالی که استریمینگ کاری است که مرورگرها به هر حال هنگام دریافت داده‌هایی مانند HTML یا ویدیوها برای نمایش در صفحات وب انجام می‌دهند، این قابلیت قبل از معرفی fetch با استریم‌ها در سال ۲۰۱۵، هرگز برای جاوا اسکریپت در دسترس نبوده است.

نکته: از نظر فنی، استریمینگ با XMLHttpRequest امکان‌پذیر است، اما راه‌حل بهینه‌ای نیست. در اینجا یک مثال خوب از 'XMLHttpRequest' در گیت‌هاب آورده شده است .

پیش از این، اگر می‌خواستید نوعی منبع (چه ویدیو، چه فایل متنی و غیره) را پردازش کنید، باید کل فایل را دانلود می‌کردید، منتظر می‌ماندید تا به فرمت مناسبی deserialize شود و سپس آن را پردازش می‌کردید. با در دسترس قرار گرفتن streamها برای جاوااسکریپت، همه این‌ها تغییر می‌کند. اکنون می‌توانید داده‌های خام را به محض اینکه در کلاینت در دسترس قرار گرفتند، بدون نیاز به تولید بافر، رشته یا blob، به تدریج با جاوااسکریپت پردازش کنید. این امر تعدادی از موارد استفاده را باز می‌کند که برخی از آنها را در زیر فهرست می‌کنم:

  • جلوه‌های ویدیویی: لوله‌کشی یک جریان ویدیویی قابل خواندن از طریق یک جریان تبدیل که جلوه‌ها را به صورت بلادرنگ اعمال می‌کند.
  • فشرده‌سازی (از)داده‌ها: لوله‌کشی یک جریان فایل از طریق یک جریان تبدیل که به صورت انتخابی آن را فشرده‌سازی (از) می‌کند.
  • رمزگشایی تصویر: لوله‌کشی یک جریان پاسخ HTTP از طریق یک جریان تبدیل که بایت‌ها را به داده‌های بیت‌مپ رمزگشایی می‌کند، و سپس از طریق یک جریان تبدیل دیگر که بیت‌مپ‌ها را به PNG تبدیل می‌کند. اگر این قابلیت درون کنترل‌کننده‌ی fetch یک سرویس ورکر نصب شود، به شما امکان می‌دهد فرمت‌های تصویری جدید مانند AVIF را به صورت شفاف چندلایه کنید.

پشتیبانی مرورگر

ReadableStream و WritableStream

Browser Support

  • کروم: ۴۳.
  • لبه: ۱۴.
  • فایرفاکس: ۶۵.
  • سافاری: ۱۰.۱.

Source

ترنس‌استریم

Browser Support

  • کروم: ۶۷.
  • لبه: ۷۹.
  • فایرفاکس: ۱۰۲.
  • سافاری: ۱۴.۱.

Source

مفاهیم اصلی

قبل از اینکه به جزئیات انواع مختلف جریان‌ها بپردازم، اجازه دهید برخی مفاهیم اصلی را معرفی کنم.

تکه‌ها

یک تکه، قطعه‌ای از داده است که در یک جریان نوشته یا از آن خوانده می‌شود. می‌تواند از هر نوعی باشد؛ جریان‌ها حتی می‌توانند شامل تکه‌هایی از انواع مختلف باشند. اغلب اوقات، یک تکه، واحد داده‌ای بسیار کوچک برای یک جریان داده مشخص نخواهد بود. برای مثال، یک جریان بایت ممکن است شامل تکه‌هایی متشکل از واحدهای Uint8Array کیلوبایتی به حجم 16 کیلوبایت باشد، نه بایت‌های تکی.

جریان‌های قابل خواندن

یک جریان خواندنی، منبعی از داده‌ها را نشان می‌دهد که می‌توانید از آن بخوانید. به عبارت دیگر، داده‌ها از یک جریان خواندنی خارج می‌شوند . به طور مشخص، یک جریان خواندنی، نمونه‌ای از کلاس ReadableStream است.

جریان‌های قابل نوشتن

یک جریان قابل نوشتن، مقصدی برای داده‌ها است که می‌توانید در آن بنویسید. به عبارت دیگر، داده‌ها به یک جریان قابل نوشتن وارد می‌شوند . به طور مشخص، یک جریان قابل نوشتن، نمونه‌ای از کلاس WritableStream است.

تبدیل جریان‌ها

یک جریان تبدیل از یک جفت جریان تشکیل شده است: یک جریان قابل نوشتن، که به عنوان سمت قابل نوشتن آن شناخته می‌شود، و یک جریان قابل خواندن، که به عنوان سمت قابل خواندن آن شناخته می‌شود. استعاره دنیای واقعی برای این، یک مترجم همزمان است که در حال اجرا از یک زبان به زبان دیگر ترجمه می‌کند. به روشی خاص برای جریان تبدیل، نوشتن در سمت قابل نوشتن منجر به در دسترس قرار گرفتن داده‌های جدید برای خواندن از سمت قابل خواندن می‌شود. به طور مشخص، هر شیء با یک ویژگی writable و یک ویژگی readable می‌تواند به عنوان یک جریان تبدیل عمل کند. با این حال، کلاس استاندارد TransformStream ایجاد چنین جفتی را که به درستی درهم تنیده شده باشد، آسان‌تر می‌کند.

زنجیرهای لوله

جریان‌ها در درجه اول با اتصال آنها به یکدیگر استفاده می‌شوند. یک جریان قابل خواندن را می‌توان مستقیماً با استفاده از متد pipeTo() از جریان قابل خواندن به یک جریان قابل نوشتن لوله‌کشی کرد، یا می‌توان ابتدا با استفاده از متد pipeThrough() از جریان قابل خواندن، آن را از طریق یک یا چند جریان تبدیل لوله‌کشی کرد. مجموعه‌ای از جریان‌ها که به این روش به هم لوله‌کشی می‌شوند ، به عنوان زنجیره لوله شناخته می‌شوند.

فشار معکوس

وقتی یک زنجیره لوله ساخته می‌شود، سیگنال‌هایی در مورد سرعت جریان قطعات از طریق آن منتشر می‌کند. اگر هر مرحله در زنجیره هنوز نتواند قطعات را بپذیرد، سیگنالی را در سراسر زنجیره لوله به عقب منتشر می‌کند، تا زمانی که در نهایت به منبع اصلی دستور داده شود که تولید قطعات را با این سرعت متوقف کند. این فرآیند عادی‌سازی جریان، فشار معکوس نامیده می‌شود.

تی کردن

یک جریان قابل خواندن را می‌توان با استفاده از متد tee() آن، teed کرد (که از شکل حرف بزرگ 'T' گرفته شده است). این کار جریان را قفل می‌کند، یعنی دیگر مستقیماً قابل استفاده نیست؛ با این حال، دو جریان جدید به نام شاخه ایجاد می‌کند که می‌توانند به طور مستقل مورد استفاده قرار گیرند. teeing همچنین مهم است زیرا جریان‌ها را نمی‌توان به عقب برگرداند یا مجدداً راه‌اندازی کرد، که بعداً در این مورد بیشتر توضیح خواهیم داد.

یک زنجیره لوله، یک جریان قابل خواندن را از یک فراخوانی به API واکشی، از طریق یک جریان تبدیل، مسیریابی می‌کند، سپس برای اولین جریان قابل خواندن حاصل، هم به مرورگر و هم به حافظه پنهان سرویس دهنده برای دومین جریان قابل خواندن حاصل، ارسال و پردازش می‌کند.

مکانیک یک جریان قابل خواندن

یک جریان خواندنی، منبع داده‌ای است که در جاوا اسکریپت توسط یک شیء ReadableStream که از یک منبع اصلی جریان می‌یابد، نمایش داده می‌شود. سازنده ReadableStream() یک شیء جریان خواندنی را از کنترل‌کننده‌های داده شده ایجاد و برمی‌گرداند. دو نوع منبع اصلی وجود دارد:

  • منابع فشاری (Push sources) دائماً وقتی به آنها دسترسی پیدا کرده‌اید، داده‌ها را به سمت شما ارسال می‌کنند و شروع، مکث یا لغو دسترسی به جریان به شما بستگی دارد. نمونه‌هایی از این موارد شامل پخش زنده ویدیو، رویدادهای ارسالی از سرور یا WebSockets است.
  • منابع pull شما را ملزم می‌کنند که پس از اتصال، صریحاً از آنها داده درخواست کنید. به عنوان مثال می‌توان به عملیات HTTP از طریق فراخوانی‌های fetch() یا XMLHttpRequest اشاره کرد.

داده‌های جریان به صورت متوالی در قطعات کوچکی به نام تکه‌ها خوانده می‌شوند. تکه‌هایی که در یک جریان قرار می‌گیرند، در صف قرار می‌گیرند. این بدان معناست که آنها در یک صف منتظر خواندن هستند. یک صف داخلی، تکه‌هایی را که هنوز خوانده نشده‌اند، پیگیری می‌کند.

یک استراتژی صف‌بندی، شیئی است که تعیین می‌کند چگونه یک جریان باید بر اساس وضعیت صف داخلی خود، فشار برگشتی را ارسال کند. استراتژی صف‌بندی به هر بخش اندازه‌ای اختصاص می‌دهد و اندازه کل تمام بخش‌های موجود در صف را با یک عدد مشخص، که به عنوان علامت بالای آب شناخته می‌شود، مقایسه می‌کند.

بخش‌های درون جریان توسط یک خواننده (reader) خوانده می‌شوند. این خواننده، داده‌ها را تک تک بازیابی می‌کند و به شما امکان می‌دهد هر نوع عملیاتی را که می‌خواهید روی آن انجام دهید. خواننده به همراه کد پردازشی دیگری که همراه آن است، مصرف‌کننده (consumer) نامیده می‌شود.

ساختار بعدی در این زمینه، کنترلر (controller) نام دارد. هر جریان قابل خواندن (readable stream) یک کنترلر مرتبط دارد که همانطور که از نامش پیداست، به شما امکان کنترل جریان را می‌دهد.

فقط یک خواننده می‌تواند همزمان یک جریان را بخواند؛ وقتی یک خواننده ایجاد می‌شود و شروع به خواندن جریان می‌کند (یعنی به یک خواننده فعال تبدیل می‌شود)، به آن قفل می‌شود. اگر می‌خواهید خواننده دیگری خواندن جریان شما را به عهده بگیرد، معمولاً قبل از انجام هر کار دیگری باید خواننده اول را آزاد کنید (اگرچه می‌توانید جریان‌ها را tee کنید ).

ایجاد یک جریان قابل خواندن

شما با فراخوانی سازنده‌ی ReadableStream() یک جریان قابل خواندن ایجاد می‌کنید. این سازنده یک آرگومان اختیاری underlyingSource دارد که نشان‌دهنده‌ی یک شیء با متدها و ویژگی‌هایی است که نحوه‌ی رفتار نمونه‌ی جریان ساخته‌شده را تعریف می‌کنند.

underlyingSource

این می‌تواند از متدهای اختیاری و تعریف‌شده توسط توسعه‌دهنده زیر استفاده کند:

  • start(controller) : بلافاصله پس از ساخت شیء فراخوانی می‌شود. این متد می‌تواند به منبع جریان دسترسی پیدا کند و هر کار دیگری را که برای راه‌اندازی عملکرد جریان لازم است، انجام دهد. اگر این فرآیند به صورت ناهمگام انجام شود، متد می‌تواند یک promise را برای نشان دادن موفقیت یا شکست برگرداند. پارامتر controller ارسالی به این متد، ReadableStreamDefaultController است.
  • pull(controller) : می‌تواند برای کنترل جریان (stream) در هنگام دریافت تکه‌های بیشتر داده‌ها استفاده شود. تا زمانی که صف داخلی تکه‌های داده پر نشده باشد، این تابع به طور مکرر فراخوانی می‌شود، تا زمانی که صف به بالاترین حد خود برسد. اگر نتیجه فراخوانی pull() یک promise باشد، pull() تا زمانی که promise مذکور برآورده نشود، دوباره فراخوانی نخواهد شد. اگر promise رد شود، جریان با خطا مواجه می‌شود.
  • cancel(reason) : زمانی فراخوانی می‌شود که مصرف‌کننده‌ی استریم، استریم را لغو کند.
const readableStream = new ReadableStream({
  start(controller) {
    /* … */
  },

  pull(controller) {
    /* … */
  },

  cancel(reason) {
    /* … */
  },
});

کنترلر ReadableStreamDefaultController از متدهای زیر پشتیبانی می‌کند:

/* … */
start(controller) {
  controller.enqueue('The first chunk!');
},
/* … */

queuingStrategy

دومین آرگومان اختیاری سازنده‌ی ReadableStream() queuingStrategy است. این شیء به صورت اختیاری یک استراتژی صف‌بندی برای استریم تعریف می‌کند که دو پارامتر می‌گیرد:

  • highWaterMark : یک عدد غیر منفی که نشانگر بالاترین نقطه آب در جریان با استفاده از این استراتژی صف‌بندی است.
  • size(chunk) : تابعی که اندازه متناهی و غیرمنفی مقدار تکه داده شده را محاسبه و برمی‌گرداند. نتیجه برای تعیین فشار برگشتی استفاده می‌شود که از طریق ویژگی مناسب ReadableStreamDefaultController.desiredSize آشکار می‌شود. همچنین زمان فراخوانی متد pull() منبع اصلی را تعیین می‌کند.
const readableStream = new ReadableStream({
    /* … */
  },
  {
    highWaterMark: 10,
    size(chunk) {
      return chunk.length;
    },
  },
);

متدهای getReader() و read()

برای خواندن از یک جریان قابل خواندن، به یک خواننده نیاز دارید که یک ReadableStreamDefaultReader خواهد بود. متد getReader() از رابط ReadableStream یک خواننده ایجاد می‌کند و جریان را به آن قفل می‌کند. در حالی که جریان قفل شده است، تا زمانی که این خواننده آزاد نشود، هیچ خواننده دیگری قابل دسترسی نیست.

متد read() از رابط ReadableStreamDefaultReader یک promise را برمی‌گرداند که دسترسی به بخش بعدی در صف داخلی جریان را فراهم می‌کند. این promise بسته به وضعیت جریان، با یک نتیجه اجرا یا رد می‌شود. احتمالات مختلف به شرح زیر است:

  • اگر یک تکه موجود باشد، promise با یک شیء از فرم زیر انجام خواهد شد.
    { value: chunk, done: false } .
  • اگر جریان بسته شود، promise با شیء‌ای به شکل زیر محقق خواهد شد.
    { value: undefined, done: true } .
  • اگر جریان دچار خطا شود، promise با خطای مربوطه رد می‌شود.
const reader = readableStream.getReader();
while (true) {
  const { done, value } = await reader.read();
  if (done) {
    console.log('The stream is done.');
    break;
  }
  console.log('Just read a chunk:', value);
}

ملک locked

شما می‌توانید با دسترسی به ویژگی ReadableStream.locked بررسی کنید که آیا یک جریان قابل خواندن قفل شده است یا خیر.

const locked = readableStream.locked;
console.log(`The stream is ${locked ? 'indeed' : 'not'} locked.`);

نمونه‌های کد جریانی قابل خواندن

نمونه کد زیر تمام مراحل عملی را نشان می‌دهد. ابتدا، یک ReadableStream ایجاد کنید که در آرگومان underlyingSource خود (یعنی کلاس TimestampSource ) یک متد start() تعریف کند. این متد controller جریان می‌گوید که هر ثانیه به مدت ده ثانیه یک timestamp را در صف enqueue() . در نهایت، به کنترلر می‌گوید که جریان را close() . شما با ایجاد یک خواننده از طریق متد getReader() و فراخوانی read() تا زمانی که جریان done شود، از این جریان استفاده می‌کنید.

class TimestampSource {
  #interval

  start(controller) {
    this.#interval = setInterval(() => {
      const string = new Date().toLocaleTimeString();
      // Add the string to the stream.
      controller.enqueue(string);
      console.log(`Enqueued ${string}`);
    }, 1_000);

    setTimeout(() => {
      clearInterval(this.#interval);
      // Close the stream after 10s.
      controller.close();
    }, 10_000);
  }

  cancel() {
    // This is called if the reader cancels.
    clearInterval(this.#interval);
  }
}

const stream = new ReadableStream(new TimestampSource());

async function concatStringStream(stream) {
  let result = '';
  const reader = stream.getReader();
  while (true) {
    // The `read()` method returns a promise that
    // resolves when a value has been received.
    const { done, value } = await reader.read();
    // Result objects contain two properties:
    // `done`  - `true` if the stream has already given you all its data.
    // `value` - Some data. Always `undefined` when `done` is `true`.
    if (done) return result;
    result += value;
    console.log(`Read ${result.length} characters so far`);
    console.log(`Most recently read chunk: ${value}`);
  }
}
concatStringStream(stream).then((result) => console.log('Stream complete', result));

تکرار ناهمزمان

بررسی اینکه آیا جریان در هر تکرار حلقه read() done است یا خیر، ممکن است راحت‌ترین API نباشد. خوشبختانه به زودی روش بهتری برای انجام این کار وجود خواهد داشت: تکرار ناهمزمان.

for await (const chunk of stream) {
  console.log(chunk);
}

یک راه حل برای استفاده از تکرار ناهمزمان در حال حاضر، پیاده‌سازی این رفتار با یک polyfill است.

if (!ReadableStream.prototype[Symbol.asyncIterator]) {
  ReadableStream.prototype[Symbol.asyncIterator] = async function* () {
    const reader = this.getReader();
    try {
      while (true) {
        const {done, value} = await reader.read();
        if (done) {
          return;
          }
        yield value;
      }
    }
    finally {
      reader.releaseLock();
    }
  }
}

ایجاد یک جریان قابل خواندن

متد tee() از رابط ReadableStream جریان قابل خواندن فعلی را tee می‌کند و یک آرایه دو عنصری حاوی دو شاخه حاصل را به عنوان نمونه‌های جدید ReadableStream برمی‌گرداند. این به دو خواننده اجازه می‌دهد تا یک جریان را همزمان بخوانند. برای مثال، اگر می‌خواهید پاسخی را از سرور دریافت کرده و آن را به مرورگر ارسال کنید، اما آن را به حافظه پنهان (cache) سرویس ورکر نیز ارسال کنید، می‌توانید این کار را در یک سرویس ورکر انجام دهید. از آنجایی که یک بدنه پاسخ نمی‌تواند بیش از یک بار مصرف شود، برای انجام این کار به دو کپی نیاز دارید. برای لغو جریان، باید هر دو شاخه حاصل را لغو کنید. tee کردن یک جریان معمولاً آن را برای مدت زمان قفل می‌کند و از قفل شدن آن توسط سایر خوانندگان جلوگیری می‌کند.

const readableStream = new ReadableStream({
  start(controller) {
    // Called by constructor.
    console.log('[start]');
    controller.enqueue('a');
    controller.enqueue('b');
    controller.enqueue('c');
  },
  pull(controller) {
    // Called `read()` when the controller's queue is empty.
    console.log('[pull]');
    controller.enqueue('d');
    controller.close();
  },
  cancel(reason) {
    // Called when the stream is canceled.
    console.log('[cancel]', reason);
  },
});

// Create two `ReadableStream`s.
const [streamA, streamB] = readableStream.tee();

// Read streamA iteratively one by one. Typically, you
// would not do it this way, but you certainly can.
const readerA = streamA.getReader();
console.log('[A]', await readerA.read()); //=> {value: "a", done: false}
console.log('[A]', await readerA.read()); //=> {value: "b", done: false}
console.log('[A]', await readerA.read()); //=> {value: "c", done: false}
console.log('[A]', await readerA.read()); //=> {value: "d", done: false}
console.log('[A]', await readerA.read()); //=> {value: undefined, done: true}

// Read streamB in a loop. This is the more common way
// to read data from the stream.
const readerB = streamB.getReader();
while (true) {
  const result = await readerB.read();
  if (result.done) break;
  console.log('[B]', result);
}

جریان‌های بایت قابل خواندن

برای جریان‌هایی که نمایانگر بایت‌ها هستند، یک نسخه توسعه‌یافته از جریان قابل خواندن ارائه شده است تا بایت‌ها را به طور موثر مدیریت کند، به ویژه با به حداقل رساندن کپی‌ها. جریان‌های بایت امکان استفاده از خواننده‌های بافر خود-خود-را (BYOB) فراهم می‌کنند. پیاده‌سازی پیش‌فرض می‌تواند طیف وسیعی از خروجی‌های مختلف مانند رشته‌ها یا بافرهای آرایه‌ای را در مورد WebSockets ارائه دهد، در حالی که جریان‌های بایت خروجی بایت را تضمین می‌کنند. علاوه بر این، خواننده‌های BYOB از مزایای پایداری برخوردارند. دلیل این امر این است که اگر یک بافر جدا شود، می‌تواند تضمین کند که دو بار در یک بافر نوشته نمی‌شود و از این رو از شرایط رقابتی جلوگیری می‌شود. خواننده‌های BYOB می‌توانند تعداد دفعاتی را که مرورگر باید جمع‌آوری زباله را اجرا کند، کاهش دهند، زیرا می‌توانند از بافرها دوباره استفاده کنند.

ایجاد یک جریان بایت قابل خواندن

شما می‌توانید با ارسال یک پارامتر type اضافی به سازنده‌ی ReadableStream() یک جریان بایت قابل خواندن ایجاد کنید.

new ReadableStream({ type: 'bytes' });

underlyingSource

منبع اصلی یک جریان بایت قابل خواندن، یک ReadableByteStreamController برای دستکاری دریافت می‌کند. متد ReadableByteStreamController.enqueue() آن، یک آرگومان chunk دریافت می‌کند که مقدار آن ArrayBufferView است. ویژگی ReadableByteStreamController.byobRequest درخواست pull فعلی BYOB را برمی‌گرداند، یا در صورت عدم وجود، null را برمی‌گرداند. در نهایت، ویژگی ReadableByteStreamController.desiredSize اندازه مورد نظر برای پر کردن صف داخلی جریان کنترل‌شده را برمی‌گرداند.

queuingStrategy

دومین آرگومان اختیاری سازنده‌ی ReadableStream() queuingStrategy است. این شیء به صورت اختیاری یک استراتژی صف‌بندی برای استریم تعریف می‌کند که یک پارامتر می‌گیرد:

  • highWaterMark : تعداد غیرمنفی از بایت‌ها که نشانگر بالاترین میزان واترمارک جریان با استفاده از این استراتژی صف‌بندی است. این برای تعیین فشار برگشتی استفاده می‌شود که از طریق ویژگی مناسب ReadableByteStreamController.desiredSize آشکار می‌شود. همچنین زمان فراخوانی متد pull() منبع اصلی را تعیین می‌کند.

متدهای getReader() و read()

سپس می‌توانید با تنظیم پارامتر mode به صورت زیر به ReadableStreamBYOBReader دسترسی پیدا کنید: ReadableStream.getReader({ mode: "byob" }) . این کار امکان کنترل دقیق‌تر بر تخصیص بافر را فراهم می‌کند تا از کپی شدن جلوگیری شود. برای خواندن از جریان بایت، باید ReadableStreamBYOBReader.read(view) را فراخوانی کنید، که در آن view یک ArrayBufferView است.

نمونه کد جریان بایت قابل خواندن

const reader = readableStream.getReader({ mode: "byob" });

let startingAB = new ArrayBuffer(1_024);
const buffer = await readInto(startingAB);
console.log("The first 1024 bytes, or less:", buffer);

async function readInto(buffer) {
  let offset = 0;

  while (offset < buffer.byteLength) {
    const { value: view, done } =
        await reader.read(new Uint8Array(buffer, offset, buffer.byteLength - offset));
    buffer = view.buffer;
    if (done) {
      break;
    }
    offset += view.byteLength;
  }

  return buffer;
}

تابع زیر جریان‌های بایتی قابل خواندنی را برمی‌گرداند که امکان خواندن کارآمد بدون کپی از یک آرایه تصادفی تولید شده را فراهم می‌کند. به جای استفاده از اندازه قطعه از پیش تعیین شده 1024، سعی می‌کند بافر ارائه شده توسط توسعه‌دهنده را پر کند و امکان کنترل کامل را فراهم کند.

const DEFAULT_CHUNK_SIZE = 1_024;

function makeReadableByteStream() {
  return new ReadableStream({
    type: 'bytes',

    pull(controller) {
      // Even when the consumer is using the default reader,
      // the auto-allocation feature allocates a buffer and
      // passes it to us via `byobRequest`.
      const view = controller.byobRequest.view;
      view = crypto.getRandomValues(view);
      controller.byobRequest.respond(view.byteLength);
    },

    autoAllocateChunkSize: DEFAULT_CHUNK_SIZE,
  });
}

مکانیک یک جریان قابل نوشتن

یک جریان قابل نوشتن، مقصدی است که می‌توانید داده‌ها را در آن بنویسید، که در جاوا اسکریپت توسط یک شیء WritableStream نمایش داده می‌شود. این به عنوان یک انتزاع بر روی یک سینک زیرین - یک سینک ورودی/خروجی سطح پایین‌تر که داده‌های خام در آن نوشته می‌شوند - عمل می‌کند.

داده‌ها از طریق یک writer ، که هر بار یک تکه است، در stream نوشته می‌شوند. یک chunk می‌تواند اشکال مختلفی داشته باشد، درست مانند chunkهای یک reader. می‌توانید از هر کدی که دوست دارید برای تولید chunkهای آماده برای نوشتن استفاده کنید؛ writer به همراه کد مرتبط، producer نامیده می‌شود.

وقتی یک نویسنده ایجاد می‌شود و شروع به نوشتن در یک جریان (یک نویسنده فعال ) می‌کند، گفته می‌شود که به آن جریان قفل شده است. فقط یک نویسنده می‌تواند همزمان در یک جریان قابل نوشتن بنویسد. اگر می‌خواهید نویسنده دیگری شروع به نوشتن در جریان شما کند، معمولاً باید آن را آزاد کنید، قبل از اینکه نویسنده دیگری را به آن متصل کنید.

یک صف داخلی، بخش‌هایی از داده‌ها را که در جریان نوشته شده‌اند اما هنوز توسط سینک زیرین پردازش نشده‌اند، ردیابی می‌کند.

یک استراتژی صف‌بندی، شیئی است که تعیین می‌کند چگونه یک جریان باید بر اساس وضعیت صف داخلی خود، فشار برگشتی را ارسال کند. استراتژی صف‌بندی به هر بخش اندازه‌ای اختصاص می‌دهد و اندازه کل تمام بخش‌های موجود در صف را با یک عدد مشخص، که به عنوان علامت بالای آب شناخته می‌شود، مقایسه می‌کند.

ساختار نهایی، کنترلر (controller) نامیده می‌شود. هر جریان قابل نوشتن (writable stream) یک کنترلر مرتبط دارد که به شما امکان می‌دهد جریان را کنترل کنید (مثلاً آن را لغو کنید).

ایجاد یک جریان قابل نوشتن

رابط WritableStream از API مربوط به Streams، یک انتزاع استاندارد برای نوشتن داده‌های استریمینگ به یک مقصد، که به عنوان sink شناخته می‌شود، ارائه می‌دهد. این شیء دارای backpressure و queuing داخلی است. شما با فراخوانی سازنده‌ی آن WritableStream() ، یک استریم قابل نوشتن ایجاد می‌کنید. این سازنده یک پارامتر اختیاری underlyingSink دارد که نشان‌دهنده‌ی یک شیء با متدها و ویژگی‌هایی است که نحوه‌ی رفتار نمونه‌ی استریم ساخته شده را تعریف می‌کنند.

سینک underlyingSink

underlyingSink می‌تواند شامل متدهای اختیاری و تعریف‌شده توسط توسعه‌دهنده زیر باشد. پارامتر controller ارسالی به برخی از متدها، یک WritableStreamDefaultController است.

  • start(controller) : این متد بلافاصله پس از ساخت شیء فراخوانی می‌شود. محتوای این متد باید با هدف دسترسی به sink زیرین باشد. اگر این فرآیند به صورت غیرهمزمان انجام شود، می‌تواند یک promise را برای نشان دادن موفقیت یا شکست برگرداند.
  • write(chunk, controller) : این متد زمانی فراخوانی می‌شود که یک تکه داده جدید (که در پارامتر chunk مشخص شده است) آماده نوشتن در sink زیرین باشد. این متد می‌تواند یک promise را برای نشان دادن موفقیت یا شکست عملیات نوشتن برگرداند. این متد فقط پس از موفقیت‌آمیز بودن نوشتن‌های قبلی فراخوانی می‌شود و هرگز پس از بسته شدن یا لغو شدن استریم فراخوانی نمی‌شود.
  • close(controller) : این متد در صورتی فراخوانی می‌شود که برنامه اعلام کند نوشتن بخش‌هایی از داده در استریم را به پایان رسانده است. محتویات باید هر کاری که برای نهایی کردن نوشتن‌ها در سینک زیرین لازم است را انجام دهند و دسترسی به آن را آزاد کنند. اگر این فرآیند ناهمزمان باشد، می‌تواند یک promise را برای اعلام موفقیت یا شکست برگرداند. این متد فقط پس از موفقیت‌آمیز بودن تمام نوشتن‌های صف‌بندی شده فراخوانی می‌شود.
  • abort(reason) : این متد در صورتی فراخوانی می‌شود که برنامه اعلام کند می‌خواهد استریم را به طور ناگهانی ببندد و آن را در حالت خطا قرار دهد. این متد می‌تواند مانند close() هر منبع نگه‌داشته شده‌ای را پاک کند، اما abort() حتی اگر نوشتن‌ها در صف انتظار باشند نیز فراخوانی می‌شود. آن تکه‌ها دور انداخته می‌شوند. اگر این فرآیند ناهمزمان باشد، می‌تواند یک promise را برای نشان دادن موفقیت یا شکست برگرداند. پارامتر reason شامل یک DOMString است که توضیح می‌دهد چرا استریم لغو شده است.
const writableStream = new WritableStream({
  start(controller) {
    /* … */
  },

  write(chunk, controller) {
    /* … */
  },

  close(controller) {
    /* … */
  },

  abort(reason) {
    /* … */
  },
});

رابط WritableStreamDefaultController از API Streams، کنترلری را نشان می‌دهد که امکان کنترل وضعیت WritableStream را در طول راه‌اندازی، با ارسال بخش‌های بیشتر برای نوشتن، یا در پایان نوشتن، فراهم می‌کند. هنگام ساخت یک WritableStream ، به sink زیرین، یک نمونه WritableStreamDefaultController مربوطه برای دستکاری داده می‌شود. WritableStreamDefaultController تنها یک متد دارد: WritableStreamDefaultController.error() ، که باعث می‌شود هرگونه تعامل آینده با جریان مرتبط با خطا مواجه شود. WritableStreamDefaultController همچنین از یک ویژگی signal پشتیبانی می‌کند که نمونه‌ای از AbortSignal را برمی‌گرداند و در صورت نیاز، امکان توقف عملیات WritableStream را فراهم می‌کند.

/* … */
write(chunk, controller) {
  try {
    // Try to do something dangerous with `chunk`.
  } catch (error) {
    controller.error(error.message);
  }
},
/* … */

queuingStrategy

دومین آرگومان اختیاری سازنده‌ی WritableStream() queuingStrategy است. این شیء به صورت اختیاری یک استراتژی صف‌بندی برای استریم تعریف می‌کند که دو پارامتر می‌گیرد:

  • highWaterMark : یک عدد غیر منفی که نشانگر بالاترین نقطه آب در جریان با استفاده از این استراتژی صف‌بندی است.
  • size(chunk) : تابعی که اندازه متناهی و غیرمنفی مقدار chunk داده شده را محاسبه و برمی‌گرداند. نتیجه برای تعیین فشار برگشتی استفاده می‌شود که از طریق ویژگی مناسب WritableStreamDefaultWriter.desiredSize آشکار می‌شود.

متدهای getWriter() و write()

برای نوشتن در یک جریان قابل نوشتن، به یک نویسنده نیاز دارید که یک WritableStreamDefaultWriter خواهد بود. متد getWriter() از رابط WritableStream یک نمونه جدید از WritableStreamDefaultWriter را برمی‌گرداند و جریان را به آن نمونه قفل می‌کند. در حالی که جریان قفل شده است، تا زمانی که نویسنده فعلی آزاد نشود، هیچ نویسنده دیگری قابل دستیابی نیست.

متد write() از رابط WritableStreamDefaultWriter یک تکه داده ارسالی را در یک WritableStream و sink زیرین آن می‌نویسد، سپس promiseای را برمی‌گرداند که موفقیت یا شکست عملیات نوشتن را نشان می‌دهد. توجه داشته باشید که معنای "موفقیت" به sink زیرین بستگی دارد؛ ممکن است نشان دهد که تکه داده پذیرفته شده است، و لزوماً به این معنی نیست که به طور ایمن در مقصد نهایی خود ذخیره شده است.

const writer = writableStream.getWriter();
const resultPromise = writer.write('The first chunk!');

ملک locked

شما می‌توانید با دسترسی به ویژگی WritableStream.locked بررسی کنید که آیا یک جریان قابل نوشتن قفل شده است یا خیر.

const locked = writableStream.locked;
console.log(`The stream is ${locked ? 'indeed' : 'not'} locked.`);

نمونه کد جریان قابل نوشتن

نمونه کد زیر تمام مراحل را در عمل نشان می‌دهد.

const writableStream = new WritableStream({
  start(controller) {
    console.log('[start]');
  },
  async write(chunk, controller) {
    console.log('[write]', chunk);
    // Wait for next write.
    await new Promise((resolve) => setTimeout(() => {
      document.body.textContent += chunk;
      resolve();
    }, 1_000));
  },
  close(controller) {
    console.log('[close]');
  },
  abort(reason) {
    console.log('[abort]', reason);
  },
});

const writer = writableStream.getWriter();
const start = Date.now();
for (const char of 'abcdefghijklmnopqrstuvwxyz') {
  // Wait to add to the write queue.
  await writer.ready;
  console.log('[ready]', Date.now() - start, 'ms');
  // The Promise is resolved after the write finishes.
  writer.write(char);
}
await writer.close();

اتصال یک جریان قابل خواندن به یک جریان قابل نوشتن

یک جریان خواندنی می‌تواند از طریق متد pipeTo() در جریان خواندنی به یک جریان نوشتنی پایپ شود. ReadableStream.pipeTo() ReadableStream فعلی را به یک WritableStream داده شده پایپ می‌کند و یک promise را برمی‌گرداند که وقتی فرآیند پایپ با موفقیت انجام شود، اجرا می‌شود یا در صورت بروز هرگونه خطا، رد می‌شود.

const readableStream = new ReadableStream({
  start(controller) {
    // Called by constructor.
    console.log('[start readable]');
    controller.enqueue('a');
    controller.enqueue('b');
    controller.enqueue('c');
  },
  pull(controller) {
    // Called when controller's queue is empty.
    console.log('[pull]');
    controller.enqueue('d');
    controller.close();
  },
  cancel(reason) {
    // Called when the stream is canceled.
    console.log('[cancel]', reason);
  },
});

const writableStream = new WritableStream({
  start(controller) {
    // Called by constructor
    console.log('[start writable]');
  },
  async write(chunk, controller) {
    // Called upon writer.write()
    console.log('[write]', chunk);
    // Wait for next write.
    await new Promise((resolve) => setTimeout(() => {
      document.body.textContent += chunk;
      resolve();
    }, 1_000));
  },
  close(controller) {
    console.log('[close]');
  },
  abort(reason) {
    console.log('[abort]', reason);
  },
});

await readableStream.pipeTo(writableStream);
console.log('[finished]');

ایجاد یک جریان تبدیل

رابط TransformStream از API Streams مجموعه‌ای از داده‌های قابل تبدیل را نشان می‌دهد. شما می‌توانید با فراخوانی سازنده‌ی آن، TransformStream() ، یک جریان تبدیل ایجاد کنید که یک شیء جریان تبدیل را از کنترل‌کننده‌های داده شده ایجاد و برمی‌گرداند. سازنده‌ی TransformStream() به عنوان اولین آرگومان خود، یک شیء جاوا اسکریپت اختیاری را می‌پذیرد که نشان‌دهنده‌ی transformer است. چنین اشیاء می‌توانند شامل هر یک از متدهای زیر باشند:

transformer

  • start(controller) : این متد بلافاصله پس از ساخت شیء فراخوانی می‌شود. معمولاً از این متد برای صف‌بندی تکه‌های پیشوند با استفاده از controller.enqueue() استفاده می‌شود. این تکه‌ها از سمت خواندنی خوانده می‌شوند اما به هیچ نوشتنی در سمت نوشتنی وابسته نیستند. اگر این فرآیند اولیه ناهمزمان باشد، به عنوان مثال به این دلیل که برای به دست آوردن تکه‌های پیشوندی تلاشی لازم است، تابع می‌تواند یک promise را برای نشان دادن موفقیت یا شکست برگرداند. یک promise رد شده باعث خطا در جریان می‌شود. هرگونه exception پرتاب شده توسط سازنده TransformStream() دوباره پرتاب می‌شود.
  • transform(chunk, controller) : این متد زمانی فراخوانی می‌شود که یک تکه جدید که در ابتدا در سمت قابل نوشتن نوشته شده است، آماده تبدیل باشد. پیاده‌سازی جریان تضمین می‌کند که این تابع فقط پس از موفقیت تبدیل‌های قبلی فراخوانی شود و هرگز قبل از تکمیل start() یا پس از flush() فراخوانی نشود. این تابع کار تبدیل واقعی جریان تبدیل را انجام می‌دهد. می‌تواند نتایج را با استفاده از controller.enqueue() در صف قرار دهد. این امر به یک تکه نوشته شده در سمت قابل نوشتن اجازه می‌دهد تا بسته به تعداد دفعات فراخوانی controller.enqueue() منجر به صفر یا چند تکه در سمت قابل خواندن شود. اگر فرآیند تبدیل ناهمزمان باشد، این تابع می‌تواند یک promise را برای نشان دادن موفقیت یا شکست تبدیل برگرداند. یک promise رد شده، هر دو سمت قابل خواندن و قابل نوشتن جریان تبدیل را با خطا مواجه می‌کند. اگر هیچ متد transform() ارائه نشود، از identity transform استفاده می‌شود که تکه‌های بدون تغییر را از سمت قابل نوشتن به سمت قابل خواندن در صف قرار می‌دهد.
  • flush(controller) : این متد پس از اینکه تمام تکه‌های نوشته شده در سمت نوشتنی با عبور موفقیت‌آمیز از transform() تبدیل شدند و سمت نوشتنی در شرف بسته شدن است، فراخوانی می‌شود. معمولاً از این برای قرار دادن تکه‌های پسوندی در صف خواندنی‌ها، قبل از بسته شدن آن، استفاده می‌شود. اگر فرآیند شستشو ناهمزمان باشد، تابع می‌تواند یک promise را برای نشان دادن موفقیت یا شکست برگرداند؛ نتیجه به فراخواننده stream.writable.write() اطلاع داده می‌شود. علاوه بر این، یک promise رد شده، هر دو سمت خواندنی و نوشتنی جریان را با خطا مواجه می‌کند. ارسال یک استثنا مانند بازگرداندن یک promise رد شده در نظر گرفته می‌شود.
const transformStream = new TransformStream({
  start(controller) {
    /* … */
  },

  transform(chunk, controller) {
    /* … */
  },

  flush(controller) {
    /* … */
  },
});

استراتژی‌های صف‌بندی writableStrategy و readableStrategy

پارامترهای اختیاری دوم و سوم سازنده‌ی TransformStream() استراتژی‌های صف‌بندی اختیاری writableStrategy و readableStrategy هستند. آن‌ها به ترتیب همانطور که در بخش‌های readable و writable stream شرح داده شده است، تعریف می‌شوند.

نمونه کد جریان تبدیل

نمونه کد زیر یک جریان تبدیل را در عمل نشان می‌دهد.

// Note that `TextEncoderStream` and `TextDecoderStream` exist now.
// This example shows how you would have done it before.
const textEncoderStream = new TransformStream({
  transform(chunk, controller) {
    console.log('[transform]', chunk);
    controller.enqueue(new TextEncoder().encode(chunk));
  },
  flush(controller) {
    console.log('[flush]');
    controller.terminate();
  },
});

(async () => {
  const readStream = textEncoderStream.readable;
  const writeStream = textEncoderStream.writable;

  const writer = writeStream.getWriter();
  for (const char of 'abc') {
    writer.write(char);
  }
  writer.close();

  const reader = readStream.getReader();
  for (let result = await reader.read(); !result.done; result = await reader.read()) {
    console.log('[value]', result.value);
  }
})();

لوله‌کشی یک جریان قابل خواندن از طریق یک جریان تبدیل

متد pipeThrough() از رابط ReadableStream یک روش زنجیره‌ای برای پایپ کردن جریان فعلی از طریق یک جریان تبدیل یا هر جفت قابل نوشتن/خواندن دیگر فراهم می‌کند. پایپ کردن یک جریان معمولاً آن را برای مدت زمان پایپ قفل می‌کند و از قفل شدن آن توسط سایر خوانندگان جلوگیری می‌کند.

const transformStream = new TransformStream({
  transform(chunk, controller) {
    console.log('[transform]', chunk);
    controller.enqueue(new TextEncoder().encode(chunk));
  },
  flush(controller) {
    console.log('[flush]');
    controller.terminate();
  },
});

const readableStream = new ReadableStream({
  start(controller) {
    // called by constructor
    console.log('[start]');
    controller.enqueue('a');
    controller.enqueue('b');
    controller.enqueue('c');
  },
  pull(controller) {
    // called read when controller's queue is empty
    console.log('[pull]');
    controller.enqueue('d');
    controller.close(); // or controller.error();
  },
  cancel(reason) {
    // called when rs.cancel(reason)
    console.log('[cancel]', reason);
  },
});

(async () => {
  const reader = readableStream.pipeThrough(transformStream).getReader();
  for (let result = await reader.read(); !result.done; result = await reader.read()) {
    console.log('[value]', result.value);
  }
})();

نمونه کد بعدی (کمی ساختگی) نشان می‌دهد که چگونه می‌توانید یک نسخه "فریاد" از fetch() را پیاده‌سازی کنید که تمام متن را با استفاده از promise پاسخ برگشتی به عنوان یک جریان و بزرگ کردن تکه به تکه حروف بزرگ، بزرگ می‌کند. مزیت این رویکرد این است که نیازی نیست منتظر دانلود کل سند باشید، که می‌تواند هنگام کار با فایل‌های بزرگ تفاوت زیادی ایجاد کند.

function upperCaseStream() {
  return new TransformStream({
    transform(chunk, controller) {
      controller.enqueue(chunk.toUpperCase());
    },
  });
}

function appendToDOMStream(el) {
  return new WritableStream({
    write(chunk) {
      el.append(chunk);
    }
  });
}

fetch('./lorem-ipsum.txt').then((response) =>
  response.body
    .pipeThrough(new TextDecoderStream())
    .pipeThrough(upperCaseStream())
    .pipeTo(appendToDOMStream(document.body))
);

نسخه آزمایشی

دموی زیر، استریم‌های خواندنی، نوشتنی و تبدیلی را در عمل نشان می‌دهد. همچنین شامل مثال‌هایی از زنجیره‌های لوله pipeThrough() و pipeTo() است و tee() را نیز نشان می‌دهد. می‌توانید به صورت اختیاری، دمو را در پنجره‌ی جداگانه اجرا کنید یا کد منبع را مشاهده کنید.

استریم‌های مفید موجود در مرورگر

تعدادی جریان مفید در مرورگر تعبیه شده‌اند. شما می‌توانید به راحتی یک ReadableStream از یک blob ایجاد کنید. متد stream() رابط Blob یک ReadableStream برمی‌گرداند که پس از خواندن، داده‌های موجود در blob را برمی‌گرداند. همچنین به یاد داشته باشید که یک شیء File نوع خاصی از Blob است و می‌تواند در هر زمینه‌ای که یک blob می‌تواند استفاده شود، مورد استفاده قرار گیرد.

const readableStream = new Blob(['hello world'], { type: 'text/plain' }).stream();

انواع جریانی (streaming) توابع TextDecoder.decode() و TextEncoder.encode() به ترتیب TextDecoderStream و TextEncoderStream نامیده می‌شوند.

const response = await fetch('https://streams.spec.whatwg.org/');
const decodedStream = response.body.pipeThrough(new TextDecoderStream());

فشرده‌سازی یا خارج کردن فایل از حالت فشرده به ترتیب با استفاده از جریان‌های تبدیل CompressionStream و DecompressionStream آسان است. نمونه کد زیر نشان می‌دهد که چگونه می‌توانید مشخصات Streams را دانلود کنید، آن را مستقیماً در مرورگر فشرده (gzip) کنید و فایل فشرده شده را مستقیماً روی دیسک بنویسید.

const response = await fetch('https://streams.spec.whatwg.org/');
const readableStream = response.body;
const compressedStream = readableStream.pipeThrough(new CompressionStream('gzip'));

const fileHandle = await showSaveFilePicker();
const writableStream = await fileHandle.createWritable();
compressedStream.pipeTo(writableStream);

جریان‌های درخواست آزمایشیِ FileSystemWritableFileStream از رابط برنامه‌نویسی کاربردی دسترسی به سیستم فایل ( File System Access API ) و fetch() نمونه‌هایی از جریان‌های قابل نوشتن در عمل هستند.

رابط برنامه‌نویسی کاربردی سریال (Serial API) به شدت از هر دو جریان خواندنی و نوشتنی استفاده می‌کند.

// Prompt user to select any serial port.
const port = await navigator.serial.requestPort();
// Wait for the serial port to open.
await port.open({ baudRate: 9_600 });
const reader = port.readable.getReader();

// Listen to data coming from the serial device.
while (true) {
  const { value, done } = await reader.read();
  if (done) {
    // Allow the serial port to be closed later.
    reader.releaseLock();
    break;
  }
  // value is a Uint8Array.
  console.log(value);
}

// Write to the serial port.
const writer = port.writable.getWriter();
const data = new Uint8Array([104, 101, 108, 108, 111]); // hello
await writer.write(data);
// Allow the serial port to be closed later.
writer.releaseLock();

در نهایت، API WebSocketStream ، استریم‌ها را با API مربوط به WebSocket ادغام می‌کند.

const wss = new WebSocketStream(WSS_URL);
const { readable, writable } = await wss.connection;
const reader = readable.getReader();
const writer = writable.getWriter();

while (true) {
  const { value, done } = await reader.read();
  if (done) {
    break;
  }
  const result = await process(value);
  await writer.write(result);
}

منابع مفید

تقدیرنامه‌ها

این مقاله توسط جیک آرچیبالد ، فرانسوا بوفورت ، سم داتون ، ماتیاس بوئلنز ، سورما ، جو مدلی و آدام رایس بررسی شده است. پست‌های وبلاگ جیک آرچیبالد در درک جریان‌ها به من کمک زیادی کرده‌اند. برخی از نمونه‌های کد از کاوش‌های کاربر گیت‌هاب @bellbind الهام گرفته شده‌اند و بخش‌هایی از نثر به شدت بر اساس اسناد وب MDN در مورد جریان‌ها ساخته شده‌اند. نویسندگان استاندارد جریان‌ها کار فوق‌العاده‌ای در نوشتن این مشخصات انجام داده‌اند.