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/laravelWhat 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 checksexamples/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:
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:
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 nextpop(), 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
triesat 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.