PageSourceSearch

https://threejsnotes.netlify.app/assets/js/d6a8d145.22bb2571.js

js threejsnotes.netlify.app collected 2026-10-03 11:33:27 UTC 57,208 bytes, 1 lines download raw bytes

1"use strict";(self.webpackChunkmy_website=self.webpackChunkmy_website||[]).push([[5814],{7767:(e,n,s)=>{s.r(n),s.d(n,{assets:()=>l,contentTitle:()=>i,default:()=>h,frontMatter:()=>a,metadata:()=>o,toc:()=>c});var r=s(4848),t=s(8453);const a={},i=void 0,o={id:"node/node-streams",title:"node-streams",description:"Event Emitter",source:"@site/docs/node/02-node-streams.md",sourceDirName:"node",slug:"/node/node-streams",permalink:"/docs/node/node-streams",draft:!1,unlisted:!1,tags:[],version:"current",sidebarPosition:2,frontMatter:{},sidebar:"tutorialSidebar",previous:{title:"how-node-works",permalink:"/docs/node/how-node-works"},next:{title:"node-processes",permalink:"/docs/node/node-processes"}},l={},c=[{value:"Event Emitter",id:"event-emitter",level:2},{value:"Streams in Node",id:"streams-in-node",level:2},{value:"Why Streams?",id:"why-streams",level:3},{value:"Stream basics",id:"stream-basics",level:3},{value:"Readable and writable streams",id:"readable-and-writable-streams",level:3},{value:"Basics of readable streams",id:"basics-of-readable-streams",level:4},{value:"<strong>consuming readable streams</strong>",id:"consuming-readable-streams",level:4},{value:"Basics of Writable streams",id:"basics-of-writable-streams",level:4},{value:"Combining readable and writable streams",id:"combining-readable-and-writable-streams",level:4},{value:"Handling backpressure",id:"handling-backpressure",level:4},{value:"Custom readable and writable streams",id:"custom-readable-and-writable-streams",level:3},{value:"Extending readable streams",id:"extending-readable-streams",level:4},{value:"Creating custom readable streams",id:"creating-custom-readable-streams",level:4},{value:"Extending writable streams",id:"extending-writable-streams",level:4},{value:"Creating custom writable streams",id:"creating-custom-writable-streams",level:4},{value:"<strong>object mode</strong>",id:"object-mode",level:4},{value:"Transform Streams and duplex streams",id:"transform-streams-and-duplex-streams",level:3},{value:"Duplex streams",id:"duplex-streams",level:4},{value:"Creating custom duplex streams",id:"creating-custom-duplex-streams",level:4},{value:"Creating transform streams",id:"creating-transform-streams",level:4},{value:"Stream compression with GZIP",id:"stream-compression-with-gzip",level:4}];function d(e){const n={code:"code",div:"div",h2:"h2",h3:"h3",h4:"h4",hr:"hr",img:"img",li:"li",ol:"ol",p:"p",path:"path",pre:"pre",strong:"strong",svg:"svg",table:"table",tbody:"tbody",td:"td",th:"th",thead:"thead",tr:"tr",ul:"ul",...(0,t.R)(),...e.components};return(0,r.jsxs)(r.Fragment,{children:[(0,r.jsx)(n.h2,{id:"event-emitter",children:"Event Emitter"}),"\n",(0,r.jsx)(n.p,{children:"The event emitter api in node allows you to partake in event-based programming and is extremely simple, just based on a single messaging event."}),"\n",(0,r.jsxs)(n.p,{children:["The basic premise is that an ",(0,r.jsx)(n.code,{children:"new EventEmitter()"})," instance has only two important methods: ",(0,r.jsx)(n.code,{children:"emitter.emit()"})," to trigger events and send a payload, and ",(0,r.jsx)(n.code,{children:"emitter.on()"})," to listen for events and access the payload."]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:'const EventEmitter = require("events");\n\nconst emitter = new EventEmitter();\n\n// 1. setup listener\nemitter.on("redevent", () => console.log("RED EVENT FIRED"));\n\n// 2. emit/trigger event\nemitter.emit("redevent");\n'})}),"\n",(0,r.jsx)(n.p,{children:"You can extend the event emitter class to create event-based programming easy:"}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:"const EventEmitter = require('events');\n\nclass MyCustomEmitter extends EventEmitter {\n  constructor() {\n    super(); // Call the parent constructor\n    this.data = [];\n  }\n\n  addItem(item) {\n    this.data.push(item);\n    this.emit('itemAdded', item); // Emit an event when an item is added\n  }
1\n\n  processData() {\n    // Simulate some processing\n    console.log('Processing data...');\n    this.emit('processingComplete', this.data.length); // Emit event on completion\n  }\n}\n\nconst myObject = new MyCustomEmitter();\n\nmyObject.on('itemAdded', (item) => {\n  console.log(`New item added: ${item}`);\n});\n\nmyObject.once('processingComplete', (count) => {\n  console.log(`Finished processing ${count} items.`);\n});\n\nmyObject.addItem('A');\nmyObject.addItem('B');\nmyObject.processData();\n\n"})}),"\n",(0,r.jsx)(n.p,{children:"We can take this a step further and make a generic, type safe class wrapper around this:"}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:'import { EventEmitter, on } from "node:events";\n\ntype EventsType = Record<string, any>;\n\nexport class EventEmitterNode<T extends EventsType> {\n  private emitter: EventEmitter;\n\n  constructor() {\n    this.emitter = new EventEmitter();\n  }\n\n  sendEvent<K extends keyof T>(event: K, data: T[K]) {\n    this.emitter.emit(event as string, data);\n  }\n\n  on<K extends keyof T>(event: K, callback: (data: T[K]) => void) {\n    this.emitter.on(event as string, callback);\n    return {\n      removeListener: () => this.off(event, callback),\n    };\n  }\n\n  once<K extends keyof T>(event: K, callback: (data: T[K]) => void) {\n    this.emitter.once(event as string, callback);\n  }\n\n  async *consumeEventStream<K extends keyof T>(event: K, signal?: AbortSignal) {\n    for await (const data of on(this.emitter, event as string, { signal })) {\n      yield data as T[K];\n    }\n  }\n\n  off<K extends keyof T>(event: K, callback: (data: T[K]) => void) {\n    this.emitter.off(event as string, callback);\n  }\n}\n\n'})}),"\n",(0,r.jsx)(n.p,{children:"And here is how you would use this:"}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:'const partyEmitter = new EventEmitterNode<{\n  "party-started": { name: string; time: string };\n  "party-ended": { bummer: boolean };\n}>();\n\npartyEmitter.onEvent("party-started", (data) => {\n  console.log(`Party started: ${data.name} at ${data.time}`);\n});\n\npartyEmitter.sendEvent("party-started", { name: "Party", time: "2025-05-24" });\n'})}),"\n",(0,r.jsxs)(n.div,{className:"markdown-alert markdown-alert-note",dir:"auto",children:["\n",(0,r.jsxs)(n.p,{className:"markdown-alert-title",dir:"auto",children:[(0,r.jsx)(n.svg,{className:"octicon",viewBox:"0 0 16 16",width:"16",height:"16","aria-hidden":"true",children:(0,r.jsx)(n.path,{d:"M0 8a8 8 0 1 1 16 0A8 8 0 0 1 0 8Zm8-6.5a6.5 6.5 0 1 0 0 13 6.5 6.5 0 0 0 0-13ZM6.5 7.75A.75.75 0 0 1 7.25 7h1a.75.75 0 0 1 .75.75v2.75h.25a.75.75 0 0 1 0 1.5h-2a.75.75 0 0 1 0-1.5h.25v-2h-.25a.75.75 0 0 1-.75-.75ZM8 6a1 1 0 1 1 0-2 1 1 0 0 1 0 2Z"})}),"NOTE"]}),"\n",(0,r.jsxs)(n.p,{children:["Event emitters are synchronous by nature, meaning that an ",(0,r.jsx)(n.code,{children:"eventEmitter.emit()"})," call is blocking until the listener for that event finishes executing."]}),"\n"]}),"\n",(0,r.jsx)(n.h2,{id:"streams-in-node",children:"Streams in Node"}),"\n",(0,r.jsx)(n.h3,{id:"why-streams",children:"Why Streams?"}),"\n",(0,r.jsx)(n.p,{children:"In Node.js, garbage collection helps manage memory by cleaning up unused data. There are two types:"}),"\n",(0,r.jsxs)(n.ul,{children:["\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.strong,{children:"Scavenge"}),": This is a quick, lightweight cleanup that handles short-term memory. It runs frequently and cleans up small amounts of memory without stopping the Node.js process, so it has minimal impact on performance."]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.strong,{children:"Mark-sweep"}),": This is a more thorough cleanup that stops the Node.js process to identify and remove larger amounts of unused memory. Because it pauses the application, it can affect performance more noticeably."]}),"\n"]}),"\n",(0,r.jsxs)(n.p,{children:["To run a node app with garbage collection information, use the ",(0,r.jsx)(n.code,{children:"--trace_gc"})," flag:"]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-bash",children:"node --trace_gc app.js\n"})}),"\n",(0,r.jsxs)(n.p,{children:["In this HTTP server example, we are loading entire files in memory. This causes a ",(0,r.jsx)(n.strong,{children:"mark-sweep"})," garbage collection to clean up large objects at a time, thus increasing memory and CPU utilization of the program, which is bad."]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:'import { createServer } from "node:http"\nimport fs from "fs/promises"\n\n\nconst server = createServer(async (req, res) => {\n\t// read large file into memory then send it all at once.\n\tconst data = await fs.readFile("largevideo.mp4")\n    res.writeHead(200, { "Content-Type" : "video/mp4" });\n    res.end(data)\n})\n\n// 4. start the server on a port\nserver.listen(3000)\n'})}),"\n",(0,r.jsxs)(n.p,{children:["But when we create a readable stream with the file as the producer using the ",(0,r.jsx)(n.code,{children:"fs.createReadStream(filename)"})," method, and then pipe it to the response, we send back small chunks of data at a time to the client, improving garbage collection tactics resulting in only ",(0,r.jsx)(n.strong,{children:"scavenge"})," garbage collection calls."]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:'import { createServer } from "node:http"\nimport fs from "fs/promises"\n\n\nconst server = createServer(async (req, res) => {\n    res.writeHead(200, { "Content-Type" : "video/mp4" });\n    \n    // create readable stream from file producer and pipe it to response\n    fs.createReadStream(\'largevideo.mp4\').pipe(res)\n})\n\n// 4. start the server on a port\nserver.listen(3000)\n'})}),"\n",(0,r.jsx)(n.h3,{id:"stream-basics",children:"Stream basics"}),"\n",(0,r.jsxs)(n.p,{children:["Streams are a way of handling data in node in a way that is memory-efficient. It allows you to ",(0,r.jsx)(n.strong,{children:"stream"})," data, never having it all in memory. This is useful for handling large files, or large amounts of data."]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-javascript",children:'import fs from "fs";\nconst readableStream = fs.createReadStream("file.txt");\nconst writableStream = fs.createWriteStream("file2.txt");\n\nreadableStream\n  .pipe(writableStream)\n  .on("finish", () => console.log("Done!"))\n  .on("error", (error) => console.log(error.message));\n'})}),"\n",(0,r.jsx)(n.p,{children:"The primary benefits of using streams are:"}),"\n",(0,r.jsxs)(n.ul,{children:["\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.strong,{children:"Memory Efficiency:"})," Process large files or data sets in chunks, avoiding the need to load the entire content into memory. This prevents out-of-memory errors and keeps your application lightweight."]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.strong,{children:"Time Efficiency:"})," Start processing data as soon as the first chunk arrives, rather than waiting for the entire resource to be available. This can lead to faster perceived performance."]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.strong,{children:"Composability:"}),' Streams are designed to be easily "piped" together, creating a clear and efficie
1nt data pipeline. This makes complex data processing workflows more manageable and readable.']}),"\n"]}),"\n",(0,r.jsx)(n.p,{children:"There are 4 types of streams;"}),"\n",(0,r.jsxs)(n.ul,{children:["\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.strong,{children:"readable streams"}),": read data in chunks, can pipe to other writable streams"]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.strong,{children:"writable streams"}),": For writing data that is accessed from a readable stream."]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.strong,{children:"duplex streams"}),": Streams that are both readable and writable"]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.strong,{children:"transform streams"}),": duplex streams that modify data in a pipeline"]}),"\n"]}),"\n",(0,r.jsxs)(n.table,{children:[(0,r.jsx)(n.thead,{children:(0,r.jsxs)(n.tr,{children:[(0,r.jsx)(n.th,{children:"Stream Type"}),(0,r.jsx)(n.th,{children:"Description"}),(0,r.jsx)(n.th,{children:"Examples"})]})}),(0,r.jsxs)(n.tbody,{children:[(0,r.jsxs)(n.tr,{children:[(0,r.jsx)(n.td,{children:(0,r.jsx)(n.strong,{children:"Readable"})}),(0,r.jsx)(n.td,{children:"Emits data (you consume it)"}),(0,r.jsxs)(n.td,{children:[(0,r.jsx)(n.code,{children:"fs.createReadStream()"}),", ",(0,r.jsx)(n.code,{children:"http.IncomingMessage"})]})]}),(0,r.jsxs)(n.tr,{children:[(0,r.jsx)(n.td,{children:(0,r.jsx)(n.strong,{children:"Writable"})}),(0,r.jsx)(n.td,{children:"Accepts data (you write to it)"}),(0,r.jsxs)(n.td,{children:[(0,r.jsx)(n.code,{children:"fs.createWriteStream()"}),", ",(0,r.jsx)(n.code,{children:"http.ServerResponse"})]})]}),(0,r.jsxs)(n.tr,{children:[(0,r.jsx)(n.td,{children:(0,r.jsx)(n.strong,{children:"Duplex"})}),(0,r.jsx)(n.td,{children:"Read + write"}),(0,r.jsx)(n.td,{children:(0,r.jsx)(n.code,{children:"net.Socket"})})]}),(0,r.jsxs)(n.tr,{children:[(0,r.jsx)(n.td,{children:(0,r.jsx)(n.strong,{children:"Transform"})}),(0,r.jsxs)(n.td,{children:["Duplex that ",(0,r.jsx)(n.strong,{children:"modifies"})," data on the fly"]}),(0,r.jsx)(n.td,{children:(0,r.jsx)(n.code,{children:"zlib.createGzip()"})})]})]})]}),"\n",(0,r.jsx)(n.h3,{id:"readable-and-writable-streams",children:"Readable and writable streams"}),"\n",(0,r.jsx)(n.h4,{id:"basics-of-readable-streams",children:"Basics of readable streams"}),"\n",(0,r.jsxs)(n.p,{children:["Readable streams inherit from the ",(0,r.jsx)(n.code,{children:"EventEmitter"})," class, so they have 4 events they can listen to:"]}),"\n",(0,r.jsxs)(n.ul,{children:["\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"'data'"}),": Emitted when a chunk of data is available."]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"'end'"}),": Emitted when there is no more data to consume."]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"'error'"}),": Emitted if an error occurs while reading."]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"'close'"}),": Emitted when the underlying resource (e.g., file descriptor) has been closed."]}),"\n"]}),"\n",(0,r.jsxs)(n.p,{children:["The main use case for using a readable stream is to pipe it to a writable stream, like ",(0,r.jsx)(n.code,{children:"process.stdout"})]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:'import fs from "fs";\n\nconst readStream = fs.createReadStream("src/streams/roadmaps.json");\n\n// pipe to a writable stream\nreadStream.pipe(process.stdout);\n'})}),"\n",(0,r.jsxs)(n.p,{children:["The basic syntax for creating a read stream using the ",(0,r.jsx)(n.code,{children:"fs.createReadStream()"})," method is like so, where we create a readable stream with a specific file as a producer."]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:"const readableStream = fs.createReadStream(filepath, options)\n"})}),"\n",(0,r.jsxs)(n.p,{children:["Here are the options you have access to in the ",(0,r.jsx)(n.code,{children:"options"})," object:"]}),"\n",(0,r.jsxs)(n.ul,{children:["\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"encoding"}),": the file encoding, like ",(0,r.jsx)(n.code,{children:'"utf8"'})]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"highWatermark"}),": the size of each chunk in the stream, in bytes"]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"start"}),": Starting byte position to read fro"]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"end"}),": Ending byte position to read to"]}),"\n"]}),"\n",(0,r.jsxs)(n.p,{children:["By default, the ",(0,r.jsx)(n.code,{children:"highWatermark"})," property when creating a readable stream reads 64kb chunks at a time. We can make it more memory-efficient by making each chunk smaller, like 1024 bytes only."]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:'const readable = fs.createReadStream(file, {\n    encoding: "utf8",\n    highWaterMark: 1024, // 1 KB chunks\n  });\n\n  readable.on("data", (chunk) =>
1 {\n    console.log("Received chunk:", chunk.length);\n  });\n\n  readable.on("end", () => {\n    console.log("Done reading file");\n  });\n'})}),"\n",(0,r.jsx)(n.h4,{id:"consuming-readable-streams",children:(0,r.jsx)(n.strong,{children:"consuming readable streams"})}),"\n",(0,r.jsx)(n.p,{children:"Readable streams operate in two modes:"}),"\n",(0,r.jsxs)(n.ol,{children:["\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.strong,{children:"Flowing mode"}),": Data is read automatically and provided as quickly as possible"]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.strong,{children:"Paused mode"}),": The\xa0",(0,r.jsx)(n.code,{children:"read()"}),"\xa0method must be called explicitly to get chunks"]}),"\n"]}),"\n",(0,r.jsx)(n.p,{children:"There are three ways to consume readable streams in node:"}),"\n",(0,r.jsxs)(n.ol,{children:["\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.strong,{children:"async iteration"}),": you can consume streams with the ",(0,r.jsx)(n.code,{children:"for await"})," syntax, which activates ",(0,r.jsx)(n.strong,{children:"flowing mode"}),"."]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.strong,{children:"event listeners"}),": By setting event listeners on the stream for the ",(0,r.jsx)(n.code,{children:'"data"'})," event, you activate ",(0,r.jsx)(n.strong,{children:"flowing mode"}),"."]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsxs)(n.strong,{children:["using ",(0,r.jsx)(n.code,{children:"readable.read()"})]}),": This imperatively reads chunks at a time and activates ",(0,r.jsx)(n.strong,{children:"paused mode"}),"."]}),"\n"]}),"\n",(0,r.jsxs)(n.p,{children:["To manually switch to paused mode, you can use the ",(0,r.jsx)(n.code,{children:"readableStream.pause()"})," method."]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:"readableStream.pause();\n"})}),"\n",(0,r.jsx)(n.p,{children:"You can now use async iteration to consume readable streams in node:"}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:"const fs = require('fs');\n\nasync function readFileInChunks() {\n  const stream = fs.createReadStream('./bigfile.txt', { encoding: 'utf8' });\n\n  for await (const chunk of stream) {\n    console.log('Chunk:', chunk.length);\n  }\n}\nreadFileInChunks();\n"})}),"\n",(0,r.jsxs)(n.p,{children:["Here's an example of where we can pipe to a node script or accept any amount of standard input by consuming the ",(0,r.jsx)(n.code,{children:"process.stdin"})," stream:"]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:"const content = []\nprocess.stdin.on('data', (chunk) => {\n    content.push(chunk.toString())\n})\n\nprocess.stdin.on('end', () => {\n    console.log('finished stream with ' + content.length + ' chunks.')\n})\n"})}),"\n",(0,r.jsx)(n.p,{children:"ANd this was what was run:"}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-bash",children:"echo 'bruh' | deno run -A streams.ts\n"})}),"\n",(0,r.jsx)(n.h4,{id:"basics-of-writable-streams",children:"Basics of Writable streams"}),"\n",(0,r.jsxs)(n.p,{children:["To create a writable stream from a file, you use the ",(0,r.jsx)(n.code,{children:"fs.createWriteStream()"})," method:"]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:"const writable = fs.createWriteStream(filepath, options)\n"})}),"\n",(0,r.jsxs)(n.p,{children:["Here are the options you have access to in the ",(0,r.jsx)(n.code,{children:"options"})," object:"]}),"\n",(0,r.jsxs)(n.ul,{children:["\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"encoding"}),": the file encoding, like ",(0,r.jsx)(n.code,{children:'"utf8"'})]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"highWatermark"}),": the size of each chunk in the stream, in bytes"]}),"\n"]}),"\n",(0,r.jsxs)(n.p,{children:["On a writable stream, you have access to the ",(0,r.jsx)(n.code,{children:"writable.write(chunk)"})," method, which writes a chunk manually to the stream. This method returns a ",(0,r.jsx)(n.code,{children:"boolean"})," that describes the status of the buffer:"]}),"\n",(0,r.jsxs)(n.ul,{children:["\n",(0,r.jsxs)(n.li,{children:["If ",(0,r.jsx)(n.code,{children:"writable.write(chunk)"})," returns ",(0,r.jsx)(n.strong,{children:"true"}),", then the buffer is OK and there is no backpressure."]}),"\n",(0,r.jsxs)(n.li,{children:["If ",(0,r.jsx)(n.code,{children:"writable.write(chunk)"})," returns ",(0,r.jsx)(n.strong,{children:"false"}),", then the buffer has backpressure and we have to relieve the backpressure."]}),"\n"]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:'const writable = fs.createWriteStream("./src/streams/output.txt", {\n\thighWaterMark: 1024,\n\tencoding: "utf8",\n});\nlet isBufferOk = false;\n\nisBufferOk = writable.write("Hello\\n");\nisBufferOk = writable.write("World\\n");\nconsole.log("isBufferOk", isBufferOk);\n\n// setus up event listener for when writable ends.\nwritable.end(() =>
1 {\n\tconsole.log("All writes are complete!");\n});\n'})}),"\n",(0,r.jsxs)(n.p,{children:["You then use ",(0,r.jsx)(n.code,{children:"writable.end()"})," to signal you're done writing and to close the writable stream automatically."]}),"\n",(0,r.jsx)(n.p,{children:(0,r.jsx)(n.strong,{children:"methods"})}),"\n",(0,r.jsxs)(n.ul,{children:["\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"writable.write(chunk)"}),": writes a chunk to the writable stream and returns a boolean where if true, there is no backpressure, and if false, there is backpressure."]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"writable.end(cb)"}),": runs a callback when you're done with all your writes, triggering the ",(0,r.jsx)(n.code,{children:"'end'"})," event on the writable stream"]}),"\n"]}),"\n",(0,r.jsx)(n.p,{children:(0,r.jsx)(n.strong,{children:"events"})}),"\n",(0,r.jsx)(n.p,{children:"Here are the 4 events you can listen to on a writable stream:"}),"\n",(0,r.jsxs)(n.ul,{children:["\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"'drain'"}),": Emitted when the internal buffer of the writable stream has been emptied and it's safe to write more data. This is crucial for backpressure handling."]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"'finish'"}),": Emitted when you invoke ",(0,r.jsx)(n.code,{children:"writable.end()"})," to signal you're done writing and all data has been flushed to the underlying system."]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"'error':"})," Emitted if an error occurs while writing."]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"'close':"})," Emitted when the underlying resource has been closed."]}),"\n"]}),"\n",(0,r.jsx)(n.p,{children:(0,r.jsx)(n.strong,{children:"built-in writable streams"})}),"\n",(0,r.jsxs)(n.p,{children:[(0,r.jsx)(n.code,{children:"process.stdout"})," and ",(0,r.jsx)(n.code,{children:"process.stderr"})," are two built-in writable streams that you can write to."]}),"\n",(0,r.jsx)(n.h4,{id:"combining-readable-and-writable-streams",children:"Combining readable and writable streams"}),"\n",(0,r.jsx)(n.p,{children:(0,r.jsx)(n.strong,{children:"level 1: piping"})}),"\n",(0,r.jsxs)(n.p,{children:["The main use case for readable and writable streams is to pipe the contents of a readable stream into a writable stream. The basic syntax is to the use ",(0,r.jsx)(n.code,{children:"readable.pipe()"})," method on a readable stream, and pass a writable stream into that."]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:"readable.pipe(writable)\n"})}),"\n",(0,r.jsx)(n.p,{children:"The rules of piping are as follows:"}),"\n",(0,r.jsxs)(n.ol,{children:["\n",(0,r.jsxs)(n.li,{children:["The ",(0,r.jsx)(n.code,{children:".pipe()"})," method only exists on readable streams, so readable streams are the only ones that can start a pipeline."]}),"\n",(0,r.jsx)(n.li,{children:"You can only pipe to transform streams and writable streams"}),"\n"]}),"\n",(0,r.jsxs)(n.p,{children:["The main problem with piping however, is that if an error occurs, the stream automatically closes and the data flow is broken. The new ",(0,r.jsx)(n.code,{children:"pipeline()"})," function handles these edge cases:"]}),"\n",(0,r.jsx)(n.p,{children:(0,r.jsx)(n.strong,{children:"level 2: pipeline"})}),"\n",(0,r.jsxs)(n.p,{children:["The basic syntax for the ",(0,r.jsx)(n.code,{children:"pipeline()"})," function from the ",(0,r.jsx)(n.code,{children:"node:stream"})," library is as follows, where you can pass in an ordered spread arguments of readable, writable, and transform streams and then for the last argument, and error callback."]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:"await pipeline(...streams, (err) => {})\n"})}),"\n",(0,r.jsxs)(n.p,{children:["Below is an example of a file read stream piping to a transform stream which then pipes to a writable stream (",(0,r.jsx)(n.code,{children:"process.stdout"}),"), but the pipeline allows for error handling."]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:"const fs = require('node:fs');\nconst { Transform, pipeline } = require('node:stream');\n\nlet errorCount = 0;\nconst upper = new Transform({\n  transform: function (data, enc, cb) {\n    if (errorCount === 10) {\n      return cb(new Error('BOOM!'));\n    }\n    errorCount++;\n    this.push(data.toString().toUpperCase());\n    cb();\n  },\n});\n\nconst readStream = fs.createReadStream(__filename, { highWaterMark: 1 });\nconst writeStream = process.stdout;\n\nreadStream.on('close', () => {\n  console.log('Readable stream closed');\n});\n\nupper.on('close', () => {\n  console.log('\\nTransform stream closed');\n});\n\nwriteStream.on('close', () => {\n  console.log('Writable stream closed');\n});\n\npipeline(readStream, upper, writeStream, err => {\n  if (err) {\n    return console.error('Pipeline error:', err.message);\n  }\n  console.log('Pipeline succeeded');\n});\n"})}),"\n",(0,r.jsx)(n.p,{children:(0,r.jsx)(n.strong,{children:"level 3: async pipeline"})}),"\n",(0,r.jsxs)(n.p,{children:["The asynchronous ",(0,r.jsx)(n.code,{children:"pipeline()"})," function also accepts async generators as transform streams:"]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:"const fs = require('node:fs');\nconst { pipeline } = require('node:stream/promises');\n\nasync function main() {\n  await pipeline(\n    fs.createReadStream(__filename), // 1. create read stream\n    async function* (source) {       // 2. async generator as transform stream\n      for await (let chunk of source) {\n        yield chunk.toString().toUpperCase();\n      }\n    },\n    process.stdout.  // 3. last in pipeline is writable stream\n  );\n}\n\nmain().catch(console.error);\n\n"})}),"\n",(0,r.jsx)(n.h4,{id:"handling-backpressure",children:"Handling backpressure"}),"\n",(0,r.jsxs)(n.p,{children:[(0,r.jsx)(n.strong,{children:"Backpressure"})," is a crucial concept in stream handling. It refers to the mechanism by which a slow consumer (writable stream) tells a fast producer (readable stream) to slow down, preventing the consumer's internal buffer from overflowing."]}),"\n",(0,r.jsxs)(n.div,{className:"markdown-alert markdown-alert-note",dir:"auto",children:["\n",(0,r.jsxs)(n.p,{className:"markdown-alert-title",dir:"auto",children:[(0,r.jsx)(n.svg,{className:"octicon",viewBox:"0 0 16 16",width:"16",height:"16","aria-hidden":"true",children:(0,r.jsx)(n.path,{d:"M0 8a8 8 0 1 1 16 0A8 8 0 0 1 0 8Zm8-6.5a6.5 6.5 0 1 0 0 13 6.5 6.5 0 0 0 0-13ZM6.5 7.75A.75.75 0 0 1 7.25 7h1a.75.75 0 0 1 .75.75v2.75h.25a.75.75 0 0 1 0 1.5h-2a.75.75 0 0 1 0-1.5h.25v-2h-.25a.75.75 0 0 1-.75-.75ZM8 6a1 1 0 1 1 0-2 1 1 0 0 1 0 2Z"})}),"NOTE"]}),"\n",(0,r.jsx)(n.p,{children:"Think of streams as a faucet (source) and drain (sink):"}),"\n",(0,r.jsxs)(n.ul,{children:["\n",(0,r.jsxs)(n.li,{children:["If the ",(0,r.jsx)(n.strong,{children:"sink is slow"}),", and you keep pouring, you\u2019ll overflow."]}),"\n",(0,r.jsxs)(n.li,{children:["Backpressure is a signal: ",(0,r.jsx)(n.strong,{children:'"Stop sending data until I\u2019m ready again."'})]}),"\n"]}),"\n"]}),"\n",(0,r.jsx)(n.p,{children:(0,r.jsx)(n.img,{src:"https://i.imgur.com/PZvkkGF.jpeg",alt:""})}),"\n",(0,r.jsx)(n.p,{children:"Let's take a deeper look at this analogy to define key terms"}),"\n",(0,r.jsxs)(n.ul,{children:["\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.strong,{children:"backpressure"}),": this is like the water in the hose. If it fills up to the top, then there is backpressure."]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.strong,{children:"water mark"}),": this is how much water the hose can handle flowing through it at once before there is backpressure. This is an attribute you set on the writable stream.","\n",(0,r.jsxs)(n.ul,{children:["\n",(0,r.jsx)(n.li,{children:"A low watermark means backpressure occurs more often and you have to wait for it to drain often."}),"\n",(0,r.jsx)(n.li,{children:"A high watermark means backpressure occurs less often but that means each chunk takes up more memory in the app, sort of defeating the purpose of streams."}),"\n"]}),"\n"]}),"\n"]}),"\n",(0,r.jsxs)(n.div,{className:"markdown-alert markdown-alert-note",dir:"auto",children:["\n",(0,r.jsxs)(n.p,{className:"markdown-alert-title",dir:"auto",children:[(0,r.jsx)(n.svg,{className:"octicon",viewBox:"0 0 16 16",width:"16",height:"16","aria-hidden":"true",children:(0,r.jsx)(n.path,{d:"M0 8a8 8 0 1 1 16 0A8 8 0 0 1 0 8Zm8-6.5a6.5 6.5 0 1 0 0 13 6.5 6.5 0 0 0 0-13ZM6.5 7.75A.75.75 0 0 1 7.25 7h1a.75.75 0 0 1 .75.75v2.75h.25a.75.75 0 0 1 0 1.5h-2a.75.75 0 0 1 0-1.5h.25v-2h-.25a.75.75 0 0 1-.75-.75ZM8 6a1 1 0 1 1 0-2 1 1 0 0 1 0 2Z"})}),"NOTE"]}),"\n",(0,r.jsxs)(n.p,{children:["If you use the ",(0,r.jsx)(n.code,{children:"readable.pipe(writable)"}),", the method of piping automatically handles backpressure for you."]}),"\n"]}),"\n",(0,r.jsx)(n.p,{children:"If instead you use the event listeners, you have to manually pause and resume streams like so:"}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:"readable.on('data', chunk => {\n  if (!writable.write(chunk)) {\n    readable.pause(); // backpressure signal\n  }\n});\n\nwritable.on('drain', () => {\n  readable.resume(); // sink is ready again\n});\n"})}),"\n",(0,r.jsxs)(n.p,{children:["To handle backpressure manually, the consumer side must pause the readable if the ",(0,r.jsx)(n.code,{children:"writable.write(chunk)"})," method returns false (signals backpressure), and also listen to the ",(0,r.jsx)(n.code,{children:"'drain'"})," event and resume the readable in the event handler."]}),"\n",(0,r.jsx)(n.p,{children:"To handle backpressure on the producer side, the readable stream has access to these methods:"}),"\n",(0,r.jsxs)(n.ul,{children:["\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"readable.pause()"}),": Pauses the ",(0,r.jsx)(n.code,{children:"'data'"})," event."]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"readable.resume()"}),": Resumes the ",(0,r.jsx)(n.code,{children:"'data'"})," event."]}),"\n"]}),"\n",(0,r.jsx)(n.p,{children:(0,r.jsx)(n.strong,{children:"summary"})}),"\n",(0,r.jsxs)(n.p,{children:["Backpressure is handled automatically when piping from a producer to a consumer with the ",(0,r.jsx)(n.code,{children:"readable.pipe(writable)"})," method, but if you want to stop backpressure manually, you must follow these steps:"]}),"\n",(0,r.jsx)(n.h3,{id:"custom-readable-and-writable-streams",children:"Custom readable and writable streams"}),"\n",(0,r.jsx)(n.h4,{id:"extending-readable-streams",children:"Extending readable streams"}),"\n",(0,r.jsxs)(n.p,{children:["You can create your own implementations of readable streams by extending from the ",(0,r.jsx)(n.code,{children:"Readable"})," class or creating an instance of it."]}),"\n",(0,r.jsxs)(n.p,{children:["The subclass must implement the ",(0,r.jsx)(n.code,{children:"_read()"})," method, intended to push a chunk to the stream, which gets called when consuming the stream either in flowing or paused mode."]}),"\n",(0,r.jsxs)(n.p,{children:["The ",(0,r.jsx)(n.code,{children:"this.push()"})," inherited method is what readable stream implementations use behind the scenes to push a chunk of data to the stream."]}),"\n",(0,r.jsxs)(n.ul,{children:["\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.strong,{children:"add chunks"}),": You add chunks to the readable stream through the ",(0,r.jsx)(n.code,{children:"this.push(chunk)"})," method."]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.strong,{children:"end stream"}),": You signify that the stream has ended through the ",(0,r.jsx)(n.code,{children:"this.push(null)"})," method."]}),"\n"]}),"\n",(0,r.jsx)(n.p,{children:"The superclass constructor takes in an object of options:"}),"\n",(0,r.jsxs)(n.ul,{children:["\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"encoding"}),": the type of encoding to use for the data chunks. By default, chunks are binary and are outputted as ",(0,r.jsx)(n.code,{children:"Buffer"})," instances."]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"objectMode"}),": a boolean configuring whether or not to enable ",(0,r.jsx)(n.strong,{children:"object mode"}),", where the stream produces chunks of JavaScript objects and you can consume the objects as is."]}),"\n"]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:"class ReadableImplementation extends Readable {\n  constructor() {\n    super({\n\t    encoding: 'utf-8'\n    });\n  }\n}\n"})}),"\n",(0,r.jsx)(n.p,{children:"Here's a more fleshed out example."}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:"const { Readable } = require('stream');\n\nclass Counter extends Readable {\n  constructor(options) {\n    super(options);\n    this.count = 0;\n    this.maxCount = 10;\n  }\n\n  _read() {\n    if (this.count < this.maxCount) {\n      const chunk = `Number: ${this.count++}\\n`;\n      // Push the data to the internal buffer\n      this.push(chunk);\n    } else {\n      // Signal that there's no more data\n      this.push(null);\n    }\n  }\n}\n\nconst counterStream = new Counter();
1\n\n// consume data in flowing mode through event listeners,\n// calls _read() implicitly._\ncounterStream.on('data', (chunk) => {\n  console.log(chunk.toString());\n});\n\ncounterStream.on('end', () => {\n  console.log('Counting finished.');\n});\n\ncounterStream.on('error', (err) => {\n  console.error('Counter stream error:', err);\n});\n"})}),"\n",(0,r.jsxs)(n.p,{children:["On any readable stream, you can manually read chunks of data through the ",(0,r.jsx)(n.code,{children:"readable.read()"})," method."]}),"\n",(0,r.jsx)(n.p,{children:"Here is an example of how we can create our own custom readable streams and consume them:"}),"\n",(0,r.jsxs)(n.div,{className:"markdown-alert markdown-alert-important",dir:"auto",children:["\n",(0,r.jsxs)(n.p,{className:"markdown-alert-title",dir:"auto",children:[(0,r.jsx)(n.svg,{className:"octicon",viewBox:"0 0 16 16",width:"16",height:"16","aria-hidden":"true",children:(0,r.jsx)(n.path,{d:"M0 1.75C0 .784.784 0 1.75 0h12.5C15.216 0 16 .784 16 1.75v9.5A1.75 1.75 0 0 1 14.25 13H8.06l-2.573 2.573A1.458 1.458 0 0 1 3 14.543V13H1.75A1.75 1.75 0 0 1 0 11.25Zm1.75-.25a.25.25 0 0 0-.25.25v9.5c0 .138.112.25.25.25h2a.75.75 0 0 1 .75.75v2.19l2.72-2.72a.749.749 0 0 1 .53-.22h6.5a.25.25 0 0 0 .25-.25v-9.5a.25.25 0 0 0-.25-.25Zm7 2.25v2.5a.75.75 0 0 1-1.5 0v-2.5a.75.75 0 0 1 1.5 0ZM9 9a1 1 0 1 1-2 0 1 1 0 0 1 2 0Z"})}),"IMPORTANT"]}),"\n",(0,r.jsxs)(n.p,{children:["Since readable streams are observables, they won't start emitting until you explicitly decide to read them with ",(0,r.jsx)(n.code,{children:"stream.read()"})," or setting up the event listeners."]}),"\n"]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:'import fs from "fs";\nimport { Writable, Transform, Duplex, Readable, ReadableOptions } from "stream";\nimport { pipeline } from "node:stream/promises";\n\nclass StreamUtils {\n  static createTransformStream(transformFunction: (chunk: Buffer) => Buffer) {\n    return new Transform({\n      transform(chunk, encoding, callback) {\n        const transformed = transformFunction(chunk);\n        this.push(transformed);\n        callback();\n      },\n    });\n  }\n\n  static createDelayStream(delay: number) {\n    return new Transform({\n      transform(chunk, encoding, callback) {\n        this.push(chunk);\n        setTimeout(() => {\n          callback();\n        }, delay);\n      },\n    });\n  }\n}\n\nclass CustomReadable extends Readable {\n  constructor(options: ReadableOptions) {\n    super(options);\n  }\n\n  addBinaryChunk(chunk: Buffer) {\n    this.push(chunk);\n  }\n\n  addTextChunk(chunk: string) {\n    this.push(chunk, "utf-8");\n  }\n\n  end() {\n    this.push(null);\n  }\n}\n\nclass FileReader extends CustomReadable {\n  private stream: fs.ReadStream;\n\n  constructor(file: string, options: ReadableOptions) {\n    super(options);\n    this.stream = fs.createReadStream(file, options);\n\n    this.stream.on("data", (chunk) => {\n      this.addTextChunk(chunk.toString());\n    });\n\n    this.stream.on("end", () => {\n      this.end();\n    });\n\n    this.stream.on("error", (err) => {\n      this.emit("error", err);\n    });\n  }\n\n  _read(size?: number): void {\n    // event listeners will handle the reading\n  }\n}\n\n// 1. create readable of file\n  const readable = new FileReader(file, {\n    highWaterMark: 256,\n    encoding: "utf-8",\n  });\n\n// 2. set transform streams\nconst upperCaseTransform = StreamUtils.createTransformStream((chunk) => {\n\treturn Buffer.from(chunk.toString().toUpperCase());\n});\nconst streamDelayTransform = StreamUtils.createDelayStream(150);\n\n// 3. run pipeline\n  await pipeline(\n    readable,\n    upperCaseTransform,\n    streamDelayTransform,\n    process.stdout\n  );\n'})}),"\n",(0,r.jsx)(n.h4,{id:"creating-custom-readable-streams",children:"Creating custom readable streams"}),"\n",(0,r.jsxs)(n.p,{children:["You can create custom readable streams which might be easier than inheriting it. These are the options and methods you provide when creating a ",(0,r.jsx)(n.code,{children:"Readable"})," instance:"]}),"\n",(0,r.jsxs)(n.ul,{children:["\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"objectMode"}),": if true, then you can deal with streaming objects instead of buffers or strings."]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"read()"}),": this implements ",(0,r.jsx)(n.code,{children:"_read()"})," under the hood, where you have access to ",(0,r.jsx)(n.code,{children:"this.push(chunk)"})," to add data to the stream or ",(0,r.jsx)(n.code,{children:"this.push(null)"})," to end data."]}),"\n"]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:'const createReadableStream = new Readable({\n  objectMode: true,\n  async read() {\n    try {\n      const rows = await fetchData(offset, CHUNK_SIZE); \n      if (rows.length) {\n        rows.forEach((row) => this.push(row));\n        offset += CHUNK_SIZE;\n      } else {\n        this.push(null); // End of stream\n      }\n    } catch (err) {\n      this.destroy(err);\n    }\n  },\n});\n\n// Log the output from the stream\ncreateReadableStream\n  .on("data", console.log)\n  .on("end", () => console.log("Stream ended"))\n  .on("error", (err) => console.error("Stream error:", err));\n'})}),"\n",(0,r.jsx)(n.h4,{id:"extending-writable-streams",children:"Extending writable streams"}),"\n",(0,r.jsxs)(n.p,{children:["You can create your own implementations of writable streams by extending from the ",(0,r.jsx)(n.code,{children:"Writable"})," class or from instantiating it"]}),"\n",(0,r.jsxs)(n.ul,{children:["\n",(0,r.jsxs)(n.li,{children:["The subclass must implement the ",(0,r.jsx)(n.code,{children:"_write()"}),", ",(0,r.jsx)(n.code,{children:"_construct()"}),", and ",(0,r.jsx)(n.code,{children:"_destroy()"})," methods"]}),"\n",(0,r.jsxs)(n.li,{children:["You add chunks to the readable stream through the ",(0,r.jsx)(n.code,{children:"this.push(chunk)"})," method."]}),"\n",(0,r.jsxs)(n.li,{children:["You signify that the stream has ended through the ",(0,r.jsx)(n.code,{children:"this.push(null)"})," method."]}),"\n"]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:"const { Writable } = require('stream');\n\nclass UppercaseWriter extends Writable {\n constructor(options) {\n   // Calls the stream.Writable() constructor\n   super(options);\n }\n\n _write(chunk, encoding, callback) {\n   // Convert the chunk to uppercase\n   const uppercase = chunk.toString().toUpperCase();\n\n   // Print to the console\n   process.stdout.write(uppercase);\n\n   // Call callback to indicate we're ready for the next chunk\n   callback();\n }\n}\n\n// Create an instance of our custom stream\nconst uppercaseWriter = new UppercaseWriter();\n\n// Write data to it\nuppercaseWriter.write('Hello, ');\nuppercaseWriter.write('world!\\n');\nuppercaseWriter.write('This text will be uppercase.\\n');\nuppercaseWriter.end();\n\nuppercaseWriter.on('finish', () => {\n console.log('All data has been processed');\n});\n"})}),"\n",(0,r.jsx)(n.h4,{id:"creating-custom-writable-streams",children:"Creating custom writable streams"}),"\n",(0,r.jsxs)(n.p,{children:["Instantiating the ",(0,r.jsx)(n.code,{children:"Writable"})," class with these options creates a custom writable stream:"]}
1),"\n",(0,r.jsxs)(n.ul,{children:["\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"objectMode"}),": a boolean, whether or not to enable object mode."]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"write(chunk, encoding, callback)"}),": calls and overrides the ",(0,r.jsx)(n.code,{children:"_write()"})," method under the hood, which is needed to create a writable. Here are the arguments:","\n",(0,r.jsxs)(n.ul,{children:["\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"chunk"}),": the current chunk to consume"]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"encoding"}),": the encoding of the chunk"]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"callback()"}),": invoking this method signals that you're ready to process the next chunk."]}),"\n"]}),"\n"]}),"\n"]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:'import { Writable } from "node:stream";\n\n// Create a writable stream function\nconst createWritableStream = () => {\n  return new Writable({\n    objectMode: true,\n    write(chunk, encoding, callback) {\n      console.log(chunk);\n      callback();\n    },\n  });\n};\n'})}),"\n",(0,r.jsx)(n.h4,{id:"object-mode",children:(0,r.jsx)(n.strong,{children:"object mode"})}),"\n",(0,r.jsx)(n.hr,{}),"\n",(0,r.jsxs)(n.p,{children:[(0,r.jsx)(n.strong,{children:"object mode"})," is when you want to deal with streaming javascript objects, which can be useful for having efficient application logic. You can read and write objects by creating readable and writable streams with the ",(0,r.jsx)(n.code,{children:"objectMode"})," property set to true:"]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:"const { Readable } = require('node:stream');\n\n// creates stream of object\nconst readable = new Readable({\n  objectMode: true,\n  read() {\n    this.push({ hello: 'world' });\n    this.push(null);\n  },\n});\n\n// consumes streams of objects\nconst writable = new Writable({\n    objectMode: true,\n    write(chunk, encoding, callback) {\n      console.log(chunk);\n      callback();\n    },\n});\n\nreadable.pipe(writable)\n"})}),"\n",(0,r.jsx)(n.p,{children:"Here's a wrapper around creating obejct streams:"}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:'class ObjectStreamManager<T extends Record<string, any>> {\n  private writable: Writable;\n  private readable: Readable;\n  constructor({\n    addChunks,\n    onWriteChunk,\n    shouldEndStream,\n  }: {\n    addChunks: () => T[];\n    shouldEndStream: () => boolean;\n    onWriteChunk: (chunk: T) => void;\n  }) {\n    this.readable = new Readable({\n      objectMode: true,\n      read(size) {\n        const chunks = addChunks();\n        chunks.forEach((chunk) => this.push(chunk));\n        if (shouldEndStream()) {\n          this.push(null);\n        }\n      },\n    });\n\n    this.writable = new Writable({\n      objectMode: true,\n      write(chunk, encoding, callback) {\n        onWriteChunk(chunk);\n        callback();\n      },\n    });\n  }\n\n  consume() {\n    this.readable.pipe(this.writable);\n  }\n}\n\nfunction* createDogs(num: number) {\n  for (let i = 0; i < num; i++) {\n    yield {\n      name: `dog ${i}`,\n      breed: "golden lab",\n    };\n  }\n}\n\nconst objectStream = new ObjectStreamManager({\n  addChunks: () => {\n    return [...createDogs(20)];\n  },\n  shouldEndStream: () => true,\n  onWriteChunk: (chunk) => {\n    console.log(chunk);\n  },\n});\n\nobjectStream.consume();\n'})}),"\n",(0,r.jsx)(n.h3,{id:"transform-streams-and-duplex-streams",children:"Transform Streams and duplex streams"}),"\n",(0,r.jsx)(n.h4,{id:"duplex-streams",children:"Duplex streams"}),"\n",(0,r.jsxs)(n.p,{children:[(0,r.jsx)(n.strong,{children:"Duplex streams"})," are middlemen streams that sit between a readable stream and a writable stream on a pipeline, where they are neither producer or consumer, but rather they just access the chunks of data as the flow in."]}),"\n",(0,r.jsx)(n.p,{children:"Here is how a pipeline with a duplex stream works:"}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{children:"readable -> duplex -> writable\n"})}),"\n",(0,r.jsxs)(n.ol,{children:["\n",(0,r.jsx)(n.li,{children:"Readable stream produces and pipes data to duplex stream"}),"\n",(0,r.jsx)(n.li,{children:"Duplex stream does something with each chunk of data as it comes in, either modifying it or reading it"}),"\n",(0,r.jsx)(n.li,{children:"Writable stream consumes data, ending the pipeline."}),"\n"]}),"\n",(0,r.jsx)(n.p,{children:"There are two types of duplex streams:"}),"\n",(0,r.jsxs)(n.ul,{children:["\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.strong,{children:"duplex streams"}),": A passthrough stream that simply reads chunks of data as they come in from a producer."]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.strong,{children:"transform streams"}),": A special type of duplex stream that performs transformations on chunks of data as they come in from a producer."]}),"\n"]}),"\n",(0,r.jsxs)(n.p,{children:["Here's an example of creating a simple duplex stream with a ",(0,r.jsx)(n.code,{children:"PassThrough"})," utility class from the ",(0,r.jsx)(n.code,{children:"node:stream"})," library:"]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:"import { PassThrough } from 'node:stream'\n\nconst loggerStream = new PassThrough()\nloggerStream.on('data', (chunk) => { \n    console.log(`chunk size: ${chunk.toString().length}`)\n})\n\nprocess.stdin  // readable\n\t\t.pipe(loggerStream) // duplex\n\t\t.pipe(process.stdout) // writable\n"})}
1),"\n",(0,r.jsx)(n.h4,{id:"creating-custom-duplex-streams",children:"Creating custom duplex streams"}),"\n",(0,r.jsxs)(n.p,{children:["You can create custom duplex streams (and even transform streams from this) by extending the ",(0,r.jsx)(n.code,{children:"Duplex"})," class from the ",(0,r.jsx)(n.code,{children:"node:stream"})," library."]}),"\n",(0,r.jsx)(n.p,{children:"This class has three methods that must be overriden:"}),"\n",(0,r.jsxs)(n.ul,{children:["\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"_read()"}),": the standard ",(0,r.jsx)(n.code,{children:"_read()"})," method to implement for readable streams, where you push chunks of data to the stream with ",(0,r.jsx)(n.code,{children:"this.push(chunk)"})," and end the producer stream with ",(0,r.jsx)(n.code,{children:"this.push(null)"})]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"_write(chunk, encoding, done)"}),": The standard ",(0,r.jsx)(n.code,{children:"_write()"})," method to implement for writable streams, where you write chunks of data at a time.","\n",(0,r.jsxs)(n.ul,{children:["\n",(0,r.jsxs)(n.li,{children:["You invoke the ",(0,r.jsx)(n.code,{children:"done()"})," callback when you are done writing that specific chunk and then are ready for the next chunk to come."]}),"\n"]}),"\n"]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"_final(done)"}),": This method is called after the writable stream portion of the duplex has finished writing all chunks.","\n",(0,r.jsxs)(n.ul,{children:["\n",(0,r.jsxs)(n.li,{children:["You invoke the ",(0,r.jsx)(n.code,{children:"done()"})," callback when you are done with all writes and you are ready to become the producer for the writable stream in the pipeline."]}),"\n"]}),"\n"]}),"\n"]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:"const { Duplex } = require('stream');\n\nclass NumberDuplex extends Duplex {\n constructor(options) {\n   super(options);\n   this.current = 0;\n   this.max = 5;\n   this.received = [];\n }\n\n _read() {\n   this.current++;\n   if (this.current <= this.max) {\n     this.push(`${this.current}\\n`);\n   } else {\n     this.push(null);\n   }\n }\n\n _write(chunk, encoding, callback) {\n   this.received.push(chunk.toString().trim());\n   callback();\n }\n\n _final(callback) {\n   console.log('Received messages:', this.received);\n   callback();\n }\n}\n\nconst duplex = new NumberDuplex();\n\n// Read side\nduplex.on('data', (chunk) => {\n console.log('Read:', chunk.toString());\n});\n\nduplex.on('end', () => {\n console.log('Read complete');\n});\n\n// Write side\nduplex.write('a');\nduplex.write('b');\nduplex.write('c');\nduplex.end();\n\n"})}),"\n",(0,r.jsx)(n.p,{children:"Here are the override methods in a bot more depth"}),"\n",(0,r.jsxs)(n.ul,{children:["\n",(0,r.jsxs)(n.li,{children:["The\xa0",(0,r.jsxs)(n.strong,{children:[(0,r.jsx)(n.code,{children:"_write()"})," method override"]}),"\xa0handles each chunk of data written to the stream. It processes the chunk (in the video, it pushes the chunk onward) and then calls a callback to signal that writing is complete. In the example, it uses a delay (throttling) before calling the callback to slow down the data flow."]}),"\n",(0,r.jsxs)(n.li,{children:["The\xa0",(0,r.jsxs)(n.strong,{children:[(0,r.jsx)(n.code,{children:"_final()"})," method override"]}),"\xa0is called when no more data will be written to the stream. It signals the end of the writable side by pushing\xa0",(0,r.jsx)(n.code,{children:"null"}),", which tells the stream pipeline that writing is finished and the stream can close properly."]}),"\n",(0,r.jsxs)(n.li,{children:["You only need to override ",(0,r.jsx)(n.code,{children:"_read()"})," if your stream needs to generate or push data to be read. Since this throttle stream just passes data through without generating new data, ",(0,r.jsx)(n.code,{children:"_read()"})," can remain empty."]}),"\n"]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:"import { PassThrough, Duplex } from 'node:stream'\n\nconst loggerStream = new PassThrough()\nloggerStream.on('data', (chunk) => { \n    console.log(`chunk size: ${chunk.toString().length}`)\n})\n\n\nclass ThrottleDuplex extends Duplex {\n    constructor(private ms: number) {\n        super()\n    }\n\n    _write(chunk, encoding, done) {\n        this.push(chunk)\n        setTimeout(done, this.ms)\n    }\n\n    _read() {}\n\n    _final(done) {\n        console.log('done writing all chunks')\n        done()\n    }\n}\n\nconst throttle = new ThrottleDuplex(250)\nprocess.stdin\n    .pipe(throttle)\n    .pipe(loggerStream)\n    .pipe(process.stdout)\n\n"})}),"\n",(0,r.jsx)(n.h4,{id:"creating-transform-streams",children:"Creating transform streams"}),"\n",(0,r.jsx)(n.p,{children:"Transform streams let you hook into readable streams and modify data in the pipeline. If a readable streams pipes to a transform stream, you can perform modifications to each chunk in a memory-efficient way."}),"\n",(0,r.jsxs)(n.p,{children:["To create a transform stream, you create a ",(0,r.jsx)(n.code,{children:"Transform"})," stream instance like so:"]}),"\n",(0,r.jsxs)(n.ol,{children:["\n",(0,r.jsxs)(n.li,{children:["Override the ",(0,r.jsx)(n.code,{children:"transform()"})," callback on a chunk"]}),"\n",(0,r.jsxs)(n.li,{children:["Use the ",(0,r.jsx)(n.code,{children:"this.push(chunk)"})," method to stream the new transformed chunks"]}),"\n",(0,r.jsxs)(n.li,{children:["Invoke ",(0,r.jsx)(n.code,{children:"callback()"})," to ensure that the modifications on the chunk are finished."]}),"\n"]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:"const transformStream = new Transform({\n\ttransform(chunk, encoding, callback) {\n\t\t// perform modifications on chunk\n\t\tthis.push(transformedChunk)\n\t\tcallback()\n\t}\n})\n"})}
1),"\n",(0,r.jsx)(n.p,{children:"Then as shown in the example below, all we do is pipe the readable stream to the transform stream (which returns a readable stream), and then pipe that to a writable stream."}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:'import { Writable, Transform, Duplex } from "stream";\n\n  const readable = fs.createReadStream(file, {\n    encoding: "utf8",\n    highWaterMark: 1024, // 1 KB chunks\n    autoClose: true,\n  });\n\n  const upperCase = new Transform({\n    transform(chunk, encoding, callback) {\n      const transformed = chunk.toString().toUpperCase();\n      this.push(transformed);\n      callback();\n    },\n  });\n\n  readable.pipe(upperCase).pipe(process.stdout);\n'})}),"\n",(0,r.jsxs)(n.p,{children:["You can also create an artificial delay with streams by delaying and invoking the ",(0,r.jsx)(n.code,{children:"callback()"})," in a ",(0,r.jsx)(n.code,{children:"setTimeout"}),":"]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:"const streamDelayTransform = new Transform({\n    transform(chunk, encoding, callback) {\n      this.push(chunk);\n      setTimeout(() => {\n        callback();\n      }, 150);\n    },\n  });\n"})}),"\n",(0,r.jsx)(n.p,{children:"Here's a utility class for creating transform streams:"}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:'import fs from "fs";\nimport { Writable, Transform, Duplex, Readable, ReadableOptions } from "stream";\nimport { pipeline } from "node:stream/promises";\n\nclass StreamUtils {\n  static createTransformStream(transformFunction: (chunk: Buffer) => Buffer) {\n    return new Transform({\n      transform(chunk, encoding, callback) {\n        const transformed = transformFunction(chunk);\n        this.push(transformed);\n        callback();\n      },\n    });\n  }\n\n  static createDelayStream(delay: number) {\n    return new Transform({\n      transform(chunk, encoding, callback) {\n        this.push(chunk);\n        setTimeout(() => {\n          callback();\n        }, delay);\n      },\n    });\n  }\n}\n'})}),"\n",(0,r.jsx)(n.h4,{id:"stream-compression-with-gzip",children:"Stream compression with GZIP"}),"\n",(0,r.jsxs)(n.p,{children:["You can use the ",(0,r.jsx)(n.code,{children:"zlib"})," library which has utility methods for transform streams that compress and decompress streams using the GZIP algorithm."]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-ts",children:'\nimport zlib from "zlib";\n\n  const readable = fs.createReadStream(file, {\n    encoding: "utf8",\n    highWaterMark: 256,\n  });\n\n  const compressed = readable.pipe(zlib.createGzip());\n  const decompressed = compressed.pipe(zlib.createGunzip());\n  decompressed.pipe(process.stdout);\n'})}),"\n",(0,r.jsxs)(n.ul,{children:["\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"zlib.createGzip()"}),": a transform stream that compresses chunks using the GZIP algorithm."]}),"\n",(0,r.jsxs)(n.li,{children:[(0,r.jsx)(n.code,{children:"zlib.createGunzip()"}),": a transform stream that decompresses chunks using the GZIP algorithm."]}),"\n"]})]})}function h(e={}){const{wrapper:n}={...(0,t.R)(),...e.components};return n?(0,r.jsx)(n,{...e,children:(0,r.jsx)(d,{...e})}):d(e)}},8453:(e,n,s)=>{s.d(n,{R:()=>i,x:()=>o});var r=s(6540);const t={},a=r.createContext(t);function i(e){const n=r.useContext(a);return r.useMemo((function(){return"function"==typeof e?e(n):{...n,...e}}),[n,e])}function o(e){let n;return n=e.disableParentContext?"function"==typeof e.components?e.components(t):e.components||t:i(e.components),r.createElement(a.Provider,{value:n},e.children)}}}]);

Line numbers count LF bytes from the start of the resource, as the search results do. Vendor segments are library code the classifier recognised; they are stored but not indexed. Bytes are shown as Latin1 characters, one per byte.