A la lliçó anterior vas consumir streams que ja existien: createReadStream produïa trossos i createWriteStream els rebia. Vas entendre la contrapressió, vas veure que pipe() la gestiona per tu i en vas descobrir el defecte: no propaga errors ni neteja el que deixa enrere.

Ara toca l'altra meitat. Escriuràs els teus propis streams —la baula intermèdia que converteix una cosa en una altra— i muntaràs la canonada completa de procés d'Escena Viva: llegir dades/vendes.csv, convertir cada línia en un objecte venda, filtrar per sala, agregar per sessió i escriure el resultat, tot amb memòria constant, tant si el fitxer és de dos kilobytes com de vuit-cents megabytes. I substituiràs pipe() per stream.pipeline, l'eina de debò: propaga errors, destrueix tota la cadena quan alguna cosa falla i té versió amb promeses. Al final veuràs la manera més llegible d'escriure transformacions, que no és una classe sinó un generador asíncron.

Contingut

  1. objectMode: streams que transporten objectes
  2. L'stream Transform: _transform i _flush
  3. Un Transform que analitza línies CSV
  4. Readable i Writable a mida
  5. stream.pipeline davant de pipe()
  6. La canonada de procés d'Escena Viva
  7. Compressió amb zlib a la canonada
  8. Readable.from() i els generadors asíncrons
  9. Errors i cancel·lació amb AbortSignal

  1. objectMode: streams que transporten objectes

Per defecte, un stream transporta bytes: Buffer o cadenes. Això és el correcte per a fitxers i sòcols, però inútil quan el que flueix per la canonada són vendes, esdeveniments o registres ja interpretats. objectMode: true canvia aquesta naturalesa:

Mode binari (per defecte) objectMode: true
Què viatja Buffer o string Qualsevol valor llevat de null
Què compta el highWaterMark Bytes (64 KB) Nombre d'objectes (16)
Trossos Es parteixen i es fusionen Cada push és un element indivisible

El mode es pot declarar per banda: new Transform({ readableObjectMode: true, writableObjectMode: false }) descriu exactament el que fa un analitzador —rep text i emet objectes—, mentre que objectMode: true a seques activa les dues bandes. I el null mereix un avís: en un stream, push(null) significa «s'ha acabat», així que no pots enviar null com a dada, ni tan sols en mode objecte; per representar l'absència d'alguna cosa, fes servir undefined o un objecte marcador.

  1. L'stream Transform: _transform i _flush

Un Transform és un stream que llegeix per una banda, fa alguna cosa i escriu per l'altra. S'implementa amb dos mètodes:

Mètode Quan es crida Per a què
_transform(tros, codificacio, callback) Per cada tros Processar i emetre amb this.push(...)
_flush(callback) En tancar-se l'entrada Emetre el que quedés pendent

La mecànica és sempre la mateixa: _transform rep un tros, fa la seva feina, emet zero o més resultats amb this.push(...) i avisa que ha acabat cridant callback(); _flush s'executa una única vegada al final, quan ja no hi entrarà res més, i és l'última oportunitat d'emetre.

Tres regles que no admeten excepció:

  • Cal cridar callback() exactament una vegada per cada _transform. Si no el crides, la canonada es penja per sempre sense cap missatge; si el crides dues vegades, Node llança ERR_MULTIPLE_CALLBACK.
  • Els errors se senyalen amb callback(error), no amb throw: una excepció llançada en un context asíncron no la captura ningú.
  • _flush és on viu l'estat acumulat. Tot el que no puguis emetre fins a veure el final —l'última línia sense salt, el total d'un agregat, el tancament d'un JSON— surt allà. Això és el que distingeix un Transform d'un simple map: pot acumular entre trossos.

  1. Un Transform que analitza línies CSV

Aquí hi ha el problema real: els trossos del fitxer no coincideixen amb les línies. Un Transform amb estat resol trossejat i conversió en un sol pas.

// src/streams/analitzar-vendes.js
// Transform: rep text de dades/vendes.csv i emet objectes venda.

const { Transform } = require('node:stream');

class AnalitzadorDeVendes extends Transform {
  #resta = '';           // tros de linia que va quedar a mitges
  #primera = true;       // per saltar la capcalera
  #descartades = 0;

