Files
Eric Allam a999d9ea3f feat(engine): Batch trigger reloaded (#2779)
New batch trigger system with larger payloads, streaming ingestion,
larger batch sizes, and a fair processing system.

This PR introduces a new `FairQueue` abstraction inspired by our own
`RunQueue` that enables multi-tenant fair queueing with concurrency
limits. The new `BatchQueue` is built on top of the `FairQueue`, and
handles processing Batch triggers in a fair manner with per-environment
concurrency limits defined per-org. Additionally, there is a global
concurrency limit to prevent the BatchQueue system from creating too
many runs too quickly, which can cause downstream issues.

For this new BatchQueue system we have a completely new batch trigger
creation and ingestion system. Previously this was a single endpoint
with a single JSON body that defined details about the batch as well as
all the items in the batch.

We're introducing a two-phase batch trigger ingestion system. In the
first phase, the BatchTaskRun record is created (and possibly rate
limited). The second phase is another endpoint that accepts an NDJSON
body with each line being a single item/run with payload and options.

At ingestion time all items are added to a queue, in order, and then
processed by the BatchQueue system.

## New batch trigger rate limits

This PR implements a new batch trigger specific rate limit, configured
on the `Organization.batchRateLimitConfig` column, and defaults using
these environment variables:

- `BATCH_RATE_LIMIT_REFILL_RATE` defaults to 10
- `BATCH_RATE_LIMIT_REFILL_INTERVAL` the duration interval, defaults to
`"10s"`
- `BATCH_RATE_LIMIT_MAX` defaults to 1200

This rate limiter is scoped to the environment ID and controls how many
runs can be submitted via batch triggers per interval. The SDK handles
the retrying side.

## Batch queue concurrency limits

The new column `Organization.batchQueueConcurrencyConfig` now defines an
org specific `processingConcurrency` value, with a backup of the env var
`BATCH_CONCURRENCY_LIMIT_DEFAULT` which defaults to 10. This controls
how many batch queue items are processed concurrently per environment.

There is also a global rate limit for the batch queue set via the
`BATCH_QUEUE_GLOBAL_RATE_LIMIT` which defaults to being disabled. If
set, the entire batch queue system won't process more than
`BATCH_QUEUE_GLOBAL_RATE_LIMIT` items per second. This allows
controlling the maximum number of runs created per second via batch
triggers.

## Batch trigger settings

- `STREAMING_BATCH_MAX_ITEMS` controls the maximum number of items in a
single batch
- `STREAMING_BATCH_ITEM_MAXIMUM_SIZE` controls the maximum size of each
item in a batch
- `BATCH_CONCURRENCY_DEFAULT_CONCURRENCY` controls the default
environment concurrency
- `BATCH_QUEUE_DRR_QUANTUM` how many credits each environment gets each
round for the DRR scheduler
- `BATCH_QUEUE_MAX_DEFICIT` the maximum deficit for the DRR scheduler
- `BATCH_QUEUE_CONSUMER_COUNT` how many queue consumers to run
- `BATCH_QUEUE_CONSUMER_INTERVAL_MS` how frequently they poll for items
in the queue

### Configuration Recommendations by Use Case

**High-throughput priority (fairness acceptable at 0.98+):**

```env
BATCH_QUEUE_DRR_QUANTUM=25
BATCH_QUEUE_MAX_DEFICIT=100
BATCH_QUEUE_CONSUMER_COUNT=10
BATCH_QUEUE_CONSUMER_INTERVAL_MS=50
BATCH_CONCURRENCY_DEFAULT_CONCURRENCY=25
```

**Strict fairness priority (throughput can be lower):**

```env
BATCH_QUEUE_DRR_QUANTUM=5
BATCH_QUEUE_MAX_DEFICIT=25
BATCH_QUEUE_CONSUMER_COUNT=3
BATCH_QUEUE_CONSUMER_INTERVAL_MS=100
BATCH_CONCURRENCY_DEFAULT_CONCURRENCY=5
```
2025-12-16 14:32:49 +00:00

253 lines
8.1 KiB
Markdown
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# Batch Queue & Fair Queue Metrics Guide
This document provides a comprehensive breakdown of all metrics emitted by the Batch Queue and Fair Queue systems, including what they mean and how to identify degraded system states.
## Overview
The batch queue system consists of two layers:
1. **BatchQueue** (`batch_queue.*`) - High-level batch processing metrics
2. **FairQueue** (`batch-queue.*`) - Low-level message queue metrics (with `name: "batch-queue"`)
Both layers emit metrics that together provide full observability into batch processing.
---
## BatchQueue Metrics
These metrics track batch-level operations.
### Counters
| Metric | Description | Labels |
|--------|-------------|--------|
| `batch_queue.batches_enqueued` | Number of batches initialized for processing | `envId`, `itemCount`, `streaming` |
| `batch_queue.items_enqueued` | Number of individual batch items enqueued | `envId` |
| `batch_queue.items_processed` | Number of batch items successfully processed (turned into runs) | `envId` |
| `batch_queue.items_failed` | Number of batch items that failed processing | `envId`, `errorCode` |
| `batch_queue.batches_completed` | Number of batches that completed (all items processed) | `envId`, `hasFailures` |
### Histograms
| Metric | Description | Unit | Labels |
|--------|-------------|------|--------|
| `batch_queue.batch_processing_duration` | Time from batch creation to completion | ms | `envId`, `itemCount` |
| `batch_queue.item_queue_time` | Time from item enqueue to processing start | ms | `envId` |
---
## FairQueue Metrics (batch-queue namespace)
These metrics track the underlying message queue operations. With the batch queue configuration, they are prefixed with `batch-queue.`.
### Counters
| Metric | Description |
|--------|-------------|
| `batch-queue.messages.enqueued` | Number of messages (batch items) added to the queue |
| `batch-queue.messages.completed` | Number of messages successfully processed |
| `batch-queue.messages.failed` | Number of messages that failed processing |
| `batch-queue.messages.retried` | Number of message retry attempts |
| `batch-queue.messages.dlq` | Number of messages sent to dead letter queue |
### Histograms
| Metric | Description | Unit |
|--------|-------------|------|
| `batch-queue.message.processing_time` | Time to process a single message | ms |
| `batch-queue.message.queue_time` | Time a message spent waiting in queue | ms |
### Observable Gauges
| Metric | Description | Labels |
|--------|-------------|--------|
| `batch-queue.queue.length` | Current number of messages in a queue | `fairqueue.queue_id` |
| `batch-queue.master_queue.length` | Number of active queues in the master queue shard | `fairqueue.shard_id` |
| `batch-queue.inflight.count` | Number of messages currently being processed | `fairqueue.shard_id` |
| `batch-queue.dlq.length` | Number of messages in the dead letter queue | `fairqueue.tenant_id` |
---
## Key Relationships
Understanding how metrics relate helps diagnose issues:
```
batches_enqueued × avg_items_per_batch ≈ items_enqueued
items_enqueued = items_processed + items_failed + items_pending
batches_completed ≤ batches_enqueued (lag indicates processing backlog)
```
---
## Degraded State Indicators
### 🔴 Critical Issues
#### 1. Processing Stopped
**Symptoms:**
- `batch_queue.items_processed` rate drops to 0
- `batch-queue.inflight.count` is 0
- `batch-queue.master_queue.length` is growing
**Likely Causes:**
- Consumer loops crashed
- Redis connection issues
- All consumers blocked by concurrency limits
**Actions:**
- Check webapp logs for "BatchQueue consumers started" message
- Verify Redis connectivity
- Check for "Unknown concurrency group" errors
#### 2. Items Stuck in Queue
**Symptoms:**
- `batch_queue.item_queue_time` p99 > 60 seconds
- `batch-queue.queue.length` growing continuously
- `batch-queue.inflight.count` at max capacity
**Likely Causes:**
- Processing is slower than ingestion
- Concurrency limits too restrictive
- Global rate limiter bottleneck
**Actions:**
- Increase `BATCH_QUEUE_CONSUMER_COUNT`
- Review concurrency limits per environment
- Check `BATCH_QUEUE_GLOBAL_RATE_LIMIT` setting
#### 3. High Failure Rate
**Symptoms:**
- `batch_queue.items_failed` rate > 5% of `items_processed`
- `batch-queue.messages.dlq` increasing
**Likely Causes:**
- TriggerTaskService errors
- Invalid task identifiers
- Downstream service issues
**Actions:**
- Check `errorCode` label distribution on `items_failed`
- Review batch error records in database
- Check TriggerTaskService logs
### 🟡 Warning Signs
#### 4. Growing Backlog
**Symptoms:**
- `batch_queue.batches_enqueued` - `batch_queue.batches_completed` is increasing over time
- `batch-queue.master_queue.length` trending upward
**Likely Causes:**
- Sustained high load
- Processing capacity insufficient
- Specific tenants monopolizing resources
**Actions:**
- Monitor DRR deficit distribution across tenants
- Consider scaling consumers
- Review per-tenant concurrency settings
#### 5. Uneven Tenant Processing
**Symptoms:**
- Some `envId` labels show much higher `item_queue_time` than others
- DRR logs show "tenants blocked by concurrency" frequently
**Likely Causes:**
- Concurrency limits too low for high-volume tenants
- DRR quantum/maxDeficit misconfigured
**Actions:**
- Review `BATCH_CONCURRENCY_*` environment settings
- Adjust DRR parameters if needed
#### 6. Rate Limit Impact
**Symptoms:**
- `batch_queue.item_queue_time` has periodic spikes
- Logs show "Global rate limit reached, waiting"
**Likely Causes:**
- `BATCH_QUEUE_GLOBAL_RATE_LIMIT` is set too low
**Actions:**
- Increase global rate limit if system can handle more throughput
- Or accept as intentional throttling
---
## Recommended Dashboards
### Processing Health
```
# Throughput
rate(batch_queue_items_processed_total[5m])
rate(batch_queue_items_failed_total[5m])
# Success Rate
rate(batch_queue_items_processed_total[5m]) /
(rate(batch_queue_items_processed_total[5m]) + rate(batch_queue_items_failed_total[5m]))
# Batch Completion Rate
rate(batch_queue_batches_completed_total[5m]) / rate(batch_queue_batches_enqueued_total[5m])
```
### Latency
```
# Item Queue Time (p50, p95, p99)
histogram_quantile(0.50, rate(batch_queue_item_queue_time_bucket[5m]))
histogram_quantile(0.95, rate(batch_queue_item_queue_time_bucket[5m]))
histogram_quantile(0.99, rate(batch_queue_item_queue_time_bucket[5m]))
# Batch Processing Duration
histogram_quantile(0.95, rate(batch_queue_batch_processing_duration_bucket[5m]))
```
### Queue Depth
```
# Current backlog
batch_queue_master_queue_length
batch_queue_inflight_count
# DLQ (should be 0)
batch_queue_dlq_length
```
---
## Alert Thresholds (Suggested)
| Condition | Severity | Threshold |
|-----------|----------|-----------|
| Processing stopped | Critical | `items_processed` rate = 0 for 5min |
| High failure rate | Warning | `items_failed` / `items_processed` > 0.05 |
| Queue time p99 | Warning | > 30 seconds |
| Queue time p99 | Critical | > 120 seconds |
| DLQ length | Warning | > 0 |
| Batch completion lag | Warning | `batches_enqueued - batches_completed` > 100 |
---
## Environment Variables Affecting Metrics
| Variable | Impact |
|----------|--------|
| `BATCH_QUEUE_CONSUMER_COUNT` | More consumers = higher throughput, lower queue time |
| `BATCH_QUEUE_CONSUMER_INTERVAL_MS` | Lower = more frequent polling, higher throughput |
| `BATCH_QUEUE_GLOBAL_RATE_LIMIT` | Caps max items/sec, increases queue time if too low |
| `BATCH_CONCURRENCY_FREE/PAID/ENTERPRISE` | Per-tenant concurrency limits |
| `BATCH_QUEUE_DRR_QUANTUM` | Credits per tenant per round (fairness tuning) |
| `BATCH_QUEUE_MAX_DEFICIT` | Max accumulated credits (prevents starvation) |
---
## Debugging Checklist
When investigating batch queue issues:
1. **Check consumer status**: Look for "BatchQueue consumers started" in logs
2. **Check Redis**: Verify connection and inspect keys with prefix `engine:batch-queue:`
3. **Check concurrency**: Look for "tenants blocked by concurrency" debug logs
4. **Check rate limits**: Look for "Global rate limit reached" debug logs
5. **Check DRR state**: Query `batch:drr:deficit` hash in Redis
6. **Check batch status**: Query `BatchTaskRun` table for stuck `PROCESSING` batches