xTaskjs 1.0 ya está disponible: cuentas por rol, documentación y una interfaz personalizable.

Paquetes

@xtaskjs/queues

Handlers de colas agnósticos del transporte, utilidades para brokers, decoradores de publicación y consumidores gestionados por el ciclo de vida.

npm install @xtaskjs/queues reflect-metadata Ruta del paquete: packages/queues

Resumen

Qué controla este paquete en el runtime

Queues incorpora entrega de mensajes con broker y en memoria a xtaskjs. Registra transports con nombre, descubre consumidores decorados desde servicios gestionados por DI, expone APIs de productores inyectables y coordina reintentos, dead-letter routing y el arranque o parada de consumidores mediante el ciclo de vida de la aplicación.

Qué ofrece

  • configureQueues() establece el comportamiento por defecto del transporte y puede desactivar el transporte automático en memoria.
  • registerQueueTransport(), registerInMemoryQueueTransport(), createRabbitMqTransport() y createMqttTransport() conectan brokers o transports locales.
  • QueueHandler(), QueuePattern(), PublishToQueue(), InjectQueueService() e InjectQueueTransport() cubren flujos de consumidores, productores e inyección de transports.
  • QueueService puede publicar mensajes, crear productores, inspeccionar consumidores y grupos, e iniciar o detener el estado del runtime de colas mediante APIs gestionadas por el ciclo de vida.

Cómo encaja

  • Se inicializa automáticamente mediante @xtaskjs/core cuando el paquete está instalado y existe configuración de colas o consumidores decorados.
  • Se combina con el transporte integrado en memoria para desarrollo local y tests, y con las utilidades de RabbitMQ o MQTT para entrega apoyada en brokers.
  • Lo demuestran los ejemplos 16-queues_memory_app y 17-queues_rabbitmq_app.

Mapa de uso

Qué peso tiene este paquete dentro del runtime

Arranque

4/5

Inyección de dependencias

4/5

Mensajería

5/5

Operaciones

5/5

Integraciones

4/5

Flujo del paquete

Cómo atraviesa este paquete las fases del runtime de xtaskjs

Antes del arranque

Configura valores por defecto de colas, registra transports y decora consumidores mientras cargan los módulos para que el runtime de colas sepa qué productores y handlers conectar.

Durante CreateApplication()

El ciclo de vida de queues publica QueueService y tokens de transport con nombre en el contenedor, descubre consumidores decorados y los prepara para arrancar en la fase ready del ciclo de vida.

Durante app.close()

Los consumidores activos se detienen y los transports conectados se desconectan automáticamente para que los recursos del broker se apaguen junto con la aplicación.

Superficie API

Exports representativos del paquete original

Configuración y transports

  • configureQueues
  • registerQueueTransport
  • registerInMemoryQueueTransport
  • createRabbitMqTransport
  • createMqttTransport

Decoradores e inyectores

  • QueueHandler
  • QueueSubscribe
  • QueuePattern
  • PublishToQueue
  • InjectQueueService
  • InjectQueueLifecycleManager
  • InjectQueueTransport

Servicio de runtime y ciclo de vida

  • QueueService
  • initializeQueueIntegration
  • shutdownQueueIntegration
  • getQueueServiceToken
  • getQueueTransportToken

Uso

Flujo típico de adopción

1. Configura transports y valores por defecto

Empieza con el transporte integrado en memoria para flujos locales, o registra las utilidades de RabbitMQ o MQTT cuando los mensajes deban cruzar límites de proceso.

2. Decora consumidores y publicadores

Usa QueueHandler y QueuePattern para consumidores, PublishToQueue para eventos basados en el resultado del método e InjectQueueService cuando los servicios necesiten APIs directas de publicación o productores.

3. Inspecciona y controla el runtime de colas

Usa QueueService para listar grupos y consumidores, crear productores con valores por defecto e iniciar o detener el procesamiento de colas para endpoints de diagnóstico u operaciones.

Ejemplo

Fragmento de referencia

Transport con nombre y consumidores de colas decorados
import { Service } from "@xtaskjs/core";
import {
  InjectQueueService,
  PublishToQueue,
  QueueHandler,
  QueueService,
  configureQueues,
  registerInMemoryQueueTransport,
} from "@xtaskjs/queues";

configureQueues({
  defaultTransportName: "memory",
  autoCreateDefaultInMemoryTransport: false,
});

registerInMemoryQueueTransport({
  name: "memory",
  kind: "in-memory",
});

@Service()
export class OrdersQueueService {
  constructor(
    @InjectQueueService()
    private readonly queues: QueueService
  ) {}

  async publishOrder(orderId: string) {
    await this.queues.publish("orders.created", { orderId });
  }

  @PublishToQueue("orders.completed", { transportName: "memory" })
  completeOrder(orderId: string) {
    return { orderId, completedAt: new Date().toISOString() };
  }

  @QueueHandler("orders.created", { name: "orders.created", transportName: "memory" })
  onOrderCreated(payload: { orderId: string }) {
    console.log("received", payload.orderId);
  }
}

Ejemplos

Ejemplos oficiales para revisar después

Ejemplos de referencia: 16-queues_memory_app and 17-queues_rabbitmq_app

Relacionados

Paquetes que suelen usarse junto a este