Summary
PgmqTransport::startConsumer() reads one batch per poll signal. A signal comes from an insert notification or from the poll timer (default 5 s). When a backlog is already in the queue, no new inserts happen, so after one batch the consumer sleeps until the next tick, even though the queue still holds messages.
A queue of N messages takes about N / batchSize × pollInterval to drain, with the worker idle. This is the main cause of slow pipelines on top of the PGMQ transport.
Where
src/MessageBus/Pgmq/PgmqTransport.php L203–L240: while ($iterator->continue()) → readBatch(count: $this->batchSize) → handle each message → wait for the next poll signal.
thesis/pgmq's own Consumer has the same loop: thesis-php/pgmq#19.
Measurements
Our workload fans out ~450 messages into a queue at once; handling a message takes ~14 ms. 2 workers per queue, defaults (batchSize 10, pollInterval 5 s).
- The pick-up rate was exactly 2 workers × 10 messages / 5 s = 40 messages per 10 s. We measured it from the message id (UUIDv7 creation time) against the "About to handle" log time. The wait in the queue was p50 56 s, p90 267 s, while the handling time of that queue summed to 131 s over the whole run.
- With the fix below (same data, same code otherwise):
| scenario |
before |
after |
| 1 unit of work (~450 messages per queue, several queues in a chain) |
60 s |
8 s |
| 10 such units at once |
421 s |
75 s |
| CPU while idle (consumers + Postgres) |
~0 % |
~1 % (unchanged) |
The results computed by the handlers were identical before and after (row counts and sums of every numeric column of the affected tables, 76 runs).
A shorter pollInterval hides the problem, but costs a lot while idle. With 100 ms on ~30 queues we measured 17 % of a core in the consumers and 23 % in Postgres, doing nothing.
Proposed fix
After a full batch, request the next poll immediately. A partial batch means the queue is (most likely) empty, and the consumer waits as before. There are no extra queries when idle.
if (\count($messages) === $this->batchSize && !$polls->isComplete()) {
$polls->pushAsync(null)->ignore();
}
PR: #6. Integration test (red without the fix): #8.
Summary
PgmqTransport::startConsumer()reads one batch per poll signal. A signal comes from an insert notification or from the poll timer (default 5 s). When a backlog is already in the queue, no new inserts happen, so after one batch the consumer sleeps until the next tick, even though the queue still holds messages.A queue of
Nmessages takes aboutN / batchSize × pollIntervalto drain, with the worker idle. This is the main cause of slow pipelines on top of the PGMQ transport.Where
src/MessageBus/Pgmq/PgmqTransport.phpL203–L240:while ($iterator->continue())→readBatch(count: $this->batchSize)→ handle each message → wait for the next poll signal.pollInterval ?? TimeSpan::fromSeconds(5),batchSize = 10;ChannelWatcheronpgmq.q_<queue>.INSERT, fired by inserts only.thesis/pgmq's ownConsumerhas the same loop: thesis-php/pgmq#19.Measurements
Our workload fans out ~450 messages into a queue at once; handling a message takes ~14 ms. 2 workers per queue, defaults (
batchSize10,pollInterval5 s).The results computed by the handlers were identical before and after (row counts and sums of every numeric column of the affected tables, 76 runs).
A shorter
pollIntervalhides the problem, but costs a lot while idle. With 100 ms on ~30 queues we measured 17 % of a core in the consumers and 23 % in Postgres, doing nothing.Proposed fix
After a full batch, request the next poll immediately. A partial batch means the queue is (most likely) empty, and the consumer waits as before. There are no extra queries when idle.
PR: #6. Integration test (red without the fix): #8.