Repository navigation
Conversation
bgentry
force-pushed
the
bg/parallel-job-completion
branch
3 times, most recently
from
October 6, 2026 18:49
b67f172 to
c744f1d
Compare
Add optional `ExecutorJobCompletionConcurrency` and `PilotJobCompletionConcurrency` interfaces through which an executor and a pilot can declare how many `JobSetStateIfRunningMany` calls they can safely run at once. Implementations that don't implement them are treated as serial. The pgx and `database/sql` pool executors and the standard pilot report two. Their transaction and subtransaction executors report one, since a single transaction can't run statements concurrently. SQLite and other drivers don't opt in. A new driver test runs the advertised number of disjoint completion batches concurrently through a pool executor and checks that a transaction never claims more than one.
The batch completer persists one `JobSetStateIfRunningMany` call at a time. Under sustained load a second full backlog builds while the first query is still in flight, so producers wait on the backlog and workers sit idle even though the database has spare capacity. When both the executor and the pilot declare completion concurrency, run up to two completion queries at once, each capped at the existing 5,000-job batch size. A second query starts only once its own backlog is full, so sparse traffic keeps the existing coalescing and connection use. SQLite, transactions, and pilots that don't opt in stay serial. Concurrent batches could otherwise hold different executions of the same job, and PostgreSQL may lock their rows out of submission order, letting an older result overwrite a newer attempt. The completer tracks in-flight and retrying job IDs so a later result for the same job waits behind the earlier batch while independent batches still persist concurrently. Shutdown waits for active queries and makes one bounded final flush attempt. Subscribers may see completion events for different jobs in a different order than before, and the completer may briefly use one more database connection. Tests block the first query and check that a sparse second batch stays queued through a flush interval and starts as soon as it reaches the batch threshold, cover capability negotiation, and force a duplicate job ID to interleave across batches.
bgentry
force-pushed
the
bg/parallel-job-completion
branch
from
October 7, 2026 13:54
c744f1d to
adba3de
Compare
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
The batch completer persists finished jobs with one
JobSetStateIfRunningManyquery at a time. That works well under light load, but under sustained heavy load it becomes the bottleneck. While one 5,000-job batch is being written, workers keep finishing jobs and a second full batch piles up behind it. When the backlog reaches its limit, completions start waiting on the completer, so workers sit idle even though the database could easily run another query.Here's an example: a client works down a large queue of short jobs. Each completion query takes a few milliseconds, and in that time workers finish thousands more jobs. The next batch is already full by the time the current query returns, so the completer is always behind, and job throughput is capped by the latency of one query at a time.
With this change, the completer can run up to two completion queries at once, each still capped at 5,000 jobs. The second query starts only after its own backlog has filled a batch, so light and sparse traffic keeps coalescing exactly as it does today and doesn't use an extra connection. Concurrency is opt-in through two new optional interfaces,
riverdriver.ExecutorJobCompletionConcurrencyandriverpilot.PilotJobCompletionConcurrency. The pgx anddatabase/sqlpool executors and the standard pilot report two. Transaction executors report one, since one transaction can't run statements concurrently. SQLite and any pilot that doesn't implement the interface stay serial.Running two batches at once introduces a hazard. Two batches could hold different executions of the same job, and PostgreSQL might lock their rows out of order, which would let an older result overwrite a newer attempt. To prevent that, the completer tracks in-flight and retrying job IDs. A later result for the same job waits until the earlier batch finishes, while batches for different jobs still persist concurrently. On shutdown, the completer waits for active queries and then makes one bounded final flush attempt, as before.
There are two user-visible side effects, both noted in the changelog. Subscribers may get completion events for different jobs in a different order than before, though events for the same job stay in order. And the completer may briefly use one extra database connection.
I compared
river benchon master (6cdaf38) and on this branch, with runs alternating between the two builds against the same local database:river bench -n 1000000(5 runs each)river bench --duration 30s(3 runs each)The
-nmode inserts all the jobs up front and then measures how fast they're worked down, which is the backlog case this change targets. The--durationmode inserts and works at the same time, and there the inserter limits throughput, so the gain is small. These numbers come from one laptop (Apple Silicon M4, 14 cores, Homebrew PostgreSQL 18.6 on localhost, Go 1.27.1). Absolute numbers will differ on other machines and on a networked database, where query latency is higher and overlapping queries may help more.This also brings Go's completion throughput back in line with River's Rust implementation, which persists completion batches the same way.
The first commit adds the capability interfaces and a driver test that runs the advertised number of disjoint completion batches concurrently through a pool executor. The second commit changes the completer, with tests for capability negotiation, a sparse second batch that stays queued until it's full, and a job whose results are forced to interleave across concurrent batches.