adaptive-batch-size-tuning-under-load · diff

v1.1.0 to v2.0.0

140 added, 62 removed. Audit A to A.

---
name: adaptive-batch-size-tuning-under-load
description: Use when writing market data or order logs to downstream databases
(TimescaleDB, ClickHouse) or message brokers to dynamically adapt write batch
- sizes and flush timeouts based on queue pressure and sink write latency.
- Hysteresis-shaped, EWMA-smoothed, back-pressure-aware. Single-writer / multi-producer
- thread-safe engine, not a generic embedder.
+ sizes and flush timeouts based on batch fullness at the flush boundary and
+ sink write latency. Hysteresis-shaped, EWMA-smoothed, back-pressure-aware.
+ Single-writer / multi-producer thread-safe engine, not a generic embedder.
domain: algorithmic-trading
subdomain: real-time-architecture
tags:
- real-time-architecture
- adaptive-batching
- dynamic-tuning
- throughput-optimization
- database-sink
- backpressure
- ewma
brokers_frameworks: []
jurisdictions: [global] # technique is jurisdiction-agnostic
- version: "1.1.0"
+ version: "2.0.0"
author: algo-trading-skills-contributors
license: Apache-2.0
---
## When to Use
Invoke this skill when persistently writing high-volume tick feeds or trading
logs into downstream databases (TimescaleDB, ClickHouse, InfluxDB) or message
queues (Kafka, Redis Streams). Hardcoding a static batch size is wrong across
regimes — it causes high persistence latency when markets are quiet (waiting
for fixed batch limits to fill) and DB I/O overload during flash crashes.
The skill produces an `AdaptiveBatchTunerEngine` whose job is to scale the
batch size `B_t` and flush timeout `T_flush` in response to two signals:
- 1. **Smoothed queue fill ratio** `R_ewma = EWMA(Q_current / Q_capacity)`
- 2. **Smoothed downstream write latency** (target: `target_write_latency_ms`)
+ 1. **Smoothed batch fill ratio at the flush boundary**
+ `F_ewma = EWMA(depth_at_flush / B_t)` — did the producer fill the batch
+ before the timeout, or did the timeout cut a half-empty batch?
+ 2. **Smoothed downstream write latency** (target: `target_write_latency_ms`),
+ which acts as the throttle and outranks signal 1.
+ ### Why fullness, and not queue depth
+
+ The obvious signal — buffer occupancy against a nominal `queue_capacity` — does
+ not work in this architecture, and getting this wrong inverts the controller.
+ `add_item` hands the batch back the instant the buffer reaches `B`, so the
+ buffer depth is bounded by `B` by construction: `depth / queue_capacity` can
+ never exceed `B / queue_capacity` and therefore measures *the tunable*, not the
+ load. Tuning on it is positive feedback on `B` itself, and under saturation it
+ drives `B` down to `B_min` — the exact opposite of the intent. Batch fullness
+ at the flush boundary is the signal that actually separates "the producer
+ outran the batch size" (`F = 1.0` ⇒ high load) from "the timeout expired
+ half-empty" (`F < 1.0` ⇒ low load).
+
## When NOT to Use
- **Single-shot or batch-mode loads** — if the total record count is bounded
and known upfront, fine-tune one batch size statically. The adaptive engine
adds no value when input is finite.
- **Synchronous request/response flows** — this is fire-and-forget throughput
tuning, not per-request optimisation. Use RPC-style tuning instead.
- **Distributed stream processing stages** (Flink, ksqlDB, Spark Structured
Streaming) — they have their own internal flush predicates; this engine is
for the **client-side** write path into the sink, not for the consumer side
of an internal stage.
- **Latency-critical tick-to-trade paths** — those need synchronously bounded
writes; the entire batch-and-flush taxonomy is wrong. See
`tick-to-trade-latency-measurement` instead.
- **When the sink is a queue with built-in batching** (Kafka producer with
`linger.ms`/`batch.size`) AND the queue is the only consumer — let Kafka's
producer do the tuning.
+ - **When strict cross-batch write ordering is required across multiple
+ producer threads** — the engine detaches each batch before returning it and
+ runs `on_flush` outside its lock, so concurrent producers may write out of
+ order. Use a single producer thread, or serialise inside your callback.
## Prerequisites
- A **downstream sink write function** accepting batches of records (caller loop).
- **Capacity constants** declared up front:
- `B_min` — minimum batch size (records; default 10)
- `B_max` — maximum batch size (records; default 1000)
- `T_min` / `T_max` — flush-timeout bounds (default 50ms / 1000ms; aligned with
ClickHouse's adaptive-busy-timeout range)
- - `L_target` — downstream write-latency target (default 50ms)
- - **Queue capacity** consistent throughout the lifecycle; pass via `TuningConfig.queue_capacity`.
+ - `L_target` — downstream write-latency target (default 50ms). This is the
+ binding constraint on expansion; set it to a latency your sink can actually
+ sustain, not an aspiration.
+ - **Queue bounds**: `max_queue_size` is the hard cap that raises
+ `QueueFullError`; `queue_capacity` is the denominator of the exported
+ backpressure gauge and must be `<= max_queue_size` (validated).
- **Sink latency instrumentation** — caller must call `record_write_latency(ms)`
- on every successful (or failed) write.
+ on every write, successful or failed. Without it the throttle never fires and
+ the batch expands until it hits `B_max`.
+ - **A scheduler tick** if the producer can go quiet — see Workflow step 5.
## Workflow
1. **Construct the engine** with a `TuningConfig`:
```python
from batch_tuner import AdaptiveBatchTunerEngine, TuningConfig
tuner = AdaptiveBatchTunerEngine(TuningConfig(
min_batch_size=10,
max_batch_size=1000,
initial_batch_size=100,
target_write_latency_ms=50.0,
queue_capacity=2000,
max_queue_size=5000,
))
```
2. **Produce-loop pattern** (the contract):
```python
try:
batch = tuner.add_item(item)
except QueueFullError as exc:
handle_overload() # backpressure / degrade / drop
continue
if batch is None:
continue # not yet a flush boundary
t0 = time.monotonic()
try:
sink_write(batch)
finally:
tuner.record_write_latency((time.monotonic() - t0) * 1000.0)
```
- 3. **Adapt batch size** (under the hood):
- - **High load** (`R_ewma > 0.70`): `B_{t+1} = min(B_max, ⌊B_t × 1.5⌋)`; reduce `T_flush`.
- - **Low load** (`R_ewma < 0.10`): `B_{t+1} = max(B_min, ⌊B_t / 1.2⌋)`; extend `T_flush`.
- - **Deadband** (`0.10 ≤ R_ewma ≤ 0.70`): no tuning — EWMA smoothing prevents
- oscillation around the boundary.
+ The `finally` matters: a failed write is still a latency observation, and
+ skipping it on the error path is how the throttle goes blind exactly when
+ the sink is sick.
- 4. **Apply latency feedback** (overrides fill-ratio tuning under sustained pressure):
- - If `EWMA(L) > L_target`: `B_{t+1} = max(B_min, ⌊B_t × 0.8⌋)`. Latency throttle
- triggers **even inside the deadband**, which is its purpose.
+ 3. **Adapt batch size** (under the hood, evaluated at each flush boundary):
+ - **High load** (`F_ewma > 0.70` **and** `EWMA(L) ≤ L_target`):
+ `B_{t+1} = min(B_max, ⌊B_t × 1.5⌋)`; reduce `T_flush`.
+ - **Low load** (`F_ewma < 0.10`): `B_{t+1} = max(B_min, ⌊B_t / 1.2⌋)`;
+ extend `T_flush`.
+ - **Deadband** (`0.10 ≤ F_ewma ≤ 0.70`): no tuning — EWMA smoothing plus the
+ deadband prevents oscillation around the boundary.
+ 4. **Apply latency feedback** — this is what closes the loop:
+ - If `EWMA(L) > L_target`: `B_{t+1} = max(B_min, ⌊B_t × 0.8⌋)`, and the
+ high-load expansion branch is **barred** until latency returns under
+ target. The bar is not decorative: expansion multiplies by 1.5 while the
+ throttle multiplies by 0.8, and `1.5 × 0.8 = 1.2 > 1`, so without it one
+ throttle per flush can never undo one expansion and `B` ratchets to
+ `B_max` with sink latency pinned above target.
+ - The throttle fires **even inside the deadband**, which is its purpose.
+
5. **Flush triggers**:
- - **Threshold**: `#items ≥ current_batch_size`.
- - **Timeout**: `elapsed ≥ current_flush_timeout_sec`.
- - **Forced flush**: `tuner.flush_now()` for shutdown / checkpoint boundaries.
+ - **Threshold**: `#items ≥ current_batch_size` (evaluated inside `add_item`).
+ - **Timeout**: `elapsed ≥ current_flush_timeout_sec`. The engine owns **no
+ timer thread**, so this is only evaluated when you call in. If the
+ producer can stall — thin instrument, feed outage, the lull after the
+ close — buffered records would otherwise sit in memory indefinitely and
+ die with the process. Drive `tuner.flush_if_due()` from a scheduler at an
+ interval at or below `min_flush_timeout_sec`:
+ ```python
+ batch = tuner.flush_if_due() # None if not yet due / nothing buffered
+ if batch:
+ sink_write_with_metric(batch)
+ ```
+
+ - **Forced flush**: `tuner.flush_now()` for checkpoint boundaries. Returns up
+ to `current_batch_size` items and deliberately does **not** tune — a batch
+ that is partial because you asked for it says nothing about producer speed.
+
6. **Shutdown**:
```python
- leftover = tuner.close()
- if leftover: sink_write(leftover)
+ leftover = tuner.close() # drains the WHOLE buffer, not one batch
+ for chunk in chunked(leftover, sink_max_rows):
+ sink_write(chunk)
```
+ `close()` is not capped at `current_batch_size`, so the returned list may be
+ larger than your sink accepts in one call — chunk it. After `close()`,
+ `add_item` raises `RuntimeError`; call `reset()` to reuse the engine.
> Full procedure with rationale: `references/workflows.md`.
> Concrete numeric thresholds: `references/standards.md`.
> Operational checklist: `assets/checklist.md`.
## Decision Points
| Situation | Action |
|-----------|--------|
- | Sink write latency consistently > `L_target` | Engine is auto-throttling (×0.8). If still failing, lower `target_write_latency_ms` to force aggressive throttle, or activate circuit breaker upstream. |
+ | Sink write latency consistently > `L_target` | Engine is auto-throttling (×0.8) and expansion is barred. If latency still does not recover, the sink — not the batch size — is the problem: check IOPS, locks, connection pool. |
+ | `current_batch_size` pinned at `B_max` with latency under target | Correct behaviour: the sink absorbs everything you can give it. Raise `B_max` only if the sink documents a larger optimal write unit. |
+ | `current_batch_size` pinned at `B_min` | Either genuinely idle traffic, or the throttle is stuck on. Compare `ewma_write_latency_ms` with `target_write_latency_ms` before assuming idleness. |
| `QueueFullError` raised | Caller is producing faster than sink drains. Implement back-pressure (block producer), degradation (drop new items), or fall back to a slower sink. |
- | Fill ratio hovers at boundary (0.10 or 0.70) | Deadband is working — there should be **no** batch-size oscillation. Verify via `tuner.get_status().total_tuning_transitions`. |
- | Overflow `.record_write_latency(0)` — i.e. hot cache made flush look instant | Engine still trusts the smoothed EWMA — single low-variance samples won't reset the throttle. |
- | Sustained high fill with no rule violations | Lower `max_batch_size` for the workload, or migrate to a faster sink. |
- | Test opposing `add_item()` from many threads | Engine is thread-safe; one engine per sink. Do **not** instantiate per add. |
+ | `batch_fill_ratio_ewma` sits inside [0.10, 0.70] | Deadband is working — there should be **no** batch-size oscillation. Verify via `total_tuning_transitions`. |
+ | `record_write_latency` reports ~0 ms for every write | The sink is not actually acknowledging durability (see the ClickHouse `wait_for_async_insert = 0` pitfall). The throttle is blind; fix the instrumentation before trusting the tuner. |
+ | Producer thread can stall | You must call `flush_if_due()` on a timer, or the timeout trigger never fires. |
+ | Many threads calling `add_item()` | Engine is thread-safe; one engine per sink. Do **not** instantiate per add. Cross-batch write ordering is not guaranteed. |
## Common Pitfalls
+ - **Tuning on queue depth instead of batch fullness** → the controller inverts
+ and drives the batch *down* under load. See "Why fullness, and not queue
+ depth" above; this is the single most important design point in the skill.
+ - **Assuming `close()` returns one batch** → it drains the entire buffer, which
+ may exceed what your sink accepts in a single call. Chunk the result. (The
+ converse bug — capping the shutdown drain at `current_batch_size` — silently
+ strands records, which for an order log is data loss.)
+ - **Never calling `flush_if_due()`** → the flush timeout is dead weight, and a
+ stalled producer leaves records buffered until the process exits.
+ - **Skipping `record_write_latency()` on the error path** → the throttle goes
+ blind exactly when the sink is failing.
- **Unbounded max batch size** → queued memory can exceed sink RAM. Always set
`max_queue_size`; the engine raises `QueueFullError` at the cap.
- - **Ignoring latency feedback** → batch escalates under DB lock contention,
- exhausting the connection pool. The `target_write_latency_ms` knob is the
- tripwire; tune it to your sink's healthy floor.
- - **Rapid oscillations** (solved by the new deadband / EWMA combination): if
- fill ratio flickers around one threshold and the engine thrashes, increase
- `fill_ewma_alpha` (lower alpha = more smoothing) — but never go to 0.
+ - **Feeding a non-finite latency** → rejected with `ValueError`. A `NaN` would
+ otherwise poison the EWMA permanently (`NaN > target` is always `False`,
+ silently disabling the throttle for the life of the process) and serialise as
+ invalid JSON in the metrics export.
+ - **Rapid oscillations**: if fullness flickers around one threshold and the
+ engine thrashes, lower `fill_ewma_alpha` (lower alpha = more smoothing) —
+ but never to 0.
- **Latency oscillation** under repeated `record_write_latency()` spikes: if
- smoothed latency is at `L_target - epsilon` and individual writes push it
- slightly above, the engine is repeatedly halving-and-climbing. Adjust
- `target_write_latency_ms` to align with the **typical** sink latency, not
- the ideal.
- - **Synchronous batch-handling logic**: do not block on a slow sink while
- holding the queue; return the batch and flush outside any locks. The engine
- holds its lock only for batch extraction + tuning decisions.
- - **Caller-driven capacity drift**: the previous API required `queue_capacity`
- on every call. The new API bakes it into `TuningConfig`; verify upstream that
- the call sites have been migrated.
- - **Cold-start EWMA** (`fill_ewma_ratio == 0`): smoothing starts at zero and the
- first burst of adds looks like "low load" until the EWMA warms up. Set
- `initial_fill_ewma` (deferred) or accept that the first minute of operation
- is slightly micro-tuned.
+ smoothed latency sits at `L_target - epsilon` and individual writes push it
+ slightly above, the engine repeatedly throttles and climbs. Set
+ `target_write_latency_ms` from the **typical** sink latency, not the ideal.
+ - **Blocking inside `on_flush`**: the callback runs outside the engine lock, so
+ it will not deadlock or block other producers — but it does run on the
+ calling producer's thread, so a slow sink write there still stalls that
+ producer. Prefer the returned-batch pattern for the actual write.
- **ClickHouse async_insert `wait_for_async_insert = 0`** (fire-and-forget):
- the engine's `record_write_latency()` will see 0 ms for "successful" flushes,
- masking real DB-side problems. Use `wait_for_async_insert = 1`.
+ `record_write_latency()` will see near-0 ms for "successful" flushes, masking
+ real DB-side problems. Use `wait_for_async_insert = 1`.
## Verification
Run the unit tests:
```bash
python -m unittest discover -s skills/adaptive-batch-size-tuning-under-load/scripts -v
```
- What they assert:
+ 43 tests. What they assert:
- - Low-load regime shrinks batch size below initial.
- - High-load regime expands batch size above initial.
- - Latency throttle fires at and above target.
- - EWMA smoothing with `alpha = 0.5` produces the expected smoothed value.
- - **Bounded queue**: `add_item` past `max_queue_size` raises `QueueFullError`.
- - **Deadband**: in the [0.10, 0.70] zone, no tuning transitions occur.
- - `flush_now()` extracts up to `current_batch_size` items; `close()` drains.
- - `reset()` returns to initial state.
- - `RecordingWriteLatency` rejects negative values.
- - Configuration validation: `min > max`, `alpha in (0, 1]`, threshold ordering.
- - Status JSON-serializable (Prometheus-ready).
+ - **Control-law direction** (the regression that matters): saturating load
+ expands the batch toward `B_max` and shortens `T_flush`; quiet,
+ timeout-driven load shrinks it and lengthens `T_flush`.
+ - **Closed loop**: with a batch-size-dependent sink latency, the equilibrium
+ batch size settles strictly inside `(B_min, B_max)` with smoothed latency at
+ or under target.
+ - **Deadband**: batches cut ~50% full produce zero tuning transitions; exact
+ boundary values (0.10, 0.70) are inside the deadband.
+ - **Latency throttle**: fires above target, not at exactly target, stops at
+ `B_min`, and the EWMA is seeded with its first observation (hand-computed
+ expected values, not a re-derivation of the implementation).
+ - **Shutdown**: `close()` drains records the batch size would have stranded,
+ preserves order, is idempotent, and `add_item` after `close()` raises.
+ - **Flush triggers**: `flush_if_due()` releases an idle buffer; `flush_now()`
+ returns a partial batch and does not tune.
+ - **Bounded queue**: `add_item` past `max_queue_size` raises `QueueFullError`
+ and does not buffer the rejected item.
+ - **Validation**: bound ordering, alpha range, non-finite/negative latency,
+ `queue_capacity > max_queue_size`, and mis-signed tuning multipliers.
+ - **Concurrency**: an `on_flush` callback may re-enter the engine without
+ deadlocking; a raising callback does not lose the batch; 8 concurrent
+ producers × 500 items lose and duplicate nothing.
+ - **Status** is strict-JSON serializable (`allow_nan=False`, Prometheus-ready).
Confirm with the operational checklist in `assets/checklist.md` before deploying.
## Success Criteria
A tuning engine is considered **healthy in production** when:
1. `total_tuning_transitions` is bounded — under steady load it should be O(tens
per hour), **not** O(thousands). Sustained high transition count ⇒ noisy
upstream or mis-tuned `target_write_latency_ms`.
2. Sink write latency P99 < `L_target` over a rolling 1-hour window.
3. `QueueFullError` rate is 0 (upstream throughput matches or exceeds sink
drain). If non-zero, page on-call.
4. `current_batch_size` and `current_flush_timeout_sec` settle inside their
- ranges within 5 minutes of traffic beginning.
+ ranges within 5 minutes of traffic beginning — and are **not** pinned at
+ `B_min` while the feed is busy, which is the signature of a mis-wired
+ control signal.
5. `get_status()` exports cleanly to a JSON metrics pipeline.
## Related Skills
- `kafka-based-tick-distribution-at-scale` — the parent batch architecture for
Kafka paths; this skill is the *client-side* tuning companion.
- `producer-consumer-tick-pipeline` — the broader produce/consume pipeline;
this skill is the leaf that decides *when* to flush to the sink.
- `tick-buffering-burst-handling` — what to do when the queue genuinely
overflows; `QueueFullError` from this skill should be the trigger.
- `backpressure-drop-degrade-policy` — design the policy that decides what
to do when `QueueFullError` fires.
+ - `graceful-shutdown-draining-in-flight-ticks` — the shutdown counterpart;
+ `close()` is this engine's contribution to that drain.
- `kill-switch-and-drawdown-circuit-breakers` — strategy-level circuit breaker
upstream of this engine; pair them so strategy stops suppress writes.
- `latency-monitoring-percentile-based-slas` — for monitoring the sink's
P99/P999 against `L_target`.
- `model-inference-latency-budget-for-live-trading` — analogous pattern for
model inference instead of DB writes.