Изградете отпорен PHP 8.3 CLI работник за продукциски оптоварувања
Работник за редица е едноставен сè до првиот незгоден пад. Во вообичаениот тек се презема ред и се повикува обработувач. Во продукција се појавуваат конкурентни процеси, отровни пораки, неизвесни потврдувања, истечени резервации, распоредувања и сигнали за прекин што пристигнуваат токму на погрешната граница.
Овој туторијал гради изворен PHP 8.3 работник со PDO и MySQL 8. Користи FOR UPDATE SKIP LOCKED за конкурентни преземања, ги потврдува резервациите пред обработка, враќа истечени закупи, применува ограничени повторни обиди, ги преместува крајните неуспеси во мртва редица и се исклучува преку pcntl.
Неговата гаранција за испорака е најмалку еднаш. Обработувач може повторно да се изврши по пад или истекување на закупот. Затоа, сигурноста зависи од заштита на деловниот ефект со стабилен клуч за идемпотентност таму каде што тој ефект се потврдува.
Предуслови и архитектонски граници
Потребни ви се PHP 8.3 CLI со pdo_mysql, mysqlnd и pcntl; MySQL 8 што користи InnoDB; и Linux-домаќин за долготрајниот процес.
php --version
php -m | grep -E '^(PDO|pdo_mysql|mysqlnd|pcntl)$'
mysql --version
Работникот набљудува три граници:
- Преземете еден подобен ред во кратка трансакција, доделете случаен токен за закуп и потврдете.
- Извршете ја деловната логика без да држите заклучувања на редови.
- Во друга кратка трансакција, потврдете сопственост, запишете го идемпотентниот деловен ефект и означете ја задачата како завршена.
Закупот е повратливо барање за сопственост, а не трансакција што се држи во текот на целото извршување. Ако процес исчезне, процесот за враќање ја враќа истечената резервација во редицата или ја означува како мртва кога не остануваат обиди.
Структура на проектот
report-worker/
├── database/
│ └── schema.sql
├── src/
│ └── config.php
└── bin/
├── enqueue.php
└── worker.php
Креирајте ги табелите за редицата и ефектите
Применете database/schema.sql со административен MySQL идентитет. Овој пример го задржува MySQL на истиот домаќин; приспособете го домаќинот на сметката само кога распоредувате преку приватна мрежа.
CREATE DATABASE report_queue
CHARACTER SET utf8mb4
COLLATE utf8mb4_0900_ai_ci;
USE report_queue;
CREATE TABLE jobs (
id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT,
queue_name VARCHAR(64) NOT NULL,
job_type VARCHAR(100) NOT NULL,
payload JSON NOT NULL,
idempotency_key CHAR(64) NOT NULL,
status ENUM('ready', 'reserved', 'done', 'dead')
NOT NULL DEFAULT 'ready',
attempts SMALLINT UNSIGNED NOT NULL DEFAULT 0,
max_attempts SMALLINT UNSIGNED NOT NULL DEFAULT 5,
available_at DATETIME(6) NOT NULL DEFAULT CURRENT_TIMESTAMP(6),
lease_until DATETIME(6) NULL,
lease_token BINARY(16) NULL,
last_error VARCHAR(4000) NULL,
created_at DATETIME(6) NOT NULL DEFAULT CURRENT_TIMESTAMP(6),
updated_at DATETIME(6) NOT NULL DEFAULT CURRENT_TIMESTAMP(6)
ON UPDATE CURRENT_TIMESTAMP(6),
PRIMARY KEY (id),
UNIQUE KEY uq_job_idempotency (queue_name, idempotency_key),
KEY ix_claim (queue_name, status, available_at, id),
KEY ix_recovery (queue_name, status, lease_until, id)
) ENGINE=InnoDB;
CREATE TABLE generated_reports (
id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT,
idempotency_key CHAR(64) NOT NULL,
tenant_id VARCHAR(100) NOT NULL,
report_month CHAR(7) NOT NULL,
result_json JSON NOT NULL,
created_at DATETIME(6) NOT NULL DEFAULT CURRENT_TIMESTAMP(6),
PRIMARY KEY (id),
UNIQUE KEY uq_report_idempotency (idempotency_key)
) ENGINE=InnoDB;
CREATE USER 'queue_worker'@'127.0.0.1'
IDENTIFIED BY 'replace-with-a-generated-secret';
GRANT SELECT, INSERT, UPDATE
ON report_queue.*
TO 'queue_worker'@'127.0.0.1';
Првиот уникатен клуч ги потиснува дупликатите барања за додавање во редица. Вториот го заштитува вистинскиот ефект во базата на податоци ако обработката се повтори. Идентитетот за извршување нема привилегии за управување со шемата или бришење.
Ограничете ги врските, мрежните читања и чекањата за заклучување
mysqlnd.net_read_timeout е PHP поставка на системско ниво, затоа конфигурирајте ја пред PHP да стартува. Најпрво лоцирајте ја CLI-конфигурацијата со php --ini. Овој пример за Debian/Ubuntu додава посебно CLI-презапишување; на друга дистрибуција, додајте ја истата директива во вчитаниот CLI php.ini или во неговиот скениран конфигурациски директориум.
php --ini
sudo install -d -m 0755 /etc/php/8.3/cli/conf.d
printf '%s\n' 'mysqlnd.net_read_timeout=10' \
| sudo tee /etc/php/8.3/cli/conf.d/99-report-worker.ini >/dev/null
php -r 'echo ini_get("mysqlnd.net_read_timeout"), PHP_EOL;'
Последната команда мора да испечати 10. Креирајте src/config.php. Временските ограничувања имаат различни намени и намерно се под 60-секундниот закуп.
<?php
declare(strict_types=1);
function database(): PDO
{
static $pdo = null;
if ($pdo instanceof PDO) {
return $pdo;
}
if (!extension_loaded('mysqlnd')) {
throw new RuntimeException('mysqlnd is required');
}
if ((int) ini_get('mysqlnd.net_read_timeout') !== 10) {
throw new RuntimeException(
'Set mysqlnd.net_read_timeout=10 in the CLI php.ini'
);
}
$password = getenv('DB_PASS');
if ($password === false || $password === '') {
throw new RuntimeException('DB_PASS is required');
}
$host = getenv('DB_HOST') ?: '127.0.0.1';
$port = getenv('DB_PORT') ?: '3306';
$name = getenv('DB_NAME') ?: 'report_queue';
$user = getenv('DB_USER') ?: 'queue_worker';
$dsn = "mysql:host={$host};port={$port};dbname={$name};charset=utf8mb4";
$pdo = new PDO($dsn, $user, $password, [
PDO::ATTR_ERRMODE => PDO::ERRMODE_EXCEPTION,
PDO::ATTR_DEFAULT_FETCH_MODE => PDO::FETCH_ASSOC,
PDO::ATTR_EMULATE_PREPARES => false,
PDO::ATTR_TIMEOUT => 5,
PDO::MYSQL_ATTR_MULTI_STATEMENTS => false,
]);
$pdo->exec("SET SESSION time_zone = '+00:00'");
$pdo->exec("SET SESSION innodb_lock_wait_timeout = 3");
$pdo->exec("SET SESSION lock_wait_timeout = 5");
return $pdo;
}
PDO::ATTR_TIMEOUT го ограничува воспоставувањето врска за PDO MySQL; не е краен рок за барање. mysqlnd.net_read_timeout ограничува поединечно мрежно читање на десет секунди. Сепак, не ја откажува работата на серверската страна, па работникот излегува по неочекуван исклучок од базата на податоци и дозволува неговиот надзорник да создаде чиста врска.
innodb_lock_wait_timeout ги ограничува чекањата за заклучување записи во InnoDB, додека lock_wait_timeout покрива заклучувања на метаподатоци. Овие поставки не го претвораат барањето за преземање со заклучување SELECT ... FOR UPDATE во строг краен рок за барање. Клиентската поставка mysqlnd.net_read_timeout спречува поединечно мрежно читање да чека засекогаш; по неочекуван истек на базата на податоци, работникот излегува, а неговиот надзорник создава чиста врска. Секое барање за преземање нека биде индексирано, а секоја трансакција за резервација кратка.
Додавајте задачи во редица идемпотентно
Креирајте bin/enqueue.php. Неговиот клуч доаѓа од деловниот идентитет на месечен извештај, така што повторувањето на командата го враќа постојниот ID на задачата.
<?php
declare(strict_types=1);
require dirname(__DIR__) . '/src/config.php';
if ($argc !== 3) {
fwrite(STDERR, "Usage: php bin/enqueue.php TENANT_ID YYYY-MM\n");
exit(64);
}
$tenant = $argv[1];
$month = $argv[2];
if (!preg_match('/^[A-Za-z0-9_-]{1,100}$/D', $tenant)) {
throw new InvalidArgumentException('Invalid tenant ID');
}
if (!preg_match('/^\d{4}-(0[1-9]|1[0-2])$/D', $month)) {
throw new InvalidArgumentException('Month must use YYYY-MM');
}
$key = hash('sha256', "monthly-report:{$tenant}:{$month}");
$payload = json_encode([
'tenant_id' => $tenant,
'report_month' => $month,
], JSON_THROW_ON_ERROR);
$sql = <<<'SQL'
INSERT INTO jobs (queue_name, job_type, payload, idempotency_key)
VALUES ('reports', 'report.generate', ?, ?)
ON DUPLICATE KEY UPDATE id = LAST_INSERT_ID(id)
SQL;
$pdo = database();
$statement = $pdo->prepare($sql);
$statement->execute([$payload, $key]);
fwrite(STDOUT, $pdo->lastInsertId() . PHP_EOL);
Имплементирајте преземања, повторни обиди и уредно ослободување
Креирајте bin/worker.php. Враќањето обработува најмногу 100 редови по премин. Индексираното преземање користи SKIP LOCKED и кратка трансакција. PDO MySQL нема универзален рок по исказ, па работникот комбинира серверски ограничувања на чекањето за заклучување со ограничено клиентско временско ограничување за мрежно читање и рестартирање од надзорник.
<?php
declare(strict_types=1);
require dirname(__DIR__) . '/src/config.php';
const QUEUE = 'reports';
const LEASE_SECONDS = 60;
function recoverExpiredLeases(PDO $pdo): int
{
$sql = <<<'SQL'
UPDATE jobs
SET status = CASE
WHEN attempts >= max_attempts THEN 'dead'
ELSE 'ready'
END,
available_at = CASE
WHEN attempts >= max_attempts THEN available_at
ELSE UTC_TIMESTAMP(6)
END,
lease_until = NULL,
lease_token = NULL,
last_error = CASE
WHEN attempts >= max_attempts
THEN 'Lease expired after final attempt'
ELSE 'Lease expired; returned to queue'
END
WHERE queue_name = ?
AND status = 'reserved'
AND lease_until <= UTC_TIMESTAMP(6)
ORDER BY lease_until, id
LIMIT 100
SQL;
$statement = $pdo->prepare($sql);
$statement->execute([QUEUE]);
return $statement->rowCount();
}
function claim(PDO $pdo): ?array
{
$pdo->beginTransaction();
try {
$select = $pdo->prepare(
"SELECT *
FROM jobs
WHERE queue_name = ?
AND status = 'ready'
AND attempts < max_attempts
AND available_at <= UTC_TIMESTAMP(6)
ORDER BY available_at, id
LIMIT 1
FOR UPDATE SKIP LOCKED"
);
$select->execute([QUEUE]);
$job = $select->fetch();
if ($job === false) {
$pdo->commit();
return null;
}
$token = random_bytes(16);
$update = $pdo->prepare(
"UPDATE jobs
SET status = 'reserved',
attempts = attempts + 1,
lease_token = ?,
lease_until = TIMESTAMPADD(
SECOND, ?, UTC_TIMESTAMP(6)
)
WHERE id = ?"
);
$update->execute([$token, LEASE_SECONDS, $job['id']]);
$pdo->commit();
$job['attempts'] = (int) $job['attempts'] + 1;
$job['lease_token'] = $token;
$testPause = (int) (getenv('WORKER_TEST_POST_CLAIM_PAUSE_MS') ?: 0);
if ($testPause > 0 && $testPause <= 10_000) {
usleep($testPause * 1000);
}
return $job;
} catch (Throwable $error) {
if ($pdo->inTransaction()) {
$pdo->rollBack();
}
throw $error;
}
}
function releaseClaim(PDO $pdo, array $job): bool
{
$statement = $pdo->prepare(
"UPDATE jobs
SET status = 'ready',
attempts = IF(attempts > 0, attempts - 1, 0),
available_at = UTC_TIMESTAMP(6),
lease_until = NULL,
lease_token = NULL,
last_error = 'Released during shutdown before processing'
WHERE id = ?
AND status = 'reserved'
AND lease_token = ?"
);
$statement->execute([$job['id'], $job['lease_token']]);
return $statement->rowCount() === 1;
}
function processJob(array $job): string
{
if ($job['job_type'] !== 'report.generate') {
throw new DomainException('Unsupported job type');
}
$payload = json_decode(
$job['payload'],
true,
flags: JSON_THROW_ON_ERROR
);
$tenant = $payload['tenant_id'] ?? null;
$month = $payload['report_month'] ?? null;
if (
!is_string($tenant)
|| !preg_match('/^[A-Za-z0-9_-]{1,100}$/D', $tenant)
|| !is_string($month)
|| !preg_match('/^\d{4}-(0[1-9]|1[0-2])$/D', $month)
) {
throw new UnexpectedValueException('Malformed job payload');
}
if (getenv('WORKER_ENABLE_TEST_FIXTURES') === '1') {
if (!empty($payload['test_pause'])) {
$until = hrtime(true) + 30_000_000_000;
while (hrtime(true) < $until) {
sleep(1);
}
}
$transientUntil = (int) (
$payload['test_transient_until_attempt'] ?? 0
);
if (
$transientUntil > 0
&& (int) $job['attempts'] <= $transientUntil
) {
throw new RuntimeException(
'Synthetic transient test failure'
);
}
}
return json_encode([
'tenant_id' => $tenant,
'report_month' => $month,
'state' => 'generated',
], JSON_THROW_ON_ERROR);
}
function ownsLease(array|false $row, string $token): bool
{
return $row !== false
&& (int) $row['lease_valid'] === 1
&& hash_equals($row['lease_token'], $token);
}
function complete(PDO $pdo, array $job, string $result): bool
{
$pdo->beginTransaction();
try {
$lock = $pdo->prepare(
"SELECT lease_token,
lease_until > UTC_TIMESTAMP(6) AS lease_valid
FROM jobs
WHERE id = ? AND status = 'reserved'
FOR UPDATE"
);
$lock->execute([$job['id']]);
$current = $lock->fetch();
if (!ownsLease($current, $job['lease_token'])) {
$pdo->rollBack();
return false;
}
$payload = json_decode(
$job['payload'],
true,
flags: JSON_THROW_ON_ERROR
);
$effect = $pdo->prepare(
"INSERT INTO generated_reports (
idempotency_key, tenant_id, report_month, result_json
) VALUES (?, ?, ?, ?)
ON DUPLICATE KEY UPDATE idempotency_key = ?"
);
$effect->execute([
$job['idempotency_key'],
$payload['tenant_id'],
$payload['report_month'],
$result,
$job['idempotency_key'],
]);
$finish = $pdo->prepare(
"UPDATE jobs
SET status = 'done',
lease_until = NULL,
lease_token = NULL,
last_error = NULL
WHERE id = ?"
);
$finish->execute([$job['id']]);
$pdo->commit();
return true;
} catch (Throwable $error) {
if ($pdo->inTransaction()) {
$pdo->rollBack();
}
throw $error;
}
}
function errorSummary(Throwable $error): string
{
$controlled = $error instanceof DomainException
|| $error instanceof UnexpectedValueException
|| $error instanceof JsonException;
return $controlled
? substr($error::class . ': ' . $error->getMessage(), 0, 4000)
: $error::class;
}
function failJob(PDO $pdo, array $job, Throwable $error): bool
{
$pdo->beginTransaction();
try {
$lock = $pdo->prepare(
"SELECT attempts, max_attempts, lease_token,
lease_until > UTC_TIMESTAMP(6) AS lease_valid
FROM jobs
WHERE id = ? AND status = 'reserved'
FOR UPDATE"
);
$lock->execute([$job['id']]);
$current = $lock->fetch();
if (!ownsLease($current, $job['lease_token'])) {
$pdo->rollBack();
return false;
}
$permanent = $error instanceof DomainException
|| $error instanceof UnexpectedValueException
|| $error instanceof JsonException;
$exhausted = $permanent
|| (int) $current['attempts']
>= (int) $current['max_attempts'];
$summary = errorSummary($error);
if ($exhausted) {
$statement = $pdo->prepare(
"UPDATE jobs
SET status = 'dead',
lease_until = NULL,
lease_token = NULL,
last_error = ?
WHERE id = ?"
);
$statement->execute([$summary, $job['id']]);
} else {
$attempt = (int) $current['attempts'];
$delay = min(
300,
(2 ** min($attempt, 8)) + random_int(0, 3)
);
$statement = $pdo->prepare(
"UPDATE jobs
SET status = 'ready',
available_at = TIMESTAMPADD(
SECOND, ?, UTC_TIMESTAMP(6)
),
lease_until = NULL,
lease_token = NULL,
last_error = ?
WHERE id = ?"
);
$statement->execute([$delay, $summary, $job['id']]);
}
$pdo->commit();
return true;
} catch (Throwable $failure) {
if ($pdo->inTransaction()) {
$pdo->rollBack();
}
throw $failure;
}
}
function logEvent(string $event, array $context = []): void
{
fwrite(STDERR, json_encode([
'event' => $event,
'time' => gmdate(DATE_ATOM),
] + $context, JSON_THROW_ON_ERROR) . PHP_EOL);
}
$stopping = false;
pcntl_async_signals(true);
$stop = static function (int $signal) use (&$stopping): void {
$stopping = true;
};
pcntl_signal(SIGTERM, $stop);
pcntl_signal(SIGINT, $stop);
try {
$pdo = database();
$nextRecoveryAt = 0.0;
logEvent('worker_started', ['pid' => getmypid()]);
while (!$stopping) {
if (microtime(true) >= $nextRecoveryAt) {
$recovered = recoverExpiredLeases($pdo);
$nextRecoveryAt = microtime(true) + 5.0;
if ($recovered > 0) {
logEvent('leases_recovered', ['count' => $recovered]);
}
}
if ($stopping) {
break;
}
$job = claim($pdo);
/*
* This check must be the first action after claim().
* No handler starts if shutdown arrived during reservation.
*/
if ($stopping) {
if ($job !== null) {
try {
$released = releaseClaim($pdo, $job);
logEvent(
$released
? 'claim_released_on_shutdown'
: 'claim_release_lost',
['job_id' => (int) $job['id']]
);
} catch (Throwable $releaseError) {
logEvent('claim_release_failed', [
'job_id' => (int) $job['id'],
'error' => errorSummary($releaseError),
]);
throw $releaseError;
}
}
break;
}
if ($job === null) {
for ($tick = 0; $tick < 10 && !$stopping; $tick++) {
usleep(250_000);
}
continue;
}
$startedAt = hrtime(true);
logEvent('job_claimed', [
'job_id' => (int) $job['id'],
'attempt' => $job['attempts'],
]);
try {
$result = processJob($job);
} catch (Throwable $handlerError) {
$owned = failJob($pdo, $job, $handlerError);
logEvent(
$owned ? 'job_failed' : 'job_lease_lost',
[
'job_id' => (int) $job['id'],
'error' => errorSummary($handlerError),
'elapsed_ms' => (int) (
(hrtime(true) - $startedAt) / 1_000_000
),
]
);
continue;
}
/*
* Do not classify finalization failures as handler failures.
* A database exception here escapes to the outer fatal handler,
* so systemd restarts the worker with a fresh connection.
*/
$owned = complete($pdo, $job, $result);
logEvent(
$owned ? 'job_completed' : 'job_lease_lost',
[
'job_id' => (int) $job['id'],
'elapsed_ms' => (int) (
(hrtime(true) - $startedAt) / 1_000_000
),
]
);
}
logEvent('worker_stopped', ['pid' => getmypid()]);
} catch (Throwable $fatal) {
logEvent('worker_fatal', ['error' => errorSummary($fatal)]);
exit(1);
}
SKIP LOCKED спречува работниците да чекаат зад редот избран од друг работник, но трансакцијата останува суштинска: редот останува заклучен сè додека не се зачуваат неговиот токен за закуп и рок.
Непосредната проверка за запирање затвора суптилна трка при исклучување. Ако SIGTERM пристигне додека claim() е блокирана или завршува, работникот условно го ажурира само редот што одговара и на ID-то на задачата и на случајниот токен за закуп. Тој ред го враќа во ready, го враќа неискористениот обид, бележи дали ослободувањето успеало и никогаш не повикува processJob().
Деловниот ефект и преминот во done делат една трансакција. За надворешен API, испратете го стабилниот клуч за идемпотентност кога е поддржан. MySQL не може атомски да координира неповрзан далечински несакан ефект; во спротивно користете усогласување надолу по текот или интеграција заснована на outbox.
Само исклучоците фрлени од processJob() влегуваат во failJob() и трошат обид. Неуспех во complete() или во трансакцијата за состојба на неуспех е инфраструктурен неуспех: стигнува до надворешниот обработувач на фатални грешки, го прекинува процесот и му дозволува на systemd да рестартира со свежа врска. Потврдената резервација потоа спречува непосредна дупликатна сопственост и останува повратлива по истекот.
Испробајте конкурентност и патеки на неуспех
Извезете ги ингеренциите во секој терминал, додадете го истиот извештај во редицата двапати и стартувајте два работника:
export DB_HOST=127.0.0.1
export DB_PORT=3306
export DB_NAME=report_queue
export DB_USER=queue_worker
export DB_PASS='replace-with-a-generated-secret'
php bin/enqueue.php acme 2042-01
php bin/enqueue.php acme 2042-01
php bin/worker.php
Командите за додавање во редица треба да го испечатат истиот ID. Додајте неколку закупци и потврдете дека одделни работници преземаат различни редови. Проверете ја трајната состојба со:
SELECT id, status, attempts, available_at, lease_until, last_error
FROM jobs
ORDER BY id DESC
LIMIT 20;
SELECT idempotency_key, tenant_id, report_month, result_json
FROM generated_reports
ORDER BY id DESC
LIMIT 20;
Тестирајте ја патеката за исклучување по преземањето на развојна инстанца. Паузата, која стандардно е оневозможена, го прави тајмингот детерминистички:
php bin/enqueue.php shutdown-test 2042-02
WORKER_TEST_POST_CLAIM_PAUSE_MS=10000 php bin/worker.php &
worker_pid=$!
case "$worker_pid" in
''|*[!0-9]*) exit 1 ;;
esac
sleep 1
kill -TERM "$worker_pid"
wait "$worker_pid"
Дневникот треба да содржи claim_released_on_shutdown, по што следи worker_stopped. Задачата треба да биде ready, без генериран извештај и без потрошен обид.
Извршете ги преостанатите фикстури за неуспех само врз изолирана развојна база на податоци. Однесувањето на обработувачот само за тестирање е оневозможено освен ако WORKER_ENABLE_TEST_FIXTURES=1.
Пад и враќање на истечен закуп
pause_id=$(
MYSQL_PWD="$DB_PASS" mysql --protocol=TCP -h "$DB_HOST" -P "$DB_PORT" -u "$DB_USER" "$DB_NAME" -Nse "
INSERT INTO jobs (
queue_name, job_type, payload, idempotency_key, max_attempts, available_at
) VALUES (
'reports',
'report.generate',
JSON_OBJECT(
'tenant_id', 'fixture-pause',
'report_month', '2042-03',
'test_pause', TRUE
),
SHA2(CONCAT('fixture:pause:', UUID()), 256),
3,
UTC_TIMESTAMP(6)
);
SELECT LAST_INSERT_ID();
"
)
WORKER_ENABLE_TEST_FIXTURES=1 php bin/worker.php 2> /tmp/report-worker-pause.log &
worker_pid=$!
for attempt in $(seq 1 50); do
grep -q '"event":"job_claimed"' /tmp/report-worker-pause.log && break
sleep 0.1
done
grep -q '"event":"job_claimed"' /tmp/report-worker-pause.log
kill -KILL "$worker_pid"
wait "$worker_pid" || true
sleep 61
timeout --signal=TERM 8s php bin/worker.php 2>> /tmp/report-worker-pause.log || test "$?" -eq 124
grep -q '"event":"leases_recovered"' /tmp/report-worker-pause.log
MYSQL_PWD="$DB_PASS" mysql --protocol=TCP -h "$DB_HOST" -P "$DB_PORT" -u "$DB_USER" "$DB_NAME" -e "SELECT id, status, attempts FROM jobs WHERE id = $pause_id"
Конечното барање мора да прикаже done. Дневникот мора да содржи и job_claimed и leases_recovered, докажувајќи дека втор работник повторно го презел истечениот 60-секунден закуп.
Трајни и привремени неуспеси
unsupported_id=$(
MYSQL_PWD="$DB_PASS" mysql --protocol=TCP -h "$DB_HOST" -P "$DB_PORT" -u "$DB_USER" "$DB_NAME" -Nse "
INSERT INTO jobs (
queue_name, job_type, payload, idempotency_key, max_attempts, available_at
) VALUES (
'reports',
'unsupported.fixture',
JSON_OBJECT(
'tenant_id', 'fixture-unsupported',
'report_month', '2042-04'
),
SHA2(CONCAT('fixture:unsupported:', UUID()), 256),
1,
UTC_TIMESTAMP(6)
);
SELECT LAST_INSERT_ID();
"
)
timeout --signal=TERM 5s php bin/worker.php 2> /tmp/report-worker-permanent.log || test "$?" -eq 124
MYSQL_PWD="$DB_PASS" mysql --protocol=TCP -h "$DB_HOST" -P "$DB_PORT" -u "$DB_USER" "$DB_NAME" -e "SELECT id, status, attempts FROM jobs WHERE id = $unsupported_id"
transient_id=$(
MYSQL_PWD="$DB_PASS" mysql --protocol=TCP -h "$DB_HOST" -P "$DB_PORT" -u "$DB_USER" "$DB_NAME" -Nse "
INSERT INTO jobs (
queue_name, job_type, payload, idempotency_key, max_attempts, available_at
) VALUES (
'reports',
'report.generate',
JSON_OBJECT(
'tenant_id', 'fixture-transient',
'report_month', '2042-05',
'test_transient_until_attempt', 2
),
SHA2(CONCAT('fixture:transient:', UUID()), 256),
4,
UTC_TIMESTAMP(6)
);
SELECT LAST_INSERT_ID();
"
)
timeout --signal=TERM 30s env WORKER_ENABLE_TEST_FIXTURES=1 php bin/worker.php 2> /tmp/report-worker-transient.log || test "$?" -eq 124
MYSQL_PWD="$DB_PASS" mysql --protocol=TCP -h "$DB_HOST" -P "$DB_PORT" -u "$DB_USER" "$DB_NAME" -e "SELECT id, status, attempts FROM jobs WHERE id = $transient_id"
grep -q '"event":"job_failed"' /tmp/report-worker-transient.log
grep -q '"event":"job_completed"' /tmp/report-worker-transient.log
Неподдржаната фикстура мора да биде dead по еден обид. Привремената фикстура мора да достигне done во третиот обид, по два синтетички неуспеси и ограничен backoff со случајно отстапување.
Распоредете со systemd
Инсталирајте ги прегледаните датотеки под /opt/report-worker, во сопственост на root и читливи за посебната сметка report-worker. Чувајте ги доделувањата на околински променливи во /etc/report-worker.env, во сопственост на root со режим 0600. Ова се операции на администратор на домаќинот, а не команди за контејнер.
Креирајте /etc/systemd/system/report-worker.service со sudoedit и користете ја токму оваа единица:
[Unit]
Description=PHP report queue worker
After=network-online.target
Wants=network-online.target
[Service]
Type=simple
User=report-worker
Group=report-worker
WorkingDirectory=/opt/report-worker
EnvironmentFile=/etc/report-worker.env
ExecStart=/usr/bin/php /opt/report-worker/bin/worker.php
Restart=on-failure
RestartSec=3
KillSignal=SIGTERM
TimeoutStopSec=75
NoNewPrivileges=true
PrivateTmp=true
PrivateDevices=true
ProtectSystem=strict
ProtectHome=true
ProtectKernelTunables=true
ProtectKernelModules=true
ProtectControlGroups=true
LockPersonality=true
RestrictAddressFamilies=AF_UNIX AF_INET AF_INET6
UMask=0077
[Install]
WantedBy=multi-user.target
Потврдете и овозможете ја точната датотека на единицата:
sudo systemd-analyze verify /etc/systemd/system/report-worker.service
sudo systemctl daemon-reload
sudo systemctl enable --now report-worker.service
sudo systemctl status report-worker.service
sudo journalctl -u report-worker.service --since today
Користете systemd-шаблон за дополнителни инстанци. Скалирајте постепено: работниците трошат врски со базата на податоци и капацитет надолу по текот дури и кога конкуренцијата за заклучување редови е мала. Буџетот за запирање од 75 секунди го надминува закупот, додека временските ограничувања за врска, мрежно читање и заклучување остануваат под двете вредности.
Набљудливост, безбедност и перформанси
Работникот емитува структуриран JSON што содржи ID-ја на задачи, обиди, изминато време, броеви на враќања, загуби на закуп, ослободувања при исклучување и настани од животниот циклус. Поставете предупредувања за раст на мртви задачи, старост на најстарата подготвена задача, повторени враќања, стапка на повторни обиди, латентност на обработка, неуспеси при ослободување и фатални грешки во базата на податоци.
Чувајте ги ингеренциите и чувствителните полиња од payload надвор од дневниците и last_error. Потврдувајте ги payload-ите при потрошувачка дури и кога производителите се доверливи. За далечински MySQL, барајте проверка на сертификати, ограничете ја портата 3306 на овластени адреси на апликации и создајте сметка ограничена на таа мрежа наместо да користите јавен домаќин со џокер.
Мерете ги плановите на барањата за преземање додека табелата расте. Чувајте ги завршените задачи само онолку долго колку што бараат оперативните или ревизорските потреби и архивирајте ги преку посебен прегледан процес. Бавните мрежни повици, генерирањето извештаи и работата со датотечниот систем припаѓаат надвор од трансакциите за резервација и им требаат сопствени рокови под закупот.
Вообичаени продукциски неуспеси
- Работниците изгледаат сериски: обработката останува во трансакцијата за преземање или недостига индексот за преземање.
- Задачите се преклопуваат по истекот: обработката го надминува закупот. Ограничете го времето на операцијата, намерно приспособете го закупот или имплементирајте обновување оградено со токен.
- Барањата сè уште висат: конфигуриран е само
PDO::ATTR_TIMEOUT. Ограничувањата за врска, мрежно читање и чекање заклучување се одделни контроли; PDO MySQL нема универзално временско ограничување за барања. - Распоредувањата трошат обиди: недостига проверка за запирање по преземањето или ослободувањето не враќа неискористен обид.
- Застарен работник потврдува: финализацијата не ги проверува и токенот и неистечениот рок под заклучување на редот.
- Се појавуваат далечински дупликати: идемпотентноста постои само во MySQL, наместо на границата на надворешниот ефект.
Конечна контролна листа за проверка
- Дупликатното додавање во редица враќа еден траен ID на задача.
- Конкурентни работници преземаат различни подобни редови без долго чекање.
- Резервацијата се потврдува пред да започне деловната обработка.
- Ограничувањата за врска, мрежно читање и чекање заклучување се под закупот, а трансакциите за преземање остануваат кратки и индексирани.
- Сигнал забележан за време на
claim()го ослободува само соодветниот токен и не стартува обработувач. - Истечените закупи се враќаат во
readyдодека остануваат обиди. - Исцрпените и трајните неуспеси влегуваат во
dead. - Истечен или застарен токен не може да финализира или повторно да закаже задача.
- Табелата за деловен ефект одбива дупликатни клучеви за идемпотентност.
- Одложувањата за повторни обиди и максималниот број обиди се конечни.
- Дневниците изложуваат латентност, повторни обиди, враќање, губење закуп, ослободување и исклучување.
- Идентитетот за извршување нема непотребни привилегии за база на податоци, датотечен систем или мрежа.
Траен ред во редица е само почеток на отпорен работник. Важното инженерство е на границите: резервирајте кратко, извршувајте без заклучувања, оградете секое конечно запишување, ослободете ја работата кога исклучувањето ја добива трката и направете го вистинскиот ефект идемпотентен. Откако овие правила се експлицитни, падовите и повторената испорака престануваат да бидат изненадувања и стануваат вообичаени состојби што системот е изграден да ги апсорбира.