Langfuse v4, día 2 (5 de 10): las treinta y nueve colas del worker, y las tres que deciden si pierdes datos

Índice

Quinto artículo de una serie de diez sobre operar Langfuse v4 en producción. El anterior cubría la migración desde la versión 3. Este entra en el worker que ya está corriendo. Verificado contra el código de Langfuse 4.35.0, commit 2ee1908 del 14 de septiembre de 2026.

TL;DR

Son treinta y nueve colas y treinta y dos interruptores. El catálogo está en packages/shared/src/server/queues.ts:406. Las variables QUEUE_CONSUMER_*_IS_ENABLED se leen una sola vez al arrancar, así que apagar una cola es un reinicio, no un cambio en caliente.

El particionado no reparte nada sin Redis Cluster. El productor solo aplica el hash si REDIS_CLUSTER_ENABLED === "true", y esa variable viene en "false". Subir LANGFUSE_INGESTION_QUEUE_SHARD_COUNT en una instalación con Redis de un nodo arranca N workers y deja a N−1 sin trabajo.

El reintentador de fallidos viene apagado. QUEUE_CONSUMER_DEAD_LETTER_RETRY_QUEUE_IS_ENABLED es la única de las treinta y dos con valor por defecto "false".

Y aunque lo enciendas, no cubre la ingesta. Su lista tiene cinco colas y ninguna es ingestion-queue. Un evento que agote sus seis intentos se queda en el conjunto de fallidos.

El escritor de ClickHouse descarta filas. Tres intentos con retardo plano de cien milisegundos y, si no, fuera. El comentario del código es explícito: // TODO - Add to a dead letter queue in Redis rather than dropping. La métrica que hay que vigilar se llama langfuse.queue.clickhouse_writer.rows_dropped.

Si apagas el consumidor de una cola, dejas de ver su profundidad. El emisor de métricas solo consulta colas con worker registrado en ese contenedor. La cola sigue creciendo y el panel se queda plano.

El borrado de un proyecto va a un trabajo cada diez minutos. Por limitador, a propósito, porque cada uno arrasa ClickHouse.

Hay dos sondas de salud opcionales y una de ellas puede provocar el reinicio en bucle que pretende arreglar si se configura con menos de sesenta segundos de retardo inicial.

Estás aquí: OBSERVE, día 2

El artículo anterior dejó la instalación migrada. Este es el mapa del proceso que hace el trabajo, y va antes que el de capacidad porque las colas son donde se ve la saturación antes de que ClickHouse se entere.

La analogía: la centralita de clasificación

Una plataforma de clasificación postal tiene cintas distintas para paquetes, cartas certificadas, devoluciones y valija interna. Cada cinta tiene su velocidad y su personal, y algunas están limitadas a propósito porque el destino no admite más ritmo.

El trabajo del que la opera no es mirar la cinta más rápida, es saber cuál se atasca primero, cuál puede pararse una tarde sin que nadie lo note y cuál no puede pararse ni diez minutos. Y sobre todo, dónde está el contenedor de los envíos que nadie pudo clasificar, porque en toda plataforma hay uno y en muchas no lo mira nadie.

El worker de Langfuse tiene treinta y nueve cintas. Este artículo es el plano, con las velocidades y el sitio donde está el contenedor.

El mapa de las treinta y nueve

Agrupadas por lo que hacen, con lo que deja de pasar si el grupo se para:

Ingesta (4). ingestion-queue, secondary-ingestion-queue, otel-ingestion-queue, secondary-otel-ingestion-queue. Si se paran, la web sigue aceptando eventos y subiéndolos al almacenamiento de objetos, y sigue encolando. Nada llega a ClickHouse. El SDK no ve error ninguno. Es el grupo que no puede pararse.

Propagación (1). event-propagation-queue, la del dual-write del artículo anterior, con concurrencia global de uno en todo el clúster.

Evaluación (8). trace-upsert, dataset-run-item-upsert-queue, create-eval-queue, evaluation-execution-queue y su secundaria, llm-as-a-judge-execution-queue, code-eval-execution-queue, experiment-create-queue. Si se paran, no se crean ni se ejecutan evaluaciones. Las trazas siguen entrando. Es el grupo que sí puede esperar.

