Skip to content

Laravel jobs in order, one customer at a time

A Laravel job that names its partition with QueenPartitionable, so each customer's or user's jobs run one at a time and in dispatch order. Three customers' jobs, interleaved on four workers, show it, while different customers run in parallel.

Updated View as Markdown

A job that implements QueenPartitionable waits until the job before it of the same customer has ended, whichever worker takes it. Jobs of different customers still run at the same time. The program dispatches 5 jobs for each of 3 customers, runs them on 4 workers, and checks both.

Laravel’s own WithoutOverlapping middleware also keeps two jobs with the same key from running at once, but it does so by releasing the second job back to the queue. The release costs the job an attempt, and a later job of the same customer can run before it. A partition gives both: one job at a time, in dispatch order.

Run it

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

What it printed against a 2.0.1 node (setup):

broker http://localhost:6632

dispatched 15 jobs, then started 4 workers (php artisan queue:work)

  customer  run order   on workers       from    to
  acme      1 2 3 4 5   w1 w1 w4 w2 w1   0.0 s   1.6 s
  globex    1 2 3 4 5   w2 w4 w2 w1 w3   0.0 s   1.5 s
  initech   1 2 3 4 5   w3 w3 w3 w4 w2   0.0 s   1.6 s

  at most 3 jobs ran at the same time, on 4 workers

checking
  ok: acme's jobs ran in dispatch order
  ok: acme's jobs never ran at the same time
  ok: globex's jobs ran in dispatch order
  ok: globex's jobs never ran at the same time
  ok: initech's jobs ran in dispatch order
  ok: initech's jobs never ran at the same time
  ok: jobs of different customers ran in parallel
  ok: at most 3 jobs at once, one per customer, on 4 workers

PASS: 8 checks
  • The order held, the worker changed. acme’s five jobs ran on w1, w1, w4, w2 and w1. Order does not come from pinning a customer to a worker.
  • Three at a time, not four. Three customers are three partitions. While all three ran, the fourth worker had nothing it was allowed to take.
  • 1.6 seconds for 15 jobs of 300 ms. The three customers ran side by side; one after another they would have taken 4.5 seconds.

The program

The job names its partition:

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

    public function __construct(
        public string $customer,
        public int $seq,
    ) {}

    // The partition of this job: every job of one customer goes to
    // the same ordered lane, which the broker leases to one worker at
    // a time.
    public function queenPartition(): string
    {
        return 'customer:' . $this->customer;
    }

    public function handle(): void
    {
        $started = microtime(true);
        usleep(300_000); // the work: 300 ms
        Journal::record($this->job, [
            'event' => 'ran',
            'customer' => $this->customer,
            'seq' => $this->seq,
            'started' => $started,
            'ended' => microtime(true),
        ]);
    }
}

