Log Sampling for High-Volume Queues
At a few hundred jobs per second, a worker fleet can emit more log data than the rest of the platform combined, and this guide shows how to reduce it without losing the lines incidents depend on, as part of Structured Logging for Workers in Observability & Monitoring for Job Queues. The principle is simple — keep every failure, keep a representative slice of successes, and keep each sampled job's lines together — but the implementation has several traps that quietly discard exactly the wrong lines.
Problem Statement
An events pipeline processes 4,000 jobs per second across 60 workers. Each job logs about eight lines (start, fetch, transform, three writes, publish, completion), so the fleet emits roughly 32,000 lines per second — 2.7 billion a day, costing more in log storage than the compute running the jobs. An earlier attempt to cut cost set the global log level to warning, which removed the context lines engineers needed during the next incident and did nothing about a retry storm that produced 50,000 warnings a minute. You want volume down by an order of magnitude, every failed job's complete story preserved, statistically useful samples of successful jobs, and protection against sudden floods.
Prerequisites
- Structured JSON logs with job context bound on every line (see correlating logs with job IDs in Celery or the BullMQ equivalent).
- Metrics for job counts, failures, and durations — sampling logs is only safe when metrics already count every job.
- A logging library that supports filters or processors (structlog, pino, Logback, zap), and optionally a collector that can sample (Vector, OpenTelemetry Collector).
- Agreement on what the logs are for: incident investigation of specific jobs, not counting.
Step 1 — Measure Where the Volume Comes From
Sample intelligently only after knowing which lines dominate. Group by message and level over a representative hour.
# Top 10 messages by line count in the last hour
topk(10, sum by (msg, level) (count_over_time({service="events-pipeline"} | json [1h])))
A typical result: "wrote batch" (debug-worthy detail logged at info) is 40% of lines, "fetched source" 20%, "job succeeded" 12%, retries 3%, and errors well under 1%. The biggest wins are usually not sampling at all but moving chatty per-step lines to debug and emitting a single completion line with the step timings as fields.
Step 2 — Collapse Per-Step Lines into One Completion Line
Replace chatty progress lines with timings recorded as fields on the completion line. This keeps the information and removes most of the lines.
import time, structlog
log = structlog.get_logger()
class StepTimer:
def __init__(self):
self.timings: dict[str, int] = {}
def step(self, name):
timer = self
class _Ctx:
def __enter__(self): self.t0 = time.perf_counter()
def __exit__(self, *exc):
timer.timings[f"{name}_ms"] = int((time.perf_counter() - self.t0) * 1000)
return _Ctx()
def process_event_batch(batch_id: str) -> None:
t = StepTimer()
with t.step("fetch"):
rows = fetch_source(batch_id)
with t.step("transform"):
out = transform(rows)
with t.step("write"):
write_batches(out)
with t.step("publish"):
publish(out)
log.info("job succeeded", rows=len(rows), **t.timings) # one line, all the detail
Per-step lines can stay at debug for local development. In the scenario, this change alone takes eight lines per successful job down to one — an 85% reduction before any sampling.
Step 3 — Sample Successes by Job, Never Failures
Sample the remaining success lines, but make the decision per job, not per line, so a sampled job's lines are all kept together. A hash of the job id gives a stable decision without shared state: every line from the same job hashes the same way, on any host.
import hashlib
SUCCESS_SAMPLE_RATE = 0.05 # keep 5% of successful jobs' info lines
def keep_job(job_id: str, rate: float) -> bool:
h = int.from_bytes(hashlib.blake2b(job_id.encode(), digest_size=8).digest(), "big")
return (h / 2**64) < rate
def sampling_processor(logger, method_name, event_dict):
level = event_dict.get("level", method_name)
if level in ("warning", "error", "critical"):
return event_dict # never sample problems
job_id = event_dict.get("job_id")
if job_id is None:
return event_dict # non-job lines: leave alone
if event_dict.get("attempt", 1) > 1:
return event_dict # retried jobs: keep full story
if keep_job(job_id, SUCCESS_SAMPLE_RATE):
event_dict["sampled"] = SUCCESS_SAMPLE_RATE # lets queries re-weight counts
return event_dict
raise structlog.DropEvent
structlog.configure(processors=[structlog.contextvars.merge_contextvars,
structlog.processors.add_log_level,
sampling_processor,
structlog.processors.TimeStamper(fmt="iso", utc=True),
structlog.processors.JSONRenderer()])
Two rules protect failures. Warning and error lines always pass. And any job on its second or later attempt keeps all its lines, so a job that eventually fails has its complete history — including the info lines from its earlier attempts that, unfortunately, were decided before anyone knew it would fail. Step 5 closes that gap.
Step 4 — Rate-Limit Repetitive Lines During Floods
Sampling by job does not help when one failure mode produces the same warning thousands of times a minute — a downstream outage makes every job log "connection refused". A per-message token bucket caps each distinct message and emits a summary of what was suppressed.
import threading, time
from collections import defaultdict
class MessageRateLimiter:
def __init__(self, per_minute: int = 300):
self.cap, self.lock = per_minute, threading.Lock()
self.window = int(time.time() // 60)
self.counts, self.dropped = defaultdict(int), defaultdict(int)
def __call__(self, logger, method_name, event_dict):
key = (event_dict.get("level"), event_dict.get("msg"), event_dict.get("error_class"))
with self.lock:
now = int(time.time() // 60)
if now != self.window:
for (lvl, msg, ec), n in self.dropped.items():
logger.warning("log lines suppressed", suppressed_msg=msg,
suppressed_level=lvl, error_class=ec, count=n)
self.window, self.counts, self.dropped = now, defaultdict(int), defaultdict(int)
self.counts[key] += 1
if self.counts[key] > self.cap and event_dict.get("final") is not True:
self.dropped[key] += 1
raise structlog.DropEvent
return event_dict
Note the exemption: lines marked final=true (a job that gave up permanently) are never rate-limited, because each one represents distinct lost work. Retry warnings are fungible; final failures are not. Metrics still count every retry, so the outage remains fully visible on dashboards — see alerting on queue backlog with Prometheus.
Step 5 — Tail-Sample in the Collector for Complete Failure Stories
Head sampling in the process decides before the job's outcome is known, so a job sampled out at attempt 1 loses its early info lines even if it later fails. Tail sampling buffers a job's lines briefly and decides after seeing whether any line was a failure. Vector can do this with a reduce transform grouped by job id.
# vector.toml — hold each job's lines until it completes, then keep or drop the group
[transforms.parse]
type = "remap"
inputs = ["worker_logs"]
source = '. = parse_json!(.message)'
[transforms.group_by_job]
type = "reduce"
inputs = ["parse"]
group_by = ["job_id"]
ends_when = '.msg == "job succeeded" || .final == true'
expire_after_ms = 120000 # flush incomplete groups after 2 min
merge_strategies.lines = "array" # keep individual lines inside the group
[transforms.keep_interesting]
type = "filter"
inputs = ["group_by_job"]
condition = '''
includes(array!(.level), "error") || includes(array!(.level), "warning")
|| (to_int(crc32(string!(.job_id))) % 100) < 5
'''
Tail sampling costs memory in the collector (in-flight jobs' lines are held until completion) and only works when all of a job's lines pass through one collector instance — route by job id or run the collector per node and accept that a job's attempts on different nodes are decided separately. Use it for the jobs where complete failure histories matter most, and keep head sampling elsewhere.
Verification
Confirm the volume reduction and, more importantly, that nothing important was lost:
# Volume before/after (lines per second)
sum(rate({service="events-pipeline"} [5m]))
# Every final failure is present: compare with the metric that counts them
sum(count_over_time({service="events-pipeline"} | json | final="true" [1h]))
sum(increase(job_failures_total{service="events-pipeline", final="true"}[1h]))
The two failure counts must match exactly. For successes, estimate totals from sampled logs by weighting (count / sampled) and check they track job_completed_total within a few percent.
Gotchas & Edge Cases
Sampling by line instead of by job. Random per-line sampling keeps fragments of many jobs and the complete story of none. Always key the decision on job id.
Changing the rate breaks comparisons. Record the rate on each kept line (sampled field) so queries can re-weight counts across a rate change.
Dropping lines needed for audits. Some jobs (payments, permission changes) may have audit requirements. Exempt them by job name, or write audit events to a separate, unsampled stream.
Suppression summaries are themselves logs. Make sure the rate limiter's summary line is exempt from rate limiting, or a flood silences its own report.
FAQ
Why not just raise the log level in production? Level is a blunt instrument: it removes context lines from failing jobs along with noise from healthy ones. Sampling by job removes the noise and keeps failing jobs complete.
What sample rate should I use? Enough successful jobs per job type per minute to see normal behaviour — often 1–10%. Low-volume job types can stay at 100%; make the rate a per-job-name setting.
How do I debug a specific customer's job if its logs were sampled out? Add a force-keep override: a small set of tenant or customer ids (from a feature flag or config) whose jobs bypass sampling entirely. When support escalates an issue, add the customer to the set for a day, reproduce, and remove them. Pair it with the metrics and traces that were never sampled for that job — its duration, outcome, and trace are still available even when its info lines were not kept.
Can I sample traces and logs together? Yes — use the trace's sampling decision (the sampled flag in the trace context) as the log keep decision for successes, so sampled traces always have their logs.
Related
- Structured Logging for Workers — the schema sampling depends on.
- Shipping Worker Logs to Loki — collection and storage after sampling.
- Prometheus Metrics for Workers — the complete counts that make sampling safe.
- Preventing Retry Storms After an Outage — the floods rate limiting protects against.