Skip to content

How Laravel queue prefetch fills a batch

One Laravel worker, three queues: with prefetch 1 every job costs a pop; with prefetch 'auto' the batch doubles after every two full pops on a backlog of short jobs, up to about 250 ms of work, and stays at one job when a job alone is longer.

Updated View as Markdown

With prefetch 1 a worker makes one pop for every job. With prefetch 'auto' it starts at one job per pop, doubles the batch after every two full pops, and stops at about 250 ms of work or 16 jobs. A job that alone takes longer than 250 ms keeps the batch at one job.

Run it

php artisan example:prefetch   # in examples/apps/laravel

What it printed against a 2.0.1 node (setup):

broker http://localhost:6632

a) prefetch 1: 40 jobs of 10 ms, one worker
  jobs           40
  pops           40
  jobs per pop   1.0
  batch sizes    1 ×40
  ok: every pop took one job

b) prefetch 'auto': a backlog of 200 jobs of 10 ms, one worker
  jobs           200
  pops           21
  jobs per pop   9.5
  batch sizes    1 1 2 2 4 4 8 8 13 13 16 ×3 14 12 ×6 10
  one job took about 18 ms, its ACK included: 250 ms holds 13
  ok: the batch doubled after every two full pops: 1 1 2 2 4 4 8 8
  ok: no batch was more than twice the one before it, or above 16
  ok: a pop served 9.5 jobs on average, 3 or more

c) prefetch 'auto': 10 jobs of 300 ms, one worker
  jobs           10
  pops           10
  jobs per pop   1.0
  batch sizes    1 ×10
  ok: every pop took one job: one of them is more than 250 ms of work

PASS: 5 checks
The jobs each pop took, for one worker with prefetch auto on a backlog of 200 jobs of 10 ms, 21 pops in order: 1, 1, 2, 2, 4, 4, 8, 8, 13, 13, 16, 16, 16, 14, then 12 six times, and 10.0481216jobs in the poppop 1pop 2pop 3pop 4pop 5pop 6pop 7pop 8pop 9pop 10pop 11pop 12pop 13pop 14pop 15pop 16pop 17pop 18pop 19pop 20pop 211122448813131616161412121212121210
The batch doubles after every two full pops until it holds about 250 ms of work. From there it follows the measured job time, between 12 and the ceiling of 16. The last pop took the last 10 jobs. Source: examples/apps/laravel, php artisan example:prefetch
Phase Pops for the jobs Why
a) prefetch 1 40 for 40 each pop asks for one job
b) 'auto', 10 ms jobs 21 for 200 the batch grew to 250 ms of work
c) 'auto', 300 ms jobs 10 for 10 one job is already more than 250 ms

Past the doubling, the batch follows the measured job time, so its size depends on the machine. Here a job took about 18 ms with its ACK, and the batch moved between 12 and the ceiling of 16. In other runs on the same laptop a job took 15 to 19 ms.

The program

The job records the lease of the pop that brought it:

examples/apps/laravel/app/Jobs/ProbeBatch.phpphp
final class ProbeBatch implements ShouldQueue
{
    use Queueable;

    public function __construct(public int $seq, public int $millis) {}

    public function handle(): void
    {
        // The broker leases a pop's jobs together, under one lease:
        // jobs that share a lease id came in the same pop.
        Journal::record($this->job, [
            'event' => 'ran',
            'seq' => $this->seq,
            'lease' => $this->job->getQueenMessage()['leaseId'],
        ]);
        usleep($this->millis * 1000);
    }
}

The command queues each phase’s jobs before it starts the worker, so phase b has a real backlog:

examples/apps/laravel/app/Console/Commands/Examples/PrefetchExample.phpphp
final class PrefetchExample extends ExampleCommand
{
    protected $signature = 'example:prefetch';

    protected $description = 'How prefetch fills a batch: prefetch 1, then auto';

    protected function example(): void
    {
        $this->line("\na) prefetch 1: 40 jobs of 10 ms, one worker");
        $sizes = $this->drain('queen', 'one', 40, 10);
        $this->check($sizes === array_fill(0, 40, 1), 'every pop took one job');

        $this->line("\nb) prefetch 'auto': a backlog of 200 jobs of 10 ms, one worker");
        $sizes = $this->drain('queen-auto', 'backlog', 200, 10);
        $this->check(
            array_slice($sizes, 0, 8) === [1, 1, 2, 2, 4, 4, 8, 8],
            'the batch doubled after every two full pops: 1 1 2 2 4 4 8 8',
        );
        $grewAtMostDouble = true;
        for ($i = 1; $i < count($sizes); $i++) {
            $grewAtMostDouble = $grewAtMostDouble && $sizes[$i] <= 2 * $sizes[$i - 1];
        }
        $this->check(
            $grewAtMostDouble && max($sizes) <= 16,
            'no batch was more than twice the one before it, or above 16',
        );
        $this->check(
            count($sizes) * 3 <= 200,
            sprintf('a pop served %.1f jobs on average, 3 or more', 200 / count($sizes)),
        );

        $this->line("\nc) prefetch 'auto': 10 jobs of 300 ms, one worker");
        $sizes = $this->drain('queen-auto', 'long', 10, 300);
        $this->check(
            $sizes === array_fill(0, 10, 1),
            'every pop took one job: one of them is more than 250 ms of work',
        );
    }