The command dispatches the jobs, starts the workers on the queen connection (one job per pop), and checks what the jobs recorded:

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

    protected $description = 'One customer at a time: per-entity ordering with QueenPartitionable';

    private const CUSTOMERS = ['acme', 'globex', 'initech'];

    private const JOBS_PER_CUSTOMER = 5;

    private const WORKERS = 4;

    protected function example(): void
    {
        $queue = $this->freshQueue('ordering');

        // Dispatched interleaved, the way events of several customers arrive:
        // acme 1, globex 1, initech 1, acme 2, ...
        for ($seq = 1; $seq <= self::JOBS_PER_CUSTOMER; $seq++) {
            foreach (self::CUSTOMERS as $customer) {
                RebuildCustomer::dispatch($customer, $seq)->onQueue($queue);
            }
        }
        $total = count(self::CUSTOMERS) * self::JOBS_PER_CUSTOMER;
        $workers = $this->startWorkers(self::WORKERS, 'queen', $queue);
        $this->line("\ndispatched {$total} jobs, then started "
            . count($workers) . ' workers (php artisan queue:work)');

        $ran = $this->waitForEvents($queue, 'ran', $total, 60);
        $this->stopProcesses();

        // What each customer saw: its jobs in the order they started, and the
        // worker (w1 to w4) that ran each one.
        $t0 = min(array_column($ran, 'started'));
        $name = array_flip($workers);
        $byCustomer = [];
        foreach ($ran as $job) {
            $byCustomer[$job['customer']][] = $job;
        }
        $row = '  %-9s %-11s %-16s %-7s %s';
        $this->line(sprintf("\n{$row}", 'customer', 'run order', 'on workers', 'from', 'to'));
        foreach (self::CUSTOMERS as $customer) {
            $jobs = $byCustomer[$customer] ?? [];
            usort($jobs, fn ($a, $b) => $a['started'] <=> $b['started']);
            $byCustomer[$customer] = $jobs;
            $this->line(sprintf(
                $row,
                $customer,
                implode(' ', array_column($jobs, 'seq')),
                implode(' ', array_map(fn ($job) => 'w' . ($name[$job['pid']] + 1), $jobs)),
                sprintf('%.1f s', $jobs[0]['started'] - $t0),
                sprintf('%.1f s', end($jobs)['ended'] - $t0),
            ));
        }
        $peak = $this->peakConcurrency($ran);
        $this->line("\n  at most {$peak} jobs ran at the same time, on "
            . self::WORKERS . ' workers');

        $this->line("\nchecking");
        foreach (self::CUSTOMERS as $customer) {
            $jobs = $byCustomer[$customer];
            $this->check(
                array_column($jobs, 'seq') === range(1, self::JOBS_PER_CUSTOMER),
                "{$customer}'s jobs ran in dispatch order",
            );
            $overlaps = 0;
            for ($i = 1; $i < count($jobs); $i++) {
                $overlaps += $jobs[$i]['started'] < $jobs[$i - 1]['ended'] ? 1 : 0;
            }
            $this->check($overlaps === 0, "{$customer}'s jobs never ran at the same time");
        }
        $this->check($peak >= 2, 'jobs of different customers ran in parallel');
        $this->check(
            $peak <= count(self::CUSTOMERS),
            "at most {$peak} jobs at once, one per customer, on " . self::WORKERS . ' workers',
        );
    }

    /** The most jobs that were running at one moment. */
    private function peakConcurrency(array $ran): int
    {
        $edges = [];
        foreach ($ran as $job) {
            $edges[] = [$job['started'], 1];
            $edges[] = [$job['ended'], -1];
        }
        // At equal times an end comes first: a job that starts as another
        // ends does not overlap it.
        usort($edges, fn ($a, $b) => [$a[0], $a[1]] <=> [$b[0], $b[1]]);
        $running = $peak = 0;
        foreach ($edges as [, $step]) {
            $running += $step;
            $peak = max($peak, $running);
        }

        return $peak;
    }
}

How it works

  1. The partition. dispatch() pushes each job to the partition its queenPartition() names, customer:acme here, created by the first push. A job without the interface goes to one of the queue’s 64 stripes instead, by a hash of its UUID (queues and stripes).
  2. One worker per partition. A pop takes jobs under a lease, and the broker leases a partition to one worker of the consumer group at a time. Until that worker acknowledges acme 1, no other worker receives acme 2 (the lease).
  3. The next job goes to whoever asks. The ACK releases the partition. The next pop of any worker can take acme 2: here the idle workers long-poll (block_for 1), so one of them takes it at once.

The partition is a property of the storage, so the order also holds across a crash. A worker that dies holding acme 2 holds its lease until retry_after ends; acme 3 waits, and acme 2 runs again first.

Limits

  • A slow job holds up its customer. acme 3 waits for acme 2, however long acme 2 takes. Other customers do not wait.
  • Parallelism is capped by the number of customers with work. Workers beyond that number stay idle on the queue: here one of the four was always idle.
  • Delivery is at least once. A job that runs again after a crash runs before the next job of its customer, so make jobs safe to repeat (at-least-once delivery).

Next

Order jobs per entity adds the interface to your own job, and partitions explains what a partition costs the broker.

Navigation

Type to search…

↑↓ navigate↵ selectEsc close