Queue (SQL)¶
Note
Not part of core. Install it separately:
composer require kinetis/queue-sql
Adds MySQL or Postgres as a backend for Queue, storing jobs in a database you already run. Switching to it changes configuration, not application code.
QUEUE_CONNECTION=sql
DB_CONNECTION=mysql # or "pgsql"
DB_HOST=127.0.0.1
DB_NAME=app
DB_USER=app
DB_PASSWORD=secret
vendor/bin/kinetis queue:work --queue=high,default
pop() reserves a job with SELECT ... FOR UPDATE SKIP LOCKED, so the
server must be MySQL 8.0+, MariaDB 10.6+ or Postgres 9.5+. An older
server rejects the clause, and pop() fails.
Configuring¶
The DB_* keys are the ones kinetis/database-bridge reads (see
Database). QUEUE_VISIBILITY_TIMEOUT_SECONDS, described below,
is the one key this package adds.
Create the table¶
kinetis/queue-sql ships one Migrations stub per dialect:
vendor/kinetis/queue-sql/resources/migrations/create_kinetis_queue_jobs_table.mysql.php.stub
vendor/kinetis/queue-sql/resources/migrations/create_kinetis_queue_jobs_table.pgsql.php.stub
Copy the one matching your database into your migrations/ directory
with a timestamp prefix, then run vendor/bin/kinetis migrate.
SQL mechanisms describes the columns the backend
runs on.
Enqueueing inside a transaction¶
Queue’s QueueInterface::push() is not enlisted in a transaction
the caller has open: it runs its INSERT on the queue’s own connection,
so the job is enqueued whether or not the surrounding transaction
commits. Kinetis\QueueSql\SqlQueue::pushOn() puts the row on a
transaction you already hold, so the job and the writes it belongs to
land together:
use Kinetis\Persistence\Contract\SqlTransaction;
$this->transactions->transaction($this->db, function (SqlTransaction $tx) use ($orderId): void {
$tx->execute('UPDATE orders SET status = ? WHERE id = ?', ['paid', $orderId]);
$this->queue->pushOn($tx, new SendReceipt($orderId));
});
The row becomes visible and durable only if that transaction commits. A
throw before the commit rolls it back with the caller’s other work, and a
COMMIT that fails leaves the outcome unknown — see Database’s
“When a write’s outcome is unknown”. Push telemetry closes when the
INSERT statement completes, so the span reports the enqueue statement
rather than the later commit.
pushOn() runs that one statement and nothing else: it never commits,
rolls back, nests a transaction, or falls back to the connection the
queue was built with. Addressing the database that holds
kinetis_queue_jobs is therefore the caller’s job —
the queue’s connection name picks the connection behind push() and
does not redirect a transaction you supply.
The signature belongs to this package, not to QueueInterface, so a
caller needs the SqlQueue itself rather than the QueueInterface the
container binds. SqlQueueFactory::fromConfig() returns that class.
Build it once and register that one object under both ids:
use Kinetis\Queue\QueueInterface;
use Kinetis\QueueSql\SqlQueue;
use Kinetis\QueueSql\SqlQueueFactory;
$queue = SqlQueueFactory::fromConfig($config);
$app->instance(SqlQueue::class, $queue);
$app->instance(QueueInterface::class, $queue);
$app->onDispose($queue->dispose(...));
Ordinary QueueInterface consumers and pushOn() callers then share
one backend instance and its one connection pool. Binding only
SqlQueue::class leaves the default QueueInterface binding in place,
and it builds a second SqlQueue with a pool of its own.
The factory opened that connection, so the queue owns it and the
onDispose() line is what closes it when the worker ends — the same
registration the default QueueInterface binding makes for the backend
it builds. See Appendix: Queue Contracts’s “Connection ownership”, which
also covers a SqlQueue constructed directly around a link the
application already has.
Appendix: Queue Contracts’s “Multiple backends” shows the same registration
where several backends run side by side. pushOn() takes a raw
Database transaction; an ORM transaction session does
not expose its transaction. If ORM work must schedule a job atomically,
map an application outbox intent as an entity with a #[BelongsTo] to
what it depends on, persist both and flush once, then publish after the
transaction returns —
Locking rows, or entities and SQL in one
transaction
covers the composition.
Visibility timeout¶
QUEUE_VISIBILITY_TIMEOUT_SECONDS (default 300, at least 1) is how
long a reserved job belongs to its worker. When a worker dies before
settling a job, the job becomes poppable again once its reservation is
older than this, with its attempt count increased. Constructing
SqlQueue directly takes the same value as its visibilityTimeoutSeconds
argument.
queue:work renews the reservation automatically while the job runs, at
half this window, so the setting sizes how long a crashed worker’s row
waits to come back rather than how long a job may take. A handler that
never yields to the event loop cannot be renewed, and neither can one
whose worker has died. Should a reservation lapse anyway, the first
worker’s late settlement is rejected rather than touching the new
reservation, and the worker reports it as a lost settlement (see
Queue’s “When a settlement is lost”).
Reservation times are written and compared with each worker’s own clock, not the database’s — a renewal included — so keep worker clocks synchronized: skew makes a reservation expire early or late by the difference.
Clearing a queue¶
SqlQueue declares ClearableQueueInterface (see Queue’s
“Clearing is a separate capability”). Clearing deletes the queue’s
unreserved rows, delayed ones included, and reports how many. A
reservation past its timeout is left in place, since its worker may
still be running the job.
Delays and retries¶
A delayed job becomes available to the first pop() after its delay,
measured with the pushing and popping hosts’ clocks, so it runs late
while every worker is busy. Retries follow Queue: maxAttempts,
QUEUE_MAX_ATTEMPTS and QUEUE_RETRY_BASE_DELAY_SECONDS. A delayed
retry moves the row’s own available_at, so it waits exactly as a
delayed push does and needs no schema change.
Named connections¶
QUEUE_CONNECTION_NAME=reports
QUEUE_REPORTS_CONNECTION=sql
DB_REPORTS_CONNECTION=mysql
DB_REPORTS_HOST=127.0.0.1
A connection’s name scopes its selector, its DB_* keys and
QUEUE_VISIBILITY_TIMEOUT_SECONDS; default reads the plain keys.
QUEUE_CONNECTION_NAME makes reports the bound QueueInterface, and
queue:work --connection=reports runs it either way. See Queue’s
“Named connections” and Configuration.
See also¶
Queue — jobs, workers, retries and delivery guarantees.
Appendix: Queue Contracts — reservation fencing and delivery contracts.
Database — connecting to MySQL and Postgres.
Migrations — running the table migration.