    /**
     * Queue $jobs jobs of $millis ms on a fresh queue, then let one worker
     * run them all, and print what its pops took.
     *
     * @return list<int> the number of jobs of each pop, in order
     */
    private function drain(string $connection, string $name, int $jobs, int $millis): array
    {
        $queue = $this->freshQueue("prefetch-{$name}");
        // The whole backlog first, in bulk requests of up to 100 jobs.
        $batch = [];
        for ($seq = 1; $seq <= $jobs; $seq++) {
            $batch[] = new ProbeBatch($seq, $millis);
        }
        Queue::connection($connection)->bulk($batch, '', $queue);

        $this->startWorkers(1, $connection, $queue);
        $ran = $this->waitForEvents($queue, 'ran', $jobs, 60);
        $this->stopProcesses();

        // One lease per pop: group the jobs by lease, in the order they ran.
        $pops = [];
        $gaps = [];
        foreach ($ran as $i => $job) {
            if (isset($pops[$job['lease']]) && $ran[$i - 1]['lease'] === $job['lease']) {
                $gaps[] = $job['at'] - $ran[$i - 1]['at'];
            }
            $pops[$job['lease']] = ($pops[$job['lease']] ?? 0) + 1;
        }
        $sizes = array_values($pops);

        $this->line(sprintf('  jobs           %d', count($ran)));
        $this->line(sprintf('  pops           %d', count($sizes)));
        $this->line(sprintf('  jobs per pop   %.1f', count($ran) / count($sizes)));
        $sizesText = wordwrap($this->runLengths($sizes), 54, "\n" . str_repeat(' ', 17));
        $this->line("  batch sizes    {$sizesText}");
        if ($gaps !== []) {
            sort($gaps);
            $each = $gaps[intdiv(count($gaps), 2)] * 1000;
            $this->line(sprintf(
                '  one job took about %d ms, its ACK included: 250 ms holds %d',
                $each,
                intdiv(250, max(1, (int) $each)),
            ));
        }

        return $sizes;
    }
}

Phase a uses the queen connection, phases b and c queen-auto (prefetch 'auto' with lease_renewal), both in config/queue.php.

How it works

Counting pops

The broker leases everything one pop returns under one lease id, even when the jobs come from several partitions. The job reads that id from $this->job->getQueenMessage()['leaseId'], so the jobs that share an id came in the same pop, and counting the ids counts the pops.

What ‘auto’ does

Laravel’s worker calls pop() for every job. With a prefetch above 1, the driver answers from the batch it holds and asks the broker only when the batch is empty. 'auto' decides how many jobs that request asks for:

  • It times each job. From the moment pop() hands the job to Laravel to the worker’s next pop(), so the ACK is included. It keeps a moving average per queue, in which the newest job weighs 30%.
  • It aims at 250 ms of work. The next pop asks for 250 ms divided by the average, from 1 to 16 jobs.
  • It grows slowly. The batch doubles only after two pops in a row came back full, which is why phase b went 1 1 2 2 4 4 8 8. One full pop can be a short burst on a quiet queue.
  • It shrinks at once. When the average rises, the next pop asks for less, as the fall from 16 to 14 and 12 in phase b shows. A pop that comes back short means the backlog is gone, so the next pop asks for what came back, and for one job after an empty pop.

With prefetch 1 none of this runs: every pop() asks the broker for one job.

Why ‘auto’ needs lease renewal

The jobs of a batch wait, leased, while the jobs before them run, and Laravel can pause a worker that holds them (maintenance mode, queue:pause). Lease renewal keeps their lease alive, so the connection refuses prefetch 'auto' or above 1 without it (refused at start).

Limits

  • A crash returns the rest of the batch. Jobs that never started come back; whether they are charged an attempt depends on who hands them back. Keep tries at 2 or more (prefetch, ack_async and pop_ahead).
  • A slow job delays its batch. The jobs behind it in the batch wait for it, though another worker is free. That is why 'auto' keeps long jobs at one per pop.
  • The example runs one worker. With several, each worker sizes its own batches, per queue, from the jobs it ran itself.

Next

The balanced profile is prefetch 'auto' in a real configuration, with block_for and sleep set to match.

Navigation

Type to search…

↑↓ navigate↵ selectEsc close