Skip to content

perf(subscriber): take arrived rows without a task per row - #101

Merged
nhobin219 merged 1 commit into
mainfrom
recv-many-no-loop-turn
Oct 4, 2026
Merged

nhobin219 merged 1 commit into
mainfrom
recv-many-no-loop-turn

Conversation

@nhobin219

Copy link
Copy Markdown
Owner

recv_many learned whether a row had arrived by starting a read (ensure_future) and giving the loop a turn (sleep(0)), once per row. That made batches() cost more per row than recv, although its docstring says it costs less.

Now, for rows from the socket, it checks websockets' queue of received frames, recv_messages.frames.queue. When a whole message is at the head of that queue, the read pops it without suspending, so no task or loop turn is needed. When it isn't, the batch ends and nothing is read ahead. Rows from a catch-up still use the read-ahead, since cancelling the generator would end it.

Client CPU per row, 100k frames already queued:

recv recv_many(500)
before 2.2–2.9 µs 10.6–13.9 µs
after 1.45 µs 1.3 µs

Private API.

  • websockets is capped at >=14,<18. The queue has the same shape in 14 through 17.1.
  • test_the_queue_it_looks_at_is_websockets runs against a real server and fails if a release moves the queue, so raising the cap only needs that test to pass.

Tests (each fails against the old loop or without the fin check):

  • test_a_batch_from_the_socket_starts_no_task
  • test_from_the_socket_nothing_is_read_ahead, which replaces the read-ahead test
  • test_a_fragmented_message_ends_the_batch_and_loses_nothing
  • test_the_queue_it_looks_at_is_websockets

Docs:

  • CHANGELOG: an Unreleased entry.
  • SPEC §4 bounds table: "plus one during a catch-up".
  • The recv_many docstring and API.md were already correct once this lands.

The OTel example still uses recv. Receiving now costs about 1.4 µs a row against about 8 µs to convert and queue each span, and OTel's processors take one record per call, so batches() wouldn't change much there.

🤖 Generated with Claude Code

https://claude.ai/code/session_01FsSDkeb5rVAxA1FSmKKfQi

recv_many learned whether a row had arrived by starting a read and
giving the loop a turn, once per row: 10.6-13.9 us of client CPU a row,
against recv's ~2.5, so batches() cost more than recv. It now looks at
websockets' queue of received frames, which a read pops without
suspending, and costs 1.3 us a row against recv's 1.45. A catch-up's
rows still go through the read-ahead.

The queue is private websockets API: websockets is capped below 18, and
test_batches checks its shape against the installed release.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01FsSDkeb5rVAxA1FSmKKfQi
@nhobin219
nhobin219 merged commit 6f89af7 into main Oct 4, 2026
6 checks passed
@nhobin219 nhobin219 mentioned this pull request Oct 4, 2026
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.

1 participant