Borrados y retención (7). trace-delete, score-delete, dataset-delete-queue, project-delete, data-retention-queue y su cola de proceso, batch-action-queue. Si se paran, ni los borrados pedidos por la interfaz ni la retención por proyecto se ejecutan. Esto no es una molestia operativa, es un incumplimiento si tienes un compromiso de retención firmado.

Exportaciones (3). batch-export-queue, core-data-s3-export-queue, metering-data-postgres-export-queue. Las descargas que pide un usuario se quedan pendientes.

Integraciones salientes (8). Las parejas de PostHog, Mixpanel y almacenamiento de objetos, más webhook-queue, entity-change-queue y notification-queue. En cada pareja, la cola a secas es el planificador que reparte por proyecto y la de proceso es la que trabaja.

Facturación de Cloud (4). Solo se registran con clave de Stripe. En una instalación propia no arrancan.

Mantenimiento (4). dead-letter-retry-queue, monitor-queue, in-app-agent-run-queue y v4-legacy-api-usage-queue, que no ingiere nada: es un cron que escanea system.query_log de ClickHouse para detectar qué proyectos siguen usando la API vieja.

Los interruptores son treinta y dos, todos QUEUE_CONSUMER_*_IS_ENABLED y todos en "true" salvo el del reintentador de fallidos. Cuatro colas no tienen interruptor propio: la de juez-modelo va dentro del de evaluación, las dos de exportación de datos de Cloud tienen su propia variable, y las colas de proceso comparten interruptor con su planificador.

El mecanismo es un if de nivel superior alrededor de cada registro de worker. Interruptor en falso quiere decir que no se instancia el worker, no que la cola no reciba trabajo: los productores siguen encolando igual. Apagar una cola es dejar de consumirla, no dejar de llenarla.

El particionado que no particiona

Nueve colas admiten particionado: las cuatro de ingesta, las tres de evaluación, la de juez-modelo y trace-upsert. Cada una tiene su variable de número de particiones, y todas valen 1 por defecto.

Subirlas parece el primer paso obvio para escalar la ingesta. No lo es, por esto (packages/shared/src/server/redis/ingestionQueue.ts:48):

    const shardIndex =
      IngestionQueue.getShardIndexFromShardName(shardName) ??
      (env.REDIS_CLUSTER_ENABLED === "true" && shardingKey
        ? getShardIndex(shardingKey, env.LANGFUSE_INGESTION_QUEUE_SHARD_COUNT)
        : 0);

El : 0 del final es todo el asunto. El productor solo calcula la partición si Redis está en modo clúster, y REDIS_CLUSTER_ENABLED viene en "false". En una instalación con Redis de un nodo, que es la mayoría de las instalaciones propias, subir el número de particiones a cuatro arranca cuatro workers, tres de los cuales se quedan mirando colas vacías, mientras todo el tráfico sigue entrando por ingestion-queue.

La clave de reparto, cuando sí se aplica, es projectId-eventBodyId pasada por SHA-256, de la que se toman los ocho primeros caracteres hexadecimales y se saca el módulo. Los nombres de partición son ingestion-queue para la cero y ingestion-queue-N para el resto.

Sin Redis Cluster, la palanca que sí funciona es la concurrencia, que es otra cosa.

Concurrencias, y las que están puestas a uno a propósito

Los valores por defecto, que hay que conocer porque varios están a uno por buenas razones:

ColaVariablePor defecto
ingestion-queue (por partición)LANGFUSE_INGESTION_QUEUE_PROCESSING_CONCURRENCY20
secondary-ingestion-queueLANGFUSE_INGESTION_SECONDARY_QUEUE_PROCESSING_CONCURRENCY5
otel-ingestion-queueLANGFUSE_OTEL_INGESTION_QUEUE_PROCESSING_CONCURRENCY5
trace-upsertLANGFUSE_TRACE_UPSERT_WORKER_CONCURRENCY25
monitor-queueLANGFUSE_MONITOR_QUEUE_PROCESSING_CONCURRENCY10
evaluación y juez-modelovarias5
create-eval-queueLANGFUSE_EVAL_CREATOR_WORKER_CONCURRENCY2
trace-delete, score-delete, dataset-delete-queue, project-deletevarias1
event-propagation-queuefijo en el código1 global de clúster

Los borrados están a uno y además llevan limitador de tasa. El de proyectos es el más agresivo (worker/src/app.ts:234):

    limiter: {
      // Process at most `max` delete jobs per LANGFUSE_CLICKHOUSE_PROJECT_DELETION_CONCURRENCY_DURATION_MS (default 10 min)
      max: env.LANGFUSE_PROJECT_DELETE_CONCURRENCY,

Un borrado de proyecto cada diez minutos. Si un cliente te pide borrar cuarenta proyectos, eso son seis horas y cuarenta minutos de reloj, y está bien que así sea porque cada uno de esos borrados es una escoba pasando por ClickHouse.

Las colas de ingesta, en cambio, no tienen limitador ninguno. Ahí el freno es la concurrencia y el escritor de ClickHouse.

Un detalle de BullMQ que aparece repetido en varias colas y que explica bastantes atascos ajenos (worker/src/app.ts:261):

        // The default lockDuration is 30s and the lockRenewTime 1/2 of that.
        // We set it to 60s to reduce the number of lock renewals and also be less sensitive to high CPU wait times.
        lockDuration: 60000, // 60 seconds
        stalledInterval: 120000, // 120 seconds
        maxStalledCount: 3,

El camino de un evento, de la API a ClickHouse

Merece seguirlo entero una vez, porque hay tres sitios donde se puede perder.

La web agrupa los eventos del lote por identificador de cuerpo, sube un fichero JSON por grupo al almacenamiento de objetos y aborta si esa subida falla: si no hay blob, no hay encolado. Después encola un trabajo por identificador de cuerpo en ingestion-queue.

El worker coge el trabajo, mira opcionalmente una caché de procesados recientes en Redis (apagada por defecto, cinco minutos de vida), decide si redirige el proyecto a la cola secundaria, y baja el fichero del almacenamiento. Aquí hay un comentario que vale por un gráfico de capacidad:

Process files in batches. If a user has 5k events, this will likely take 100 seconds.

Luego mezcla y escribe. Y escribir no es escribir: el servicio no habla con ClickHouse, mete la fila en un buffer con 1.000 filas, 1.000 milisegundos y 3 intentos de valores por defecto.

Ese buffer es el tercer sitio donde se pierde, y el más serio del sistema (worker/src/services/ClickhouseWriter/index.ts:535):

      // Re-add the records to the queue with incremented attempts
      let droppedCount = 0;
      queueItems.forEach((item) => {
        if (item.attempts < this.maxAttempts) {
          entityQueue.push({ ...item, attempts: item.attempts + 1 });
        } else {
          // TODO - Add to a dead letter queue in Redis rather than dropping
          recordIncrement("langfuse.queue.clickhouse_writer.error");
          droppedCount++;
        }
      });

Agotados los tres intentos, la fila se descarta. Y no hay reintento de cola que la salve, porque el trabajo de BullMQ ya terminó correctamente: para la cola, ese evento se procesó. Los reintentos internos, además, no son exponenciales de verdad: el retardo es plano de cien milisegundos, de modo que los tres intentos ocurren en trescientos milisegundos y un ClickHouse indispuesto medio segundo se lleva por delante el lote.

La consecuencia para el runbook es directa. langfuse.queue.clickhouse_writer.rows_dropped tiene que estar en un panel y tiene que tener alerta a cero. Cualquier valor distinto de cero es pérdida de datos silenciosa.

Reintentos, y el contenedor que nadie mira

Cada cola define su propia política. Las que importan:

ColaIntentosRetroceso
ingestion-queue6exponencial, 5 s
secondary-ingestion-queue5exponencial, 5 s
evaluation-execution-queue10exponencial, 1 s
code-eval-execution-queue3exponencial, 30 s
trace-upsert2exponencial, 5 s, con 30 s de retardo inicial
trace-delete, score-delete2exponencial, 30 s
monitor-queue, in-app-agent-run-queue1sin retroceso, a propósito

Los dos de un solo intento lo son por decisión de diseño, y el comentario lo explica bien: “Postgres owns correctness (claim CAS + reconcile-on-read); BullMQ is delivery-only, so never redeliver on its own.”

Y ahora el contenedor. No hay una cola de fallidos dedicada. Lo que hay es el conjunto failed nativo de BullMQ de cada cola, con retención propia, más una cola que reintenta desde ahí. Esa cola tiene dos problemas para el que opera:

Primero, viene desactivada. Es la única de las treinta y dos con valor por defecto "false".

Segundo, aunque la enciendas, su alcance está escrito a mano y tiene cinco entradas (worker/src/services/dlq/dlqRetryService.ts:9):

  private static retryQueues = [
    QueueName.ProjectDelete,
    QueueName.TraceDelete,
    QueueName.ScoreDelete,
    QueueName.BatchActionQueue,
    QueueName.DataRetentionProcessingQueue,
  ] as const;

ingestion-queue no está. Los eventos de ingesta que agoten sus seis intentos se quedan en el conjunto de fallidos, hasta cien mil, y nadie los reintenta nunca. Para sacarlos hay una ruta de administración, POST /api/admin/bullmq con acción retry, que el propio fichero describe como la que usa el servicio gestionado para esto mismo.

Las métricas, con sus nombres exactos

El nombre de métrica de cada cola se construye en tiempo de ejecución a partir del nombre de la cola, cambiando guiones por barras bajas y quitando el sufijo _queue. Así que ingestion-queue produce langfuse.queue.ingestion, y secondary-ingestion-queue produce langfuse.queue.secondary_ingestion. En las colas particionadas se añade siempre una etiqueta shard.

Lo que sale de ahí:

  • langfuse.queue.<cola>.rate con type en request, completed, failed, error y stalled.
  • langfuse.queue.<cola>.time_distribution con type en wait y processing.
  • langfuse.queue.<cola>.depth con type en waiting, failed y active, más una serie agregada con shard: "all".
  • langfuse.queue.<cola>.dlq_oldest_age, la edad del fallido más viejo, que emite cero cuando el conjunto se vacía para que el panel se resetee.
  • Del escritor: langfuse.queue.clickhouse_writer.rows_dropped, .error, .processing_time y ingestion_clickhouse_insert_queue_length.

Dos cosas que hay que saber de esas métricas. La profundidad waiting suma los pausados, por coherencia con el contador de BullMQ. Y la trampa gorda (worker/src/features/queue-metrics-runner/index.ts:70):

    // Only poll queues that have registered workers. This avoids calling
    // getInstance() on queues this worker doesn't consume, which would
    // create unnecessary Redis connections and can trigger side effects

Solo se emiten métricas de las colas que ese contenedor consume. Si separas roles y apagas el consumidor de una cola en todos los contenedores, esa cola deja de aparecer en los paneles mientras sigue creciendo en Redis. Es el modo perfecto de no enterarte.

La métrica que satura primero sigue siendo la misma: langfuse.queue.ingestion.depth con type: "waiting". Y el intervalo de emisión es de un segundo por defecto.

Las dos sondas, y la que puede reiniciarte en bucle

El worker expone /api/health y /api/ready, idénticas salvo en una cosa: la segunda devuelve 500 tras recibir SIGTERM, que es como se drena. Las dos comprueban siempre Postgres con un SELECT 1 y Redis con un ping con dos segundos de tope, fijos en el código.

Después hay dos comprobaciones opcionales que se activan por parámetro:

?failIfEventPropagationStuck=true mira el latido del propagador y devuelve 503 si lleva sin refrescarse más de 35 minutos. Un latido ausente no cuenta como atascado, a propósito, para no reiniciar en bucle un contenedor recién arrancado. Y la condición que hay que respetar sí o sí:

Probes using this flag MUST set initialDelaySeconds >= 60s (one cron cycle) […] A shorter delay can crash-loop the very restart this check triggers.

?failIfQueueConsumptionStuck=true es una señal en memoria del contenedor, sin Redis de por medio: se marca en los eventos active y completed de cualquier worker. Devuelve 503 si ese contenedor lleva 60 minutos sin coger ni terminar un solo trabajo. Existe por un modo de fallo muy concreto:

After Redis lock loss BullMQ workers can wedge permanently: the process stays alive and connectivity checks pass, but no queue picks up jobs ever again.

El umbral de una hora se apoya en que los crones por defecto mantienen ocupado a un worker sano al menos una vez por hora. En despliegues con varias réplicas, cada tic de planificador cae en una sola réplica, así que el propio comentario avisa de subir el umbral si una réplica puede estar legítimamente parada.

El drenado, por cierto, no tiene tope. El orden es cerrar el servidor HTTP, parar los runners periódicos, cerrar los workers de BullMQ y después vaciar el escritor de ClickHouse, que en el apagado vuelca la cola entera y no solo un lote. Lo único que acota ese drenado es el periodo de gracia de Kubernetes, así que más vale que sea generoso.

Los crones

ColaPatrónQué hace
event-propagation-queue* * * * *propaga una partición de la tabla de paso
dead-letter-retry-queue0 */10 * * * *reintenta fallidos de cinco colas
v4-legacy-api-usage-queue*/15 * * * *escanea system.query_log buscando uso de la API vieja
blobstorage-integration-queue*/20 * * * *reparte exportaciones por proyecto
cloud-usage-metering-queue5 * * * *facturación de Cloud
posthog-integration-queue, mixpanel-integration-queue30 * * * *reparto horario
metering-data-postgres-export-queue30 2 * * *exportación diaria
core-data-s3-export-queue, data-retention-queue15 3 * * *exportación y retención diarias

Y un aviso que el propio fichero del planificador escribe en mayúsculas: el registro de los crones es fire-and-forget desde los constructores de las colas, y si falla solo se registra en el log. Un fallo transitorio de Redis al arrancar puede dejar un cron sin programar hasta el siguiente reinicio, sin que el contenedor se marque como enfermo.

Separar roles, que es para lo que sirve todo esto

Con treinta y dos interruptores, lo que se puede montar es un worker por familia. El reparto que tiene sentido en una instalación con carga agéntica:

# Pool 1: ingesta. Es el que escala con el trafico.
QUEUE_CONSUMER_INGESTION_QUEUE_IS_ENABLED=true
QUEUE_CONSUMER_OTEL_INGESTION_QUEUE_IS_ENABLED=true
QUEUE_CONSUMER_EVENT_PROPAGATION_QUEUE_IS_ENABLED=true    # ojo: concurrencia global 1
LANGFUSE_INGESTION_QUEUE_PROCESSING_CONCURRENCY=20
# el resto a false

# Pool 2: evaluaciones. Escala con el numero de evaluadores, no con el trafico.
QUEUE_CONSUMER_EVAL_EXECUTION_QUEUE_IS_ENABLED=true
QUEUE_CONSUMER_CREATE_EVAL_QUEUE_IS_ENABLED=true
QUEUE_CONSUMER_TRACE_UPSERT_QUEUE_IS_ENABLED=true
QUEUE_CONSUMER_EXPERIMENT_CREATE_QUEUE_IS_ENABLED=true

# Pool 3: mantenimiento. Una replica basta.
QUEUE_CONSUMER_TRACE_DELETE_QUEUE_IS_ENABLED=true
QUEUE_CONSUMER_PROJECT_DELETE_QUEUE_IS_ENABLED=true
QUEUE_CONSUMER_DATA_RETENTION_QUEUE_IS_ENABLED=true
QUEUE_CONSUMER_BATCH_EXPORT_QUEUE_IS_ENABLED=true
QUEUE_CONSUMER_DEAD_LETTER_RETRY_QUEUE_IS_ENABLED=true    # encenderla

Dos avisos sobre este reparto. Cada registro de worker abre su propia conexión a Redis, así que el número de conexiones crece con colas por particiones por réplicas, y eso se mira en el maxclients de Redis antes de multiplicar réplicas. Y en cuanto separas roles pierdes las métricas de las colas que ese contenedor no consume, de modo que necesitas que alguna réplica consuma cada cola, o te quedas ciego en esa.

Checklist

  • Hay alerta sobre langfuse.queue.clickhouse_writer.rows_dropped, con umbral cero.
  • Hay alerta sobre la profundidad de langfuse.queue.ingestion con type: "waiting".
  • Hay alerta sobre dlq_oldest_age de las colas que importan, y alguien sabe usar POST /api/admin/bullmq.
  • El reintentador de fallidos está encendido en alguna réplica, sabiendo que no cubre la ingesta.
  • No se ha subido el número de particiones sin tener Redis Cluster.
  • Las sondas de salud opcionales están puestas con initialDelaySeconds de al menos 60 segundos.
  • El periodo de gracia de terminación de Kubernetes da tiempo a vaciar el escritor de ClickHouse.
  • Cada cola la consume al menos una réplica, para no perder sus métricas.
  • El maxclients de Redis aguanta colas por particiones por réplicas.

Trampas

El particionado sin Redis Cluster no reparte. Arranca los workers y los deja sin trabajo.

El reintentador de fallidos viene apagado y no cubre la ingesta. Dos hechos independientes y los dos incómodos.

El escritor de ClickHouse tira filas tras tres intentos. Con un TODO en el código reconociéndolo y trescientos milisegundos de margen total.

Apagar un consumidor apaga también sus métricas. La cola sigue creciendo, el panel se queda plano.

Un cron puede quedarse sin programar y el contenedor sigue sano. Fallo transitorio de Redis al arrancar, y hasta el siguiente reinicio.

Las integraciones de PostHog y de almacenamiento se desactivan solas cuando un error se clasifica como fallo de configuración del cliente, y si además falla la notificación se quedan apagadas y calladas.

La sonda de propagación mal configurada reinicia en bucle. Menos de sesenta segundos de retardo inicial y te comes el ciclo.

El drenado no tiene tope propio. Lo acota el periodo de gracia de Kubernetes y nada más.

La serie: los diez artículos

  1. Qué entra en una traza: modelo de datos de la versión 4, límites, precedencias, scores, enmascarado e índices.
  2. Poner LangGraph delante: la instrumentación de una plataforma agéntica y el coste medido de un turno.
  3. Dos agentes conocidos instrumentados: Open Deep Research y GPT Researcher, con las cifras medidas de una petición real.
  4. Migrar de la versión 3 a la 4 sin ventana: los tres pasos del modo de escritura, las migraciones de fondo reanudables y dónde está el punto de no retorno del retroceso.
  5. Las colas del worker (este artículo): el mapa de las treinta y nueve, qué pool dedicar a cada grupo, los interruptores por cola, el particionado y la concurrencia.
  6. Capacidad y coste real de ClickHouse: cómo medir los bytes por observación con las tablas del sistema, la diferencia entre la tabla completa y la de listados, y el coste de fusión de los índices de texto completo.
  7. Retención, borrado y protección de datos: por qué un borrado no libera disco, el limpiador de máscaras que viene desactivado, la cola de borrados pendientes y el ciclo de vida de S3 que hay que implementar a mano.
  8. Copias de seguridad y recuperación cruzada: orden de restauración entre Postgres, ClickHouse y el almacenamiento de objetos, qué rompe cada desajuste, y hasta dónde llega la reproducción de eventos.
  9. Runbook de saturación: qué alertar de las métricas de cola, las sondas de atasco, el drenado por el endpoint de preparación y la cola de mensajes fallidos.
  10. Sacar los datos fuera: la integración de almacenamiento de objetos a Parquet, las exportaciones por lotes y la API de métricas, para montar el lago de datos.

Ver también

Fuentes

  • Código de Langfuse 4.35.0, commit 2ee1908 del 14 de septiembre de 2026: packages/shared/src/server/queues.ts, packages/shared/src/server/redis/, worker/src/env.ts, worker/src/app.ts, worker/src/queues/, worker/src/services/ClickhouseWriter/, worker/src/services/dlq/, worker/src/features/health/, worker/src/features/queue-metrics-runner/, worker/src/utils/shutdown.ts.
  • Documentación de despliegue de Langfuse, consultada el 14 de septiembre de 2026.
  • BullMQ: going to production, consultado el 14 de septiembre de 2026.