Skip to content

Lease tasks beside the payload instead of rewriting it - #69

Merged
codingjoe merged 11 commits into
mainfrom
codingjoe-redis-cpu-per-task
Oct 6, 2026
Merged

codingjoe merged 11 commits into
mainfrom
codingjoe-redis-cpu-per-task

Conversation

@codingjoe

@codingjoe codingjoe commented Oct 1, 2026 •

Copy link
Copy Markdown
Owner

The acquire script decoded and re-encoded every task payload to stamp it RUNNING. That is two full JSON passes on the Redis server for each fetch. It was the largest per-task cost in the backend, and it grew with the payload size.

  • The task hash keeps the payload as it was enqueued, and stores the lease in two fields beside it: lease_worker and lease_started_at.
  • acquire(), peek() and the reaper apply that lease on read.
  • requeue() clears the lease when a task returns to the deferred set.
  • A hash from an older release still reads, and keeps its attempt count and start time, so a rolling upgrade does not reset the retry limit.
  • A reaped task now records started_at from its lease. The old acquire script intended that and never did it.
  • EVALSHA per acquire drops from 10.30 to 6.01 µs at 447 B, and from 34.14 to 12.11 µs at 5.4 KB. Server CPU per task in a 20,000-task drain drops from 23.01 to 19.70 µs.
  • The telemetry PUBLISH, the history eviction and the stored result stay as they are. Each was measured before that decision.

Fixes #66

The acquire script decoded and re-encoded every task payload to stamp it
RUNNING, so a fetch cost two full JSON passes on the Redis server — the
single largest per-task cost in the backend, and one that grows with the
payload.

The task hash now keeps the payload exactly as it was enqueued and stores
the lease in two fields next to it: the worker holding it and the time it
started. acquire() applies that lease to the task result it returns, peek()
applies it when reading RUNNING tasks, and the reaper folds it into the
failure it writes, so lease visibility is unchanged for the inspector and
the reaper. requeue() clears the lease when a task returns to the deferred
set.

Measured on a dedicated Redis 8.10 instance, one worker process, 20,000
tasks, EVALSHA service time per acquire:

  payload    before     after
  447 B      10.90 us    6.23 us   (-43%)
  5.4 KB     36.03 us   12.38 us   (-66%)

End-to-end drains of 20,000 tasks drop from 30.7 to 21.2 us of server CPU
per task (commandstats), and from 172 to 149 us per task of Redis process
CPU, with the same 9.0 Redis commands per task.

Refs #66
…er-task

# Conflicts:
#	threadmill/backends/lua/reaper.lua
The fetch script popped a task, leased it, and returned the payload for the
worker to deserialize. A payload this code cannot read — a document another
version wrote, a task module or retry callback gone from the code base — then
raised out of acquire(), which the consumer loop does not catch, so the
consumer thread died while the id stayed leased where every reaper pass
re-claimed and re-skipped it. The reaper's own copy of that read had the same
hole, and enumerating exception types twice already missed ImportError.

Reading a stored payload is now one rule instead of a list of failures: the
acquire script returns the task id beside the payload, and both the fetch path
and the reap path catch any read failure, log it with the task ID, and leave
the task stored. Nothing is deleted: dropping the hash would turn a deploy or
an environment skew into permanent, unaudited task loss, and the fetch path can
run before any attempt exists. The reap path distinguishes a payload that is
already gone, which a task finishing inside the claim window explains, from one
that will not parse.

peek() applies the same rule, skipping an entry it cannot read so the inspector
still renders the rest of the page, and _decode_text accepts the str replies a
REDIS_URL carrying decode_responses=True produces.

Refs #66, #69
Comment-only: the rule the guards apply lives in the log messages and the
method docstrings, so the paragraphs explaining it at each call site are gone.
No behaviour change; the token streams are identical once comments and
docstrings are removed.
The guards caught Exception with a BLE001 suppression, which hid the rule
behind a linter pragma. They now catch UNREADABLE_PAYLOAD_ERRORS: the failures
a stored payload raises when this code cannot read it as a task, each one
reachable from a payload and pinned by a test, including the AttributeError a
document that parses but is not a TaskResult raises.
A Redis reply is bytes by default, and the URL query overrides any
decode_responses kwarg, so the type cannot be pinned from the client
construction. The decoder now states the one type it takes, and every payload
read decodes before parsing: handing raw bytes to json.loads was the slowest of
the three variants measured, because json.loads(bytes) runs its own encoding
detection, while decode-then-parse measures the same as a str-configured client.

Each decode sits inside the read's guard, so a reply of the wrong shape is
reported with its task id and left as stored, rather than raising out of
acquire() or the reaper on a path the caller does not catch.
decode_responses=False is the supported client, so a reply is bytes and the
one-expression `_decode_text` wrapper only hid a call to `.decode()`. Every
reply now decodes at its read site, which removes the last str tolerance in the
module: worker_telemetry decodes its pubsub data directly.

The unreadable-payload matrix keeps its nine cases; the three that wrapped a
literal payload in a one-line function carry the literal instead, applied by
one `_apply_poison` for both shapes.
…er-task

# Conflicts:
#	tests/backends/test_redis.py
#	threadmill/backends/redis.py
EAFP on a field whose absent, cleared and unparseable shapes all mean the same
thing here: no start time to apply. Asking permission first only restated the
answer.
@codingjoe
codingjoe merged commit d74d4e7 into main Oct 6, 2026
4 checks passed
@codingjoe
codingjoe deleted the codingjoe-redis-cpu-per-task branch October 6, 2026 20:21
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Threadmill uses 4 to 5 times more Redis CPU than dramatiq for each task

1 participant