  constructor(opcions = {}) {
    // Entra text, surten objectes: les dues bandes son diferents.
    super({ ...opcions, writableObjectMode: false, readableObjectMode: true });
  }

  #convertir(linia) {
    const [codiEntrada, sessioId, dataVenda, preuCentims, canal] = linia.split(',');

    if (!codiEntrada || !sessioId || Number.isNaN(Number(preuCentims))) {
      this.#descartades += 1;       // linia corrupta: s'ignora, no avorta
      return;
    }
    this.push({ codiEntrada, sessioId, dataVenda, preuCentims: Number(preuCentims), canal });
  }

  _transform(tros, codificacio, callback) {
    const linies = (this.#resta + tros).split('\n');
    this.#resta = linies.pop();     // l'ultima pot estar tallada

    for (const linia of linies) {
      if (this.#primera) { this.#primera = false; continue; }  // capcalera
      if (linia.trim() !== '') this.#convertir(linia);
    }
    callback();
  }

  _flush(callback) {
    // L'ultima linia, si el fitxer no acaba en salt.
    if (this.#resta.trim() !== '') this.#convertir(this.#resta);
    if (this.#descartades > 0) console.error(`[vendes] ${this.#descartades} descartades`);
    callback();
  }
}

module.exports = { AnalitzadorDeVendes };

Fixa't en el paper de #resta: és el problema que readline resolia per nosaltres a la lliçó anterior, ara resolt a mà perquè necessitem emetre objectes, no línies. I a _flush es processa aquesta resta pendent: sense aquesta línia, l'última venda es perdria en silenci sempre que el fitxer no acabi en salt. És un error clàssic i difícil de detectar, perquè falla en un registre d'un milió.

  1. Readable i Writable a mida

Un Readable produeix dades. S'implementa amb _read(), que Node crida quan vol més material, i a dins s'emet amb this.push(valor); quan ja no hi ha res a produir, this.push(null) senyala el final de l'stream. A la pràctica gairebé mai no cal escriure'n cap: Readable.from() (apartat 8) cobreix gairebé tots els casos amb una línia.

Un Writable consumeix dades. Implementa _write(tros, codificacio, callback), i aquest callback és la peça clau: fins que no el crides, l'stream considera que continues ocupat. Aquí és on neix la contrapressió que vas sentir a la lliçó anterior.

const { Writable } = require('node:stream');

class AcumuladorPerSessio extends Writable {
  #perSessio = new Map();

  constructor() { super({ objectMode: true }); }

  _write(venda, codificacio, callback) {
    const acumulat = this.#perSessio.get(venda.sessioId) ??
      { sessioId: venda.sessioId, entrades: 0, recaptacioCentims: 0 };

    acumulat.entrades += 1;
    acumulat.recaptacioCentims += venda.preuCentims;
    this.#perSessio.set(venda.sessioId, acumulat);
    callback();                     // a punt per al seguent
  }

  get resultat() {
    return [...this.#perSessio.values()].sort((a, b) => a.sessioId.localeCompare(b.sessioId));
  }
}

Si la feina de _write fos asíncrona —guardar a base de dades, cridar una API—, n'hi hauria prou de cridar el callback en acabar: la canonada sencera es frena sola mentrestant. Existeix a més _writev per processar diversos elements de cop, útil en inserir per lots.

  1. stream.pipeline davant de pipe()

Ja vas veure el problema de pipe(). stream.pipeline l'arregla, i la seva versió amb promeses —const { pipeline } = require('node:stream/promises')— és la que farem servir sempre.

pipe() pipeline()
Contrapressió Sí Sí
Propaga errors per la cadena No Sí: la fallada d'una etapa ho avorta tot
Destrueix els streams en fallar No: queden descriptors oberts Sí: destroy() a tots
Avisa que ha acabat Escoltant finish a mà La promesa es resol
Cancel·lació / errors No / un on('error') per stream AbortSignal / un try/catch

La diferència pràctica és enorme:

// Amb pipe: un on('error') per etapa i neteja manual.
origen.on('error', gestionar);
transformacio.on('error', gestionar);
desti.on('error', gestionar);
origen.pipe(transformacio).pipe(desti);

// Amb pipeline: una linia i un try/catch.
try {
  await pipeline(origen, transformacio, desti);
} catch (error) {
  console.error(`[canonada] fallada: ${error.message}`);
}

Regla del curs: mai pipe() en codi de producció. pipe serveix per explicar el concepte; pipeline és el que s'escriu.

  1. La canonada de procés d'Escena Viva

Amb les peces anteriors muntem el procés complet: del CSV cru a l'informe en JSON, filtrant per sala i sense carregar mai el fitxer sencer.

// src/informes/canonada-vendes.js
// CSV de vendes -> objectes -> filtre per sala -> agregacio -> JSON.

const fs = require('node:fs');
const path = require('node:path');
const { Transform } = require('node:stream');
const { pipeline } = require('node:stream/promises');

const { AnalitzadorDeVendes } = require('../streams/analitzar-vendes.js');
const { obtenirCataleg } = require('../cataleg-dades.js');
const { FITXER_VENDES, DIRECTORI_INFORMES } = require('../config/rutes.js');
const { escriureAtomic } = require('../utils/escriptura-atomica.js');

// Transform d'objecte a objecte: deixa passar nomes les sessions permeses.
// Amb la forma abreujada: opcions i metode transform sense subclasse.
const filtrarPerSessions = (permeses) => new Transform({
  objectMode: true,
  transform(venda, codificacio, callback) {
    if (permeses.has(venda.sessioId)) this.push(venda);
    callback();
  }
});

async function generarVendesPerSessio({ sala = null } = {}) {
  // El creuament amb el cataleg es fa UN cop, no per linia del CSV.
  const esdeveniments = await obtenirCataleg();
  const rellevants = sala ? esdeveniments.filter((e) => e.sala === sala) : esdeveniments;
  const sessionsPermeses = new Set(rellevants.flatMap((e) => e.sessions.map((s) => s.id)));

  const acumulador = new AcumuladorPerSessio();
  await pipeline(
    fs.createReadStream(FITXER_VENDES, { encoding: 'utf8' }),
    new AnalitzadorDeVendes(),
    filtrarPerSessions(sessionsPermeses),
    acumulador
  );

  const sessions = acumulador.resultat;
  const totals = sessions.reduce((t, s) => ({
    entrades: t.entrades + s.entrades,
    recaptacioCentims: t.recaptacioCentims + s.recaptacioCentims
  }), { entrades: 0, recaptacioCentims: 0 });

  const informe = { generatEl: new Date().toISOString(), sala: sala ?? 'totes', sessions, totals };
  const ruta = path.join(DIRECTORI_INFORMES, 'vendes-per-sessio.json');
  await escriureAtomic(ruta, JSON.stringify(informe, null, 2));
  return { ruta, ...totals };
}

module.exports = { generarVendesPerSessio };
node -e "require('./src/informes/canonada-vendes.js').generarVendesPerSessio().then(console.log)"
# { ruta: '.../informes/vendes-per-sessio.json', entrades: 1811, recaptacioCentims: 5389800 }

Les 1811 entrades quadren amb la llavor, que és la comprovació que buscàvem. I observa la forma de la canonada: quatre etapes, cadascuna amb una sola responsabilitat, encadenades en un pipeline que es llegeix de dalt a baix; canviar el filtre, afegir una validació o escriure en una altra destinació és tocar una línia. Un detall de disseny: l'AcumuladorPerSessio va al final perquè una agregació no pot emetre res fins a veure l'últim element; si necessitessis continuar la canonada després d'agregar, seria un Transform que emet el seu resultat a _flush.

  1. Compressió amb zlib a la canonada

Arxivar l'històric és intercalar una etapa més: zlib.createGzip() és un Transform com els teus.

// src/informes/arxivar-vendes.js
const fs = require('node:fs');
const zlib = require('node:zlib');
const { pipeline } = require('node:stream/promises');

async function arxivarVendes(rutaOrigen, rutaDesti) {
  await pipeline(
    fs.createReadStream(rutaOrigen),
    zlib.createGzip({ level: 9 }),      // 1 = rapid, 9 = maxima compressio
    fs.createWriteStream(rutaDesti)
  );
}
ls -l dades/vendes*
# vendes.csv 112340   |   vendes-2026.csv.gz 14208   <- 87 % menys

Aquí no s'ha indicat encoding al lector: les dades viatgen com a Buffer de principi a fi, que és el correcte per comprimir —convertir-les a text i tornar-les a bytes no aportaria res i costaria CPU—. La descompressió és la mateixa canonada al revés: n'hi ha prou d'intercalar zlib.createGunzip() entre el lector del .gz i l'analitzador per processar l'històric comprimit sense descomprimir-lo a disc. Això és el que fa potent el model: les etapes es combinen sense que cap sàpiga res de les altres, i l'analitzador no té ni idea que els seus bytes vénen d'un fitxer comprimit.

  1. Readable.from() i els generadors asíncrons

Readable.from() converteix qualsevol iterable —array, Map, generador— en un stream:

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

// Un array qualsevol, ja com a stream en mode objecte.
const vendes = Readable.from([{ sessioId: 'ses-001-1' }, { sessioId: 'ses-002-1' }]);

És utilíssim per a les proves del Mòdul 9: alimentes una canonada amb dades de mentida sense tocar el disc. Però el millor ve ara: un generador asíncron es pot fer servir directament com a etapa d'un pipeline, i és la manera més llegible d'escriure una transformació.

// La mateixa logica de l'AnalitzadorDeVendes, sense classe i sense callbacks.
// convertir(linia) torna l'objecte venda, com el metode privat d'abans.
async function* analitzarVendes(origen) {
  let resta = '';
  let primera = true;

  for await (const tros of origen) {
    const linies = (resta + tros).split('\n');
    resta = linies.pop();

    for (const linia of linies) {
      if (primera) { primera = false; continue; }
      if (linia.trim() !== '') yield convertir(linia);
    }
  }

  if (resta.trim() !== '') yield convertir(resta);   // l'equivalent de _flush
}

// Filtre: tres linies, sense classe, sense objectMode explicit.
const filtrarPerSessions = (permeses) => async function* (origen) {
  for await (const venda of origen) {
    if (permeses.has(venda.sessioId)) yield venda;
  }
};

await pipeline(
  fs.createReadStream(FITXER_VENDES, { encoding: 'utf8' }),
  analitzarVendes,
  filtrarPerSessions(sessionsPermeses),
  acumulador
);

La comparació és reveladora:

Classe Transform Generador asíncron
Línies de codi Constructor, _transform, _flush Moltes menys
objectMode Cal declarar-lo Implícit
Fi de l'entrada _flush El codi després del for await
Estat entre trossos Camp de la classe Variable local
await a dins / reutilitzable com a objecte Incòmode / sí Natural / no

Per a transformacions noves, comença sempre per un generador asíncron. Recorre a la classe Transform quan necessitis un objecte amb estat accessible des de fora —com el nostre AcumuladorPerSessio, el resultat del qual es llegeix en acabar— o quan t'hagis d'integrar amb una API que espera un stream concret.

  1. Errors i cancel·lació amb AbortSignal

Una fallada en qualsevol etapa avorta la canonada i rebutja la promesa; el que arriba al catch és l'error original, amb el seu code intacte:

try {
  await pipeline(origen, analitzarVendes, acumulador);
} catch (error) {
  // Casos previstos: s'informen i es degraden. La resta es propaga.
  if (error.code === 'ENOENT') return avisar('no existeix el fitxer de vendes');
  if (error.name === 'AbortError') return avisar('proces cancellat');
  throw error;
}

I pipeline accepta un senyal de cancel·lació, igual que fetch:

const controlador = new AbortController();
// Un informe no pot trigar mes de trenta segons.
const limit = setTimeout(() => controlador.abort(), 30_000);

try {
  await pipeline(origen, analitzarVendes, acumulador, { signal: controlador.signal });
} finally {
  clearTimeout(limit);   // no deixar temporitzadors vius
}

En avortar, pipeline destrueix tots els streams de la cadena: els descriptors s'alliberen i la promesa es rebutja amb un AbortError. És exactament la feina que a la lliçó anterior calia escriure a mà amb on('error') creuats i destroy(). Al Mòdul 4 hi tornaràs: quan un client HTTP tanca la connexió a mitja descàrrega, el correcte és avortar la canonada que generava la resposta en comptes de continuar treballant per a ningú.

Errors Comuns i Consells

  • Oblidar callback() a _transform o _write. La canonada es penja en silenci, sense error ni traça. És la fallada número u.
  • Fer servir throw dins de _transform en comptes de callback(error), o oblidar _flush, amb la qual cosa tot l'acumulat —l'última línia, el total, el tancament— es perd sense avisar.
  • No declarar objectMode en emetre objectes: Node protesta amb ERR_INVALID_ARG_TYPE perquè espera Buffer o cadena. I push(null) no és una dada: significa fi d'stream, sempre.
  • Continuar fent servir pipe(). Si un error de la primera etapa no arriba a l'última, tens descriptors perduts i fallades invisibles.
  • Consell: una etapa, una responsabilitat. Un Transform que analitza, filtra i agrega alhora és impossible de provar; tres etapes encadenades es proven per separat amb Readable.from(). I si dubtes entre classe i generador, comença pel generador: convertir-lo en classe després és fàcil, al revés no tant.

Exercicis

Exercici 1: validador a la canonada

Escriu un generador asíncron validarVendes que s'intercali entre l'analitzador i el filtre, i que comprovi cada venda: codiEntrada amb el format EV-<any>-<6 dígits>, sessioId amb la forma ses-NNN-M, preuCentims enter positiu i canal entre web, taquilla i app. Les vendes vàlides passen; les invàlides es compten i es registren per stderr amb el seu número de línia. Al final ha d'informar de quantes se'n van descartar.

Exercici 2: informe per canal comprimit

Escriu src/informes/vendes-per-canal.js amb una canonada que llegeixi dades/vendes.csv, agregui per canal (entrades, recaptació i tiquet mitjà) i escrigui informes/vendes-per-canal.json.gz comprimit, sense generar el JSON sense comprimir a disc. Pista: Readable.from([JSON.stringify(informe, null, 2)]) converteix el resultat en l'origen d'una segona canonada.

Exercici 3: canonada amb límit de temps i reintent

Escriu una funció generarAmbReintent(opcions) que executi la canonada de generarVendesPerSessio amb un AbortSignal de 10 segons i, si falla per AbortError o per un error d'E/S transitori, la reintenti fins a tres vegades fent servir reintentar i dormir del Mòdul 2. No ha de reintentar davant d'ENOENT ni davant de dades corruptes: aquestes fallades no s'arreglen repetint.

Solucions

Solució 1. El validador com a generador queda gairebé declaratiu:

const PATRO_CODI = /^EV-\d{4}-\d{6}$/;
const PATRO_SESSIO = /^ses-\d{3}-\d+$/;
const CANALS = new Set(['web', 'taquilla', 'app']);

async function* validarVendes(origen) {
  let numero = 0;
  let descartades = 0;

  for await (const venda of origen) {
    numero += 1;
    const problemes = [];
    if (!PATRO_CODI.test(venda.codiEntrada)) problemes.push('codi');
    if (!PATRO_SESSIO.test(venda.sessioId)) problemes.push('sessio');
    if (!Number.isInteger(venda.preuCentims) || venda.preuCentims <= 0) problemes.push('preu');
    if (!CANALS.has(venda.canal)) problemes.push('canal');

    if (problemes.length === 0) { yield venda; continue; }
    descartades += 1;
    console.error(`[validar] venda ${numero} descartada (${problemes.join(', ')})`);
  }

  if (descartades > 0) console.error(`[validar] ${descartades} vendes descartades`);
}

El bloc final, després del for await, és l'equivalent exacte del _flush d'una classe Transform: s'executa en exhaurir-se l'entrada. Si necessitessis parametritzar el validador, embolcalla'l en una funció que torni el generador, com el filtre de l'apartat 8.

Solució 2. Dues canonades encadenades, la segona amb Readable.from:

const acumulador = new AcumuladorPerCanal();
await pipeline(fs.createReadStream(FITXER_VENDES, { encoding: 'utf8' }), analitzarVendes, acumulador);

const informe = acumulador.resultat.map((c) => ({
  ...c,
  tiquetMitjaCentims: Math.round(c.recaptacioCentims / c.entrades)
}));

await pipeline(
  Readable.from([JSON.stringify(informe, null, 2)]),
  zlib.createGzip(),
  fs.createWriteStream(path.join(DIRECTORI_INFORMES, 'vendes-per-canal.json.gz'))
);

El JSON no toca mai el disc sense comprimir: surt de memòria com a stream, passa pel compressor i arriba al fitxer. Amb un informe petit tant li fa; amb un de centenars de megabytes és la diferència entre necessitar el doble d'espai lliure o no.

Solució 3. La clau és classificar els errors abans de decidir si es reintenta:

const NO_REINTENTABLES = new Set(['ENOENT', 'EACCES', 'DADES_CORRUPTES']);

const generarAmbReintent = (opcions = {}) => reintentar(async () => {
  const controlador = new AbortController();
  const limit = setTimeout(() => controlador.abort(), 10_000);
  try {
    return await generarVendesPerSessio({ ...opcions, signal: controlador.signal });
  } catch (error) {
    // Marquem les fallades definitives perque reintentar no insisteixi.
    if (NO_REINTENTABLES.has(error.code ?? error.codi)) error.definitiu = true;
    throw error;
  } finally {
    clearTimeout(limit);
  }
}, { intents: 3, esperaMs: 500, haDeReintentar: (e) => !e.definitiu });

Reintentar un ENOENT és inútil: el fitxer no apareixerà tot sol i només aconsegueixes endarrerir el missatge d'error tres vegades. Un AbortError per límit de temps o un EBUSY transitori, en canvi, sí que es poden resoldre al segon intent. Reintentar sense classificar l'error és de les maneres més cares d'amagar un problema.

Conclusió

Ja no només consumeixes streams: els escrius. Saps que objectMode converteix una canonada de bytes en una canonada d'objectes, que el highWaterMark passa a comptar elements i que readableObjectMode i writableObjectMode es declaren per separat perquè un analitzador rep text i emet objectes. Domines el Transform amb _transform —cridant callback() exactament una vegada i senyalant els errors amb callback(error)— i _flush, on s'emet tot l'acumulat, inclosa aquella última línia sense salt que es perd en silenci quan falta. Has construït un analitzador de CSV amb estat, un Writable agregador i un filtre, i els has encadenat a la canonada completa d'Escena Viva: CSV → objectes → filtre per sala → agregació → JSON, amb les 1811 vendes quadrant amb la llavor i memòria constant en tot el recorregut. I has substituït pipe() per stream.pipeline de node:stream/promises, que propaga errors per tota la cadena, destrueix tots els streams en fallar, s'integra amb try/catch i accepta AbortSignal. La regla és ferma: pipe per explicar, pipeline per treballar.

Has vist a més com zlib.createGzip s'intercala com una etapa més —i com una canonada pot llegir directament d'un .gz sense descomprimir a disc—, com Readable.from() converteix qualsevol iterable en un stream, i per què un generador asíncron és gairebé sempre millor que una classe Transform: menys codi, objectMode implícit, await natural i el final de l'entrada tractat sense cerimònies. I saps cancel·lar una canonada amb AbortSignal sense deixar descriptors vius.

Queda una peça al fons de tot això. Cada vegada que hem escrit { encoding: 'utf8' } estàvem demanant una traducció de bytes a text, i cada vegada que l'hem omesa —en comprimir, en copiar— les dades viatjaven com a Buffer. A Buffers i Dades Binàries mirarem per fi dins d'aquesta caixa: què és un Buffer i per què viu fora del munt de V8, per què allocUnsafe pot mostrar-te memòria d'un altre, les codificacions i els seus usos, l'ordre de bytes, l'error clàssic de slice compartint memòria, per què un emoji trenca un trossejat ingenu i com ho evita StringDecoder. I ho aplicarem a Escena Viva: detectar pels seus nombres màgics si el cartell que puja un organitzador és de debò un PNG, i codificar en base64url el codi QR d'una entrada.

Curs de Node.js: De Principiant a Avançat

Mòdul 1: Introducció a Node.js

Mòdul 2: Conceptes Bàsics

Mòdul 3: Sistema de Fitxers i E/S

Mòdul 4: HTTP i Servidors Web

Mòdul 5: NPM i Gestió de Paquets

Mòdul 6: Framework Express.js

Mòdul 7: Bases de Dades i ORMs

Mòdul 8: Autenticació i Autorització

Mòdul 9: Proves i Depuració

Mòdul 10: Temes Avançats

Mòdul 11: Desplegament i DevOps

Mòdul 12: Projectes del Món Real

© Copyright 2026. Tots els drets reservats