Tema representativo de entrevista

Entrevista Backend: ¿Cómo Diseñar un Pipeline Cancelable con stream.compose en Node.js?

BackendDifícil
Equipo editorial de Offer.ccPublicado Actualizado

Pregunta

Un servicio de Node.js debe procesar subidas en stream mediante descompresión, escaneo, validación y almacenamiento de objetos. ¿Cómo garantiza la contrapresión, la cancelación, la propagación de errores y la limpieza con stream.compose?

Consigna y contexto

Un servicio de Node.js recibe cargas de archivos grandes y debe ejecutar en secuencia la descompresión, el escaneo de virus, la validación de formato y la escritura en el almacenamiento de objetos. La implementación anterior escucha manualmente los eventos data, a veces incrementa el consumo de memoria, sigue ejecutándose después de que un cliente se desconecta y deja archivos temporales cuando falla un paso intermedio.

Utilice stream.compose() o un enfoque equivalente de composición de pipelines. Explique cómo se conectan las etapas legibles (readable), escribibles (writable) y Transform, y cómo la contrapresión, AbortSignal, la propagación de errores y la limpieza final mantienen un trabajo bajo control.

Qué evalúa el entrevistador

  • Si distingue los límites y ciclos de vida de pipe, pipeline y compose.
  • Si puede explicar cómo la contrapresión limita la producción en lugar de simplemente aumentar los límites de las colas.
  • Si la cancelación, las excepciones y las desconexiones de clientes se propagan a través de todo el pipeline.
  • Si maneja generadores asíncronos (async generators), liberación de recursos, escrituras idempotentes y observabilidad.

Preguntas para clarificar

  1. ¿La entrada proviene de una solicitud HTTP, un archivo o un SDK de almacenamiento de objetos? ¿Se admiten salidas multipart y reintentos?
  2. ¿Es cada etapa un stream de Node, un Web Stream, un AsyncIterable o una función ordinaria?
  3. ¿El escaneo y la validación crean un proceso secundario (child process), un archivo temporal o un estado en la base de datos?
  4. ¿Puede el almacenamiento de objetos abortar una subida multipart y debe reanudarse un trabajo cancelado?

Respuesta de 30 segundos

Definiría cada etapa como un readable, writable, Transform o AsyncIterable y las compondría en un Duplex con stream.compose(), dejando luego que pipeline controle el destino final. Los productores continúan únicamente cuando el flujo descendente (downstream) puede aceptar datos, por lo que la contrapresión acota la memoria. Las desconexiones de solicitudes, los límites de tiempo (deadlines) y la cancelación de negocio comparten un único AbortSignal que se pasa a las etapas compatibles; cualquier error hace fallar la cadena. Los archivos temporales, procesos secundarios y subidas multipart se limpian en bloques finally o manejadores de abort, y las escrituras usan claves de idempotencia. Las métricas cubren el rendimiento (throughput), la profundidad de la cola, el pico de RSS, el motivo de cancelación y el resultado de la limpieza.

Análisis detallado paso a paso

Definir el contrato de cada etapa

Especifique los tipos de entrada y salida, el tamaño de chunk, si se permite null, el comportamiento de bloqueo y la propiedad de los recursos al finalizar. Un generador asíncrono debe consumir su origen y emitir bajo demanda; una función ordinaria no debe leer silenciosamente todo el archivo en memoria.

Componer etapas reutilizables

compose conecta streams, AsyncIterables o funciones en un nuevo Duplex y maneja etapas adyacentes con la semántica de pipelines. Ejemplo:

js
import { compose } from 'node:stream';

async function* validate(source) {
  for await (const chunk of source) {
    checkChunk(chunk);
    yield chunk;
  }
}

const processing = compose(decompress(), validate, scan());

La implementación real debe conectar el procesamiento con el writable de destino y centralizar la finalización y los errores en lugar de silenciarlos en cada etapa.

Hacer de la contrapresión el plano de control predeterminado

Cuando el flujo descendente no está listo, las etapas Readable y Transform deben detener la producción. No haga push sin límites en un callback de data ni oculte un consumidor lento aumentando indefinidamente el límite de agua alta (high-water mark). Las pruebas de carga deben registrar cada cola, el throughput y el RSS para que la etapa más lenta determine la velocidad total.

Propagar la cancelación y las desconexiones

Combine la desconexión de la solicitud, el tiempo límite y la cancelación manual en un solo AbortController. Pase su señal a las etapas compatibles de compose y a los SDK externos; tras abortar, detenga la lectura, destruya el flujo descendente y espere el cierre. Los registros deben distinguir la cancelación normal, el fallo de negocio y el error de red.

Centralizar errores y limpieza

Un único propietario debe invocar pipeline/compose, recibir el primer error y destruir la cadena. Los archivos temporales, escáneres, sockets y subidas multipart deben liberarse en caso de éxito, fallo y cancelación; los fallos de limpieza generan alertas con el ID del trabajo sin reemplazar el error original.

Diseñar la idempotencia y las métricas

Genere claves de idempotencia a partir del ID de subida y la versión de la etapa, y luego confirme el estado de negocio solo después de que se complete la escritura del objeto. Registre los bytes de entrada y salida, la duración, el pico de RSS, el motivo de cancelación, la etapa fallida y el tiempo de limpieza. Reintente únicamente desde un límite recuperable para evitar escrituras duplicadas.

Respuesta modelo

Definiría la subida, descompresión, escaneo, validación y almacenamiento como etapas con entradas, salidas y responsabilidades explícitas, las compondría en un pipeline readable, writable o async-iterable, y dejaría que un único propietario del pipeline lo controle. El consumo descendente regula la contrapresión; el almacenamiento en búfer ilimitado en manejadores data está prohibido. La desconexión, el timeout y la cancelación manual comparten un AbortSignal pasado a las etapas compatibles y al SDK de almacenamiento. El propietario maneja el primer error y destruye la cadena; los archivos temporales, escáneres y subidas multipart se limpian tanto en éxito, como en fallo y cancelación. Las claves de idempotencia protegen las escrituras, las métricas cubren throughput, colas, RSS, cancelación y limpieza, y los reintentos inician solo en límites seguros.

Errores comunes

  • Acumular todo el stream en un Buffer y aun así afirmar que se utiliza compose.
  • Escuchar únicamente el error del writable final, pasando por alto fallos en generadores asíncronos o procesos secundarios.
  • Continuar leyendo y escribiendo después de que el cliente se ha desconectado.
  • Resolver la contrapresión aumentando indefinidamente el high-water mark.
  • Eliminar archivos temporales solo en caso de éxito, ignorando abortos y excepciones.
  • Reintentar sin claves de idempotencia y duplicar objetos o estados de negocio.

Preguntas de seguimiento

¿En qué se diferencian compose y pipeline?

Compose crea un Duplex reutilizable a partir de etapas; pipeline controla la conexión de extremo a extremo, propaga errores y espera el cierre. Pueden combinarse, pero la propiedad y el manejo de errores deben estar centralizados.

¿Qué sucede cuando un generador asíncrono lanza una excepción?

El stream compuesto debe fallar y destruir las etapas adyacentes. El llamador aún espera la limpieza del pipeline y registra la etapa fallida en lugar de depender únicamente de la salida del proceso.

¿Cómo se verifica que la contrapresión funciona?

Realice pruebas de carga con un productor más rápido que un destino lento controlado, observe una profundidad de cola y un RSS acotados, y confirme que el throughput sigue a la etapa más lenta en lugar de acumular búfer ilimitado.

¿Se puede reintentar inmediatamente después de una cancelación?

Primero confirme que todos los recursos estén cerrados y que el estado temporal sea identificable. Reintente desde un límite recuperable e idempotente; compense o marque el estado desconocido para aquellas escrituras externas que no se puedan interrumpir.

Fuentes públicas

Preguntas relacionadas