Luettavien, kirjoitettavien, kaksisuuntaisten ja muuntavien streamien sisäiset rakenteet
Ymmärrä neljä stream-tyyppiä sekä se, miten sisäinen puskuri ja highWaterMark ohjaavat niiden toimintaa.
Luettavien, kirjoitettavien, kaksisuuntaisten ja muuntavien streamien sisäiset rakenteet on ilmainen Node.js-taustakehityksen bootcamp-oppitunti CoddyKitissä. Tämä on oppitunti 1/4. Voit lukea koko oppitunnin alta ilmaiseksi ja harjoitella sen jälkeen käytännössä selaimessa sisäänrakennetulla koodieditorilla ja ympäri vuorokauden käytettävissä olevan tekoälytuutorin avulla. Oppitunti kuuluu Node.js-taustakehityksen bootcamp-oppimispolkuun, ja edistymisesi synkronoituu verkon ja CoddyKit-sovelluksen välillä. Node.js-taustakehityksen bootcamp-kurssilla on yhteensä 4 oppituntia.
Miksi streamit ovat olemassa
Node.js-streamien avulla dataa voidaan käsitellä pala kerrallaan sen sijaan, että kaikki ladattaisiin kerralla muistiin. Tämä on välttämätöntä backend-työssä, kuten suurten tiedostojen tarjoamisessa, HTTP-sisältöjen välittämisessä tai tietokantavientien putkittamisessa.
- Readable — lähde, JOSTA luetaan (tiedoston luku, HTTP-pyyntö)
- Writable — kohde, JOHON kirjoitetaan (tiedoston kirjoitus, HTTP-vastaus)
- Duplex — sekä luettava että kirjoitettava, toisistaan riippumattomat kanavat (TCP-socket)
- Transform — Duplex, jossa tuloste on syötteen funktiona (gzip, salaus)
Jokaisen taustalla on sisäinen puskuri, jota ohjaa yksi luku: highWaterMark.
Sisäinen puskuri ja highWaterMark
Jokainen stream ylläpitää sisäistä puskuria kohteessa _readableState tai _writableState. highWaterMark (HWM) on kynnysarvo, ei ehdoton raja, jonka kohdalla stream ilmoittaa puskuroineensa "riittävästi".
- Tavujonojen oletus-HWM on 16 KB (16384 tavua).
- Object mode -tilassa HWM laskee objekteja, ja oletusarvo on 16.
Kun Readable-streamin puskuri täyttyy HWM-arvoon asti, se lakkaa noutamasta dataa lähteestä. Kun Writablen puskuri ylittää HWM-arvon, write() palauttaa arvon false — tätä signaalia kutsutaan nimellä backpressure.
const fs = require('fs');
const rs = fs.createReadStream('/etc/hostname', { highWaterMark: 4 });
console.log('configured HWM:', rs.readableHighWaterMark);
rs.on('data', (chunk) => {
console.log('chunk of', chunk.length, 'bytes:', JSON.stringify(chunk.toString()));
});
rs.on('end', () => console.log('done'));Readable-streamit: virtaava ja keskeytetty tila
Readable toimii toisessa kahdesta tilasta:
- Paused (oletus): datan noutamiseksi on kutsuttava
read()-metodia erikseen. - Flowing: dataa työnnetään teille
'data'-tapahtumien kautta niin nopeasti kuin sitä saapuu.
'data'-kuuntelijan liittäminen tai .pipe()-kutsun tekeminen vaihtaa streamin flowing-tilaan. Kutsu .pause() vaihtaa sen takaisin. Tämän ymmärtäminen on olennaista muistin käytön hallinnassa.
const { Readable } = require('stream');
const r = Readable.from(['a', 'b', 'c']);
// Paused mode: pull explicitly
r.on('readable', () => {
let chunk;
while ((chunk = r.read()) !== null) {
console.log('pulled:', chunk);
}
});
r.on('end', () => console.log('stream finished'));Mukautetun Readablen toteuttaminen
Rakentaaksenne oman lähteen laajentakaa Readable-luokkaa ja toteuttakaa _read(size). Kutsukaa sen sisällä this.push(chunk)-metodia datan syöttämiseksi puskuriin ja this.push(null)-metodia ilmoittamaan streamin päättymisestä (EOF).
Ratkaiseva yksityiskohta on tämä: kun push() palauttaa arvon false, sisäinen puskuri on saavuttanut HWM-arvon. Hyvin toimiva tuottaja lopettaa datan työntämisen, kunnes _read kutsutaan uudelleen.
const { Readable } = require('stream');
class Counter extends Readable {
constructor(max) {
super({ objectMode: true, highWaterMark: 2 });
this.max = max;
this.current = 1;
}
_read() {
if (this.current > this.max) {
this.push(null); // EOF
return;
}
const keepGoing = this.push({ n: this.current++ });
console.log('pushed, buffer wants more:', keepGoing);
}
}
Readable.from([]); // noop
const c = new Counter(5);
c.on('data', (obj) => console.log('consumed:', obj.n));
c.on('end', () => console.log('all consumed'));Writable-streamit ja write()-palautusarvo
Writable puskuroi saapuvat palat ja tyhjentää puskurin komennolla _write(chunk, encoding, callback). Teidän TÄYTYY kutsua callback-takaisinkutsua, kun jokainen pala on käsitelty — näin stream tietää tyhjentää puskuriaan ja hyväksyä lisää dataa.
write()-metodin palautusarvo on backpressure-signaali:
true— puskuri on HWM-arvon alapuolella, jatkakaa kirjoittamista.false— puskuri on HWM-arvossa tai sen yläpuolella, teidän PITÄISI lopettaa ja odottaa'drain'-tapahtumaa.
const { Writable } = require('stream');
class SlowSink extends Writable {
constructor() {
super({ highWaterMark: 8 });
}
_write(chunk, enc, cb) {
console.log('writing', chunk.length, 'bytes');
setTimeout(cb, 50); // simulate slow I/O
}
}
const sink = new SlowSink();
const ok = sink.write(Buffer.alloc(16));
console.log('write returned:', ok); // false -> over HWM
sink.once('drain', () => console.log('drained, safe to write again'));
sink.end(() => console.log('finished'));Backpressuren huomioiminen manuaalisesti
Jos ohitatte write()-metodilta saadun arvon false ja jatkatte kirjoittamista, sisäinen puskuri kasvaa rajatta ja prosessilta voi loppua muisti. Oikea manuaalinen toimintamalli on keskeyttää tuotanto, kunnes 'drain'-tapahtuma käynnistyy.
Käytännössä tätä tarvitsee harvoin kirjoittaa itse — .pipe() ja pipeline() hoitavat sen puolestanne — mutta mekanismin tunteminen selittää, MIKSI putkittaminen on turvallista.
const { Writable } = require('stream');
const sink = new Writable({
highWaterMark: 4,
write(chunk, enc, cb) { setTimeout(cb, 20); }
});
let i = 0;
function writeMore() {
let ok = true;
while (i < 10 && ok) {
ok = sink.write(String(i++));
}
if (i < 10) {
console.log('backpressure at i =', i, '-> wait for drain');
sink.once('drain', writeMore);
} else {
sink.end(() => console.log('done'));
}
}
writeMore();pipe(): automaattinen vuonhallinta
readable.pipe(writable) yhdistää lähteen kohteeseen ja huomioi backpressuren automaattisesti: kun kohde palauttaa arvon false, pipe kutsuu source.pause()-metodia; 'drain'-tapahtuman yhteydessä se kutsuu source.resume()-metodia.
Pelkkä .pipe()-metodi ei kuitenkaan käsittele virheitä kattavasti: jos lähteessä tapahtuu virhe, kohdetta EI suljeta automaattisesti, mikä voi vuotaa tiedostokahvoja. Suosikaa tuotannossa stream.pipeline()-metodia.
const fs = require('fs');
const zlib = require('zlib');
// gzip a file: Readable -> Transform -> Writable
fs.createReadStream('input.txt')
.pipe(zlib.createGzip())
.pipe(fs.createWriteStream('input.txt.gz'))
.on('finish', () => console.log('compressed'));Duplex-streamit: kaksi toisistaan riippumatonta kanavaa
Duplex-stream on sekä Readable että Writable, mutta puolet ovat toisistaan riippumattomia — kirjoitettu data ei ilmesty automaattisesti lukupuolelle. TCP-socket on tästä tyypillinen esimerkki: kirjoittamanne tavut lähtevät vertaisosapuolelle ja lukemanne tavut saapuvat vertaisosapuolelta.
Toteuttakaa se määrittämällä sekä _read että _write. Kummallakin puolella on oma puskurinsa ja oma highWaterMark-arvonsa.
const { Duplex } = require('stream');
class Echo extends Duplex {
constructor() {
super();
this.queue = [];
}
_write(chunk, enc, cb) {
this.queue.push(chunk.toString().toUpperCase());
cb();
}
_read() {
const item = this.queue.shift();
this.push(item !== undefined ? item : null);
}
}
const d = new Echo();
d.on('data', (c) => console.log('read side:', c.toString()));
d.write('hello');
d.write('world');
d.end();Transform-streamit: syötteestä johdettu tuloste
Transform on erityinen Duplex, jossa luettava puoli lasketaan kirjoitettavan puolen perusteella. Erillisten _read- ja _write-metodien sijaan toteutatte yhden metodin _transform(chunk, encoding, callback) ja tuotatte tulokset komennolla this.push() tai takaisinkutsun toisella argumentilla.
Valinnainen _flush(callback) suoritetaan kerran lopussa — se sopii erinomaisesti lopputietojen tuottamiseen, kuten viimeisen tarkistussumman tai sulkevan hakasulkeen lisäämiseen.
const { Transform } = require('stream');
class UpperCase extends Transform {
_transform(chunk, enc, cb) {
cb(null, chunk.toString().toUpperCase());
}
_flush(cb) {
this.push('\n-- END --\n');
cb();
}
}
const t = new UpperCase();
t.on('data', (c) => process.stdout.write(c.toString()));
t.write('node ');
t.write('streams');
t.end();Object mode ja HWM:n laskeminen
Oletusarvoisesti streamit siirtävät Buffereita/merkkijonoja ja HWM laskee tavuja. Välittäkää { objectMode: true }, jolloin stream siirtää mielivaltaisia JS-arvoja ja HWM laskee niiden sijaan objekteja.
- Tavutilan oletus-HWM: 16384 tavua
- Object mode -tilan oletus-HWM: 16 objektia
Tämä on tärkeää backend-putkissa: NDJSON:ää jäsentävä Transform voi lukea raakaa tavudataa (kirjoitettava puoli, tavutila) mutta tuottaa jäsennettyjä objekteja (luettava puoli, object mode) käyttämällä asetusta readableObjectMode.
const { Transform } = require('stream');
// Bytes in, objects out
class NdjsonParse extends Transform {
constructor() {
super({ writableObjectMode: false, readableObjectMode: true });
this.buf = '';
}
_transform(chunk, enc, cb) {
this.buf += chunk.toString();
const lines = this.buf.split('\n');
this.buf = lines.pop();
for (const line of lines) {
if (line.trim()) this.push(JSON.parse(line));
}
cb();
}
}
const p = new NdjsonParse();
p.on('data', (o) => console.log('parsed object:', o));
p.write('{"id":1}\n{"id":2}\n');
p.end();pipeline(): tuotantokäyttöön sopiva koostaminen
stream.pipeline() ketjuttaa minkä tahansa määrän streameja ja toisin kuin .pipe() välittää virheet ja siivoaa kaikki streamit (tuhoamalla ne), kun jokin niistä epäonnistuu tai päättyy. Tämä estää vuotaneet tiedostokahvat ja jumiutuneet socketit.
Promise-pohjainen muoto (require('stream/promises')) integroituu sujuvasti async/await-syntaksiin reitinkäsittelijöissä.
const { pipeline } = require('stream/promises');
const fs = require('fs');
const zlib = require('zlib');
async function gzipFile(src, dest) {
await pipeline(
fs.createReadStream(src),
zlib.createGzip(),
fs.createWriteStream(dest)
);
console.log('pipeline complete:', dest);
}
gzipFile('access.log', 'access.log.gz').catch((err) => {
console.error('pipeline failed, all streams destroyed:', err.message);
});Pikatarkistus: backpressure-signaali
Kirjoitatte silmukassa suurta tietojoukkoa mukautettuun Writable-streamiin. Haluatte välttää muistin rajoittamattoman kasvun huomioimalla backpressuren. Mikä signaali kertoo, että kirjoittaminen on lopetettava ja odotettava?
Kertaus
Ymmärrätte nyt neljä stream-tyyppiä ja niiden taustalla olevan puskurimekaniikan:
- Readable — lähde; toteuttakaa
_read, työntäkää dataa ja EOF:n merkiksipush(null); flowing- ja paused-tilat. - Writable — kohde; toteuttakaa
_writeja kutsukaa sen callback-funktiota;write()-metodin palauttamafalsetarkoittaa backpressurea, joten odottakaa'drain'-tapahtumaa. - Duplex — toisistaan riippumattomat luku- ja kirjoituskanavat, joilla kummallakin on oma puskurinsa ja HWM-arvonsa (esimerkiksi TCP-socket).
- Transform — syötteestä
_transform-metodin avulla johdettu tuloste, ja valinnaisesti_flush.
highWaterMark (oletuksena 16 kilotavua tavuina / 16 objektia) on kynnysarvo, ei ehdoton yläraja. Se määrittää, milloin puskurit ilmoittavat olevansa täynnä. Käyttäkää tuotannossa aina mieluummin pipeline()-metodia kuin paljasta .pipe()-metodia, jotta virheet välittyvät oikein ja resurssit siivotaan.
Opi JavaScript tekoälytuutorin avulla — ilmaiseksi
Kirjoita ja suorita oikeaa koodia selaimessa, saa välitöntä apua tekoälytuutorilta ympäri vuorokauden ja jatka siitä, mihin jäit, verkossa tai sovelluksessa.
- Kurssit
- 22
- Oppitunnit
- 92
Usein kysytyt kysymykset
Onko oppitunti ”Luettavien, kirjoitettavien, kaksisuuntaisten ja muuntavien streamien sisäiset rakenteet” ilmainen?
Kyllä – oppitunnin ”Luettavien, kirjoitettavien, kaksisuuntaisten ja muuntavien streamien sisäiset rakenteet” koko tekstin voi lukea täällä verkossa ilmaiseksi. Jos haluat harjoitella interaktiivisesti sisäänrakennetulla koodieditorilla ja ympäri vuorokauden käytettävissä olevan tekoälytuutorin avulla sekä avata koko Node.js-taustakehityksen bootcamp-kurssin, päivitä CoddyKit PROhon. Node.js-taustakehityksen bootcamp-kurssilla on yhteensä 4 oppituntia.
Mitä opin oppitunnilla ”Luettavien, kirjoitettavien, kaksisuuntaisten ja muuntavien streamien sisäiset rakenteet”?
Ymmärrä neljä stream-tyyppiä sekä se, miten sisäinen puskuri ja highWaterMark ohjaavat niiden toimintaa. Harjoittelet Node.js-taustakehityksen bootcamp-aihetta koodilla, jonka suoritat suoraan selaimessa. Ympäri vuorokauden käytettävissä oleva tekoälytuutori vastaa kysymyksiisi oppitunnin aikana.
Tarvitsenko kokemusta aloittaakseni Node.js-taustakehityksen bootcamp-opiskelun?
Aiempi kokemus ei ole tarpeen. CoddyKitin Node.js-taustakehityksen bootcamp-oppimispolku sopii vasta-alkajista edistyneisiin, joten voit aloittaa tästä tai alusta ja edetä omaan tahtiisi. Tämä on oppitunti 1/4.
Kuinka kauan ”Luettavien, kirjoitettavien, kaksisuuntaisten ja muuntavien streamien sisäiset rakenteet”-oppitunnin suorittaminen kestää?
Useimmat CoddyKitin oppitunnit kestävät noin 5–10 minuuttia. Jokainen oppitunti on lyhyt ja interaktiivinen, joten edistyt tasaisesti ja voit jatkaa siitä, mihin jäit – sekä verkossa että sovelluksessa.
Voinko kirjoittaa ja suorittaa koodia tällä Node.js-taustakehityksen bootcamp-oppitunnilla?
Kyllä. Jokainen Node.js-taustakehityksen bootcamp-oppitunti sisältää sisäänrakennetun koodieditorin, joten voit kirjoittaa ja suorittaa oikeaa koodia suoraan selaimessa ja saada välitöntä palautetta tekoälyltä – paikallista asennusta ei tarvita.
Kaikki tämän kurssin oppitunnit
- Luettavien, kirjoitettavien, kaksisuuntaisten ja muuntavien streamien sisäiset rakenteet
- Mukautettujen Transform-streamien toteuttaminen _transformilla ja _flushilla
- Backpressure, pipe() ja pipeline()-apuohjelma
- Asynkroniset iteraattorit ja for-await-of streamien yli