Skip to content

PgmqTransport: consumer reads one batch per poll signal, a backlog waits pollInterval between batches #2

Description

@iGrog

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.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions