Skip to content

Repository files navigation

BackgroundQueue

Komponenta umožňuje zpracovávat úkoly na pozadí pomocí cronu nebo AMQP brokera (např. RabbitMQ). Vhodné pro dlouhotrvající requesty, komunikaci s API nebo odesílání webhooků či e-mailů.

Komponenta využívá vlastní doctrine entity manager pro ukládání záznamů do fronty. Tím pádem fungování komponenty není ovlivněno aplikačním entity managerem a naopak.

1. Instalace a konfigurace

1.1 Instalace

composer require adt/background-queue

1.2 Registrace a konfigurace

BackgroundQueue přebírá pole následujících parametrů:

$connection = [
	'serverVersion' => '8.0',
	'driver' => 'pdo_mysql',
	'host' => $_ENV['DB_HOST'],
	'port' => $_ENV['DB_PORT'],
	'user' => $_ENV['DB_USER'],
	'password' => $_ENV['DB_PASSWORD'],
	'dbname' => $_ENV['DB_DBNAME'],
];

$backgroundQueue = new \ADT\BackgroundQueue\BackgroundQueue([
	'callbacks' => [
		'processEmail' => [$mailer, 'process'],
		'processEmail2' => [ // možnost specifikace jiné fronty pro tento callback
			'callback' => [$mailer, 'process'],
			'queue' => 'general',
		],
	]
	'notifyOnNumberOfAttempts' => 5, // počet pokusů o zpracování záznamu před zalogováním
	'tempDir' => $tempDir, // cesta pro uložení zámku proti vícenásobnému spuštění commandu
	'connection' => $connection, // Doctrine\Dbal\Connection
	'queue' => $_ENV['PROJECT_NAME'], // název fronty, do které se ukládají a ze které se vybírají záznamy
	'bulkSize' => 1, // velikost dávky při vkládání více záznamů najednou
	'tableName' => 'background_job', // nepovinné, název tabulky, do které se budou ukládat jednotlivé joby
	'logger' => $logger, // nepovinné, musí implementovat psr/log LoggerInterface
	'onBeforeProcess' => function(array $parameters) {...}, // nepovinné
	'onError' => function(Throwable $e, array $parameters) {...},  // nepovinné
	'onAfterProcess' => function(array $parameters) {...}, // nepovinné
	'onProcessingGetMetadata' => function(array $parameters): ?array {...}, // nepovinné
	'stalledJobTimeout' => 3600, // nepovinné, po kolika sekundách bez zápisu se job ve stavu PROCESSING považuje za osiřelý (0 = vypnuto), viz 2.4
	'heartbeatInterval' => 60, // nepovinné, minimální prodleva v sekundách mezi dvěma zápisy heartbeat(), viz 2.4
	'lostMessageTimeout' => 3600, // nepovinné, po kolika sekundách bez zápisu se u jobu ve stavu READY/TEMPORARILY_FAILED považuje zpráva v brokeru za ztracenou a publikuje se znovu (0 = vypnuto), viz 2.5
	'parametersFormat' => \ADT\BackgroundQueue\Entity\BackgroundJob::PARAMETERS_FORMAT_SERIALIZE, // nepovinné, určuje v jakém formátu budou do DB ukládána data v `background_job.parameters` (@see \ADT\BackgroundQueue\Entity\BackgroundJob::setParameters),
]);

Potřebné databázové schéma se vytvoři při prvním použití fronty automaticky a také se automaticky aktualizuje, je-li třeba.

1.3 Broker (optional)

You can use this package with any message broker. Your producer or consumer just need to implement ADT\BackgroundQueue\Broker\Producer or ADT\BackgroundQueue\Broker\Consumer.

Or you can use php-amqplib/php-amqplib, for which this library has an ready to use implementation.

1.3.1 php-amqplib installation

Because using of php-amqplib/php-amqplib is optional, it doesn't check your installed version against the version with which this package was tested. That's why it's recommended to add to your composer:

{
  "conflict": {
    "php-amqplib/php-amqplib": "<3.0.0 || >=4.0.0"
  }
}

This version of php-amqplib/php-amqplib also need ext-sockets. You can add it to your Dockerfile like this:

docker-php-ext-install sockets

and then run:

composer require php-amqplib/php-amqplib

This makes sure you avoid BC break when upgrading php-amqplib/php-amqplib in the future.

1.3.1 php-amqplib configuration

$connectionParams = [
    'host' => $_ENV['RABBITMQ_HOST'],
    'user' => $_ENV['RABBITMQ_USER'],
    'password' => $_ENV['RABBITMQ_PASSWORD']
];
$queueParams = [
    'arguments' => ['x-queue-type' => ['S', 'quorum']]
];

// nepovinné: per-queue AMQP argumenty
// klíč = část názvu fronty, hodnota = argumenty aplikované na každou frontu, jejíž název daný řetězec obsahuje
// příklad: fronta "transcribe" poběží přes single active consumer (jen jeden aktivní consumer napříč všemi)
$queueArguments = [
    'transcribe' => ['x-single-active-consumer' => ['t', true]]
];

$manager = new \ADT\BackgroundQueue\Broker\PhpAmqpLib\Manager($connectionParams, $queueParams, $queueArguments);
$producer = new \ADT\BackgroundQueue\Broker\PhpAmqpLib\Producer();
$consumer = new \ADT\BackgroundQueue\Broker\PhpAmqpLib\Consumer();

$backgroundQueue = new \ADT\BackgroundQueue\BackgroundQueue([
	...
	'producer' => $producer,
	'waitingJobExpiration' => 1000, // nepovinné, délka v ms, po které se job pokusí znovu provést, když čeká na dokončení předchozího
]);

2 Použití

2.1 Přidání záznamu do fronty a jeho zpracování

namespace App\Presenters;

use ADT\BackgroundQueue\BackgroundQueue;

class Mailer
{
    private BackgroundQueue $backgroundQueue

    public function __construct(BackgroundQueue $backgroundQueue)
    {
        $this->backgroundQueue = $backgroundQueue;
    }

	public function send(Invoice $invoice) 
	{
		$callbackName = 'processEmail';
		$parameters = [
			'to' => 'hello@appsdevteam.com',
			'subject' => 'Background queue test'
			'text' => 'Anything you want.'
		];
		$serialGroup = 'invoice-' . $invoice->getId();
		$identifier = 'sendEmail-' . $invoice->getId();
		$isUnique = true; // always set to true if a callback on an entity should be performed only once, regardless of how it can happen that it is added to your queue twice
		$availableAt = new \DateTimeImmutable('+1 hour'); // earliest time when the record should be processed

		$this->backgroundQueue->publish($callbackName, $parameters, $serialGroup, $identifier, $isUnique, $availableAt);
	}
	
	public function process(string $to, string $subject, string $text) 
	{
	    // own implementation
	}
}

Záznam se uloží ve stavu READY.

Parametr $parameters může přijímat jakýkoliv běžný typ (pole, objekt, string, ...) či jejich kombinace (pole objektů), a to dokonce i binární data.

Parametr $serialGroup je nepovinný - jeho zadáním zajistítě, že všechny joby se stejným serialGroup budou provedeny sériově. Job, který ve své skupině narazí na dosud nedokončeného předchůdce, se odloží do stavu WAITING; jakmile předchůdce dojede, jeho konzument sám vrátí hlavu skupiny (nejvyšší priorita, při shodě nejnižší ID) na READY a znovu ji zařadí do brokera. background-queue:process totéž dělá jako záchrannou síť pro případ, že by k probuzení nedošlo (zabitý konzument, ztracená zpráva).

Parametr $identifier je nepovinný - pomocí něj si můžete označit joby vlastním identifikátorem a následně pomocí metody getUnfinishedJobIdentifiers(array $identifiers = []) zjistit, které z nich ještě nebyly provedeny. Za nedokončený se nepovažuje job ve stavu FINISHED, REDUNDANT ani PERMANENTLY_FAILED - trvale selhaný job už znovu nepoběží, takže by jinak zůstal „nedokončený" navždy.

Parametr $coalesceThreshold je nepovinný a slouží ke slučování (coalescingu) překrývajících se jobů ve stejné serialGroup. Vyžaduje, aby byl nastaven i $serialGroup (jinak se vyhodí výjimka) - ten určuje rozsah slučování. Když se job se zadaným prahem spustí, označí všechny ostatní dosud nezpracované joby téže serialGroup, jejichž práh je vyšší nebo stejný, za REDUNDANT - tedy ty, které svým během pokryje.

Typický příklad je úloha typu "přepočítej od X dál": job "přepočítej od dokladu 150" pohltí čekající joby "přepočítej od dokladu 151", "od 152" atd., protože je svým během stejně zpracuje. Rozhodnutí o redundanci dělá vždy ten "širší" job (s nižším prahem) v okamžiku svého spuštění - ví totiž jistě, co všechno přepočítá. Ruší se jen joby, které ještě nezačaly běžet; už běžící (PROCESSING) ani dokončené joby zůstanou nedotčené.

$this->backgroundQueue->publish(
	'recalculateStock',
	['itemId' => 1, 'fromDocument' => 150],
	serialGroup: 'stock-item-1', // rozsah slučování (např. jedna skladová položka)
	coalesceThreshold: 150,       // job pohltí všechny nedokončené joby téže skupiny s prahem >= 150
);

Pozor: pořadí zpracování v rámci serialGroup se řídí prioritou a ID (pořadím vložení), nikoli prahem. Coalescing tedy spolehlivě šetří práci, dokud joby s nižším prahem vznikají dříve (mají nižší ID). Pokud výjimečně vznikne dříve job s vyšším prahem, oba joby doběhnou - výsledek je korektní, jen bez úspory.

Pokud callback vyhodí ADT\BackgroundQueue\Exception\PermanentErrorException, záznam se uloží ve stavu PERMANENTLY_FAILED a je potřeba jej zpracovat ručně.

Pokud callback vyhodí ADT\BackgroundQueue\Exception\WaitingException, původní záznam se uzavře jako FINISHED a místo něj se publikuje jeho klon s odkladem waitingJobExpiration. Počítadlo pokusů se tedy nezvyšuje - job se zkusí znovu jako nový záznam. (Nepleťte si to se stavem WAITING, do kterého se odkládají joby čekající na předchůdce ve své serialGroup.)

Pokud callback vyhodí ADT\BackgroundQueue\Exception\DieException, zpracuje se vše dál podle exception, která je v ->getPrevious() (a pokud žádná není, tak jako při ADT\BackgroundQueue\Exception\PermanentErrorException). Poté je (už před další iterací) konzumer ukončen. Toho lze využít například v onError, pokud v aplikaci dojde k uzavření Doctrine Entity manageru a další iterace konzumera by opět skončili chybou.

public function onError(\Throwable $exception) {

	// Příklad 1: Bude to neopakovatelná chyba a konzumer se před další iterací ukončí.
	if (!$this->entityManager->isOpen()) {
		throw new \ADT\BackgroundQueue\Exception\DieException('EM is closed.');
	}

	// Příklad 2: Bude se zpracovávat dle toho, co je v $exception, tedy pokud je $exception instanceof TemporaryErrorException, tak to bude opakovatelná chyba, ale konzumer se také před další iterací ukončí. Používá se například pri deadlocku.
	if (!$this->entityManager->isOpen()) {
		throw new \ADT\BackgroundQueue\DieException('EM is closed. Reason: ' . $exception->getMessage(), $exception->getCode(), $exception);
	}
}

Pokud callback vyhodí jakýkoliv jiný error nebo exception implementující Throwable, záznam se uloží ve stavu TEMPORARILY_FAILED a zkusí se zpracovat při přištím spuštění background-queue:process commandu (viz níže). Po notifyOnNumberOfAttempts je zaslána notifikace. Prodleva mezi každým dalším opakováním je prodloužena o dvojnásobek času, maximálně však na dobu 16 minut.

Ve všech ostatních případech se záznam uloží jako úspěšně dokončený ve stavu STATE_FINISHED.

2.2 Commandy

background-queue:process Bez využití brokera zpracuje všechny záznamy ve stavu READY, TEMPORARILY_FAILED, WAITING a BROKER_FAILED. V případě využití brokera zařadí znovu do brokera (a přepne na READY) záznamy ve stavu STATE_BACK_TO_BROKER, zpracuje rovnou záznamy ve stavu BROKER_FAILED (těm se publikace do brokera nepovedla), pustí do hry hlavu každé skupiny, která má nějaký job ve stavu WAITING (záchranná síť pro případ, že by ji neprobudil konzument předchůdce), a znovu publikuje joby se ztracenou zprávou (viz 2.5). Command je ideální spouštět cronem každou minutu. Stav STATE_BACK_TO_BROKER je typicky nastaven ručně v databázi těm záznamům, které chceme nechat znovu zpracovat.

background-queue:monitor Jednorázově spočítá zaseklé joby: stavy TEMPORARILY_FAILED a PERMANENTLY_FAILED — jediné dva, ve kterých se joby drží dlouhodobě (ostatní stavy mají aktivní pojistku, viz 2.4 a 2.5) — a k tomu joby ve stavu PROCESSING běžící déle než 24 hodin. Ty kryjí jediný slepý bod reaperu: callback zaseknutý v nekonečné smyčce, který si přes middleware pořád posílá tep, takže pro reaper nikdy nezestárne. Počty vypíše na výstup a je-li co hlásit, pošle report i do loggeru: s úrovní critical, obsahuje-li PERMANENTLY_FAILED nebo dlouhoběžící PROCESSING joby (ani jedno se bez ručního zásahu nespraví), jinak warning. Ideální spouštět cronem jednou denně o půlnoci (0 0 * * *).

background-queue:clear-finished Smaže všechny úspěšně zpracované záznamy.

background-queue:clear-finished 14 Smaže všechny úspěšně zpracované záznamy starší 14 dní.

background-queue:reload-consumers QUEUE NUMBER Reloadne NUMBER consumerů pro danou QUEUE.

background-queue:update-schema Aktualizuje databázové schéma, pokud je potřeba.

Všechny commandy jsou chráněny proti vícenásobnému spuštění.

2.3 Callbacky

Využivát můžete také 2 callbacky onBeforeProcess a onAfterProcess, v nichž například můžete provést přepinání databází.

2.4 Osiřelé joby (zabitý consumer)

Když consumer umře uprostřed callbacku, nezapíše výsledek a záznam zůstane ve stavu PROCESSING. Ten stav není mezi zpracovatelnými, takže ho už nikdo nikdy nevyzvedne — a u jobů se serialGroup je to horší: běžící job se bere jako překážka, takže se celá skupina natrvalo zasekne ve stavu WAITING.

Dokud jde o chybu, kterou PHP stihne obsloužit (výjimka, fatální chyba), řeší to standardní stavový automat. Proti tvrdému ukončení procesu (SIGKILL od OOM killeru, pád kontejneru, reboot hostu) se ale uvnitř procesu zajistit nedá nic — žádný PHP kód se už nespustí. Na to je reaper: background-queue:process na začátku každého běhu vrátí joby, které jsou ve stavu PROCESSING a déle než stalledJobTimeout sekund se do nich nikdo nezapsal, do stavu TEMPORARILY_FAILED. Odtud je vezme běžná cesta opakování.

Živost se záměrně nezjišťuje ze sloupce pid. Ten platí jen v rámci jednoho kontejneru, takže cron na jiném hostu o něm nemůže nic tvrdit, a při jednom procesu na job se navíc rychle recykluje. Místo toho se posuzuje sloupec updated_at, který posune dopředu každý zápis do jobu.

Reaper tedy potřebuje odlišit „běží dlouho" od „consumer umřel". K tomu slouží tep: updated_at se posouvá dopředu i během běhu callbacku, takže stalledJobTimeout nemusí být nastavený podle nejdelšího možného běhu.

Máte-li nainstalovaný BackgroundQueueMiddleware (viz 6), děje se to samo a nemusíte udělat nic. Middleware sedí na aplikačním DBAL spojení, takže každý dotaz callbacku pošle tep — automaticky pro všechny joby, bez jediného řádku v jejich kódu. Throttling drží zátěž na jednom UPDATE za heartbeatInterval sekund na běžící job bez ohledu na to, kolik dotazů callback udělá.

Zbývají dva případy, kdy tep sám nedojde:

  • callback má dlouhý úsek bez dotazů do databáze (generování souboru v paměti, čekání na cizí API),
  • callback visí v jednom dlouhém dotazu — tep se pošle až po jeho dokončení.

Pro ty je heartbeat() veřejná metoda, kterou si callback může zavolat sám:

foreach ($tisiceZaznamu as $zaznam) {
	$backgroundQueue->heartbeat();
	// ... vlastní práce ...
}

Volat ji lze jakkoli často — zapisuje nejvýš jednou za heartbeatInterval sekund a mimo zpracování jobu nedělá nic.

Bez middlewaru a bez volání heartbeat() tep nechodí vůbec a reaper se řídí jen tím, kdy do jobu naposledy zapsala samotná fronta (typicky claim). V takovém případě nechte stalledJobTimeout s velkou rezervou.

2.5 Ztracené zprávy (job bez zprávy v brokeru)

Řádek v databázi je zdroj pravdy, ale k životu ho v broker módu probouzí jediná zpráva v RabbitMQ — a ta může zaniknout: proces umře mezi DB zápisem a publishem, publish zůstane viset v transakčním bufferu bez nainstalovaného middlewaru, consumer spadne mezi ackem a claimem, výpadek sítě těsně po basic_publish, ruční purge fronty. Job ve stavu READY nebo TEMPORARILY_FAILED by pak visel navždy — a job se serialGroup by s sebou blokoval i celou svou skupinu.

Pojistkou je background-queue:process: joby v těchto dvou stavech, do kterých déle než lostMessageTimeout sekund nikdo nezapsal a kterým už uplynul případný odklad (availableFrom), publikuje znovu. Falešný poplach je neškodný — pokud zpráva jen dlouho čekala v zaplněné frontě, duplicitní doručení zahodí podmíněný claim. Nastavte proto lostMessageTimeout s rezervou nad běžnou dobu čekání zprávy ve frontě, ať duplicity nevznikají zbytečně. Každý republish se zaloguje s úrovní warning.

3 Monitoring

Při spuštění consumera se do tabulky background_job do sloupce pid uloží aktuální PID procesu. Nejedná se o PID z pohledu systému, ale o PID uvnitř docker kontejneru.

Při dokončení callbacku se do tabulky background_job do sloupce memory uloží informace o využité paměti před a po dokončení. Pokud v commandu background-queue:consume využíváme parametr jobs, od verze PHP 8.2 se před každým jednotlivým zpracováním resetuje "memory peak" (metoda memory_reset_peak_usage()).

'notRealActual' => memory_get_usage(),
'realActual' => memory_get_usage(true),
'notRealPeak' => memory_get_peak_usage(),
'realPeak' => memory_get_peak_usage(true),

4 Vkládání po dávkách

Pokud vkládáme větší množství záznamů, může BackgroundQueue vkládat záznamy do DB efektivněji po dávkách pomocí INSERT INTO table () VALUES (...), (...), .... Velikost dávky se nastavuje parametrem bulkSize. Začátek a konec dávkového vkládání záznamů uvedeme metodami starBulk a endBulk. Bez započetí dávkového vkládání metodou startBulk bude vždy dávka velikosti 1, nehledě na parametr bulkSize.

$this->backgroundQueueService->startBulk();
foreach ($data as $oneJobData) {
    $this->backgroundQueue->publish(...);
}
$this->backgroundQueueService->endBulk();

5 Prioritizace záznamů

Vkládaným záznamům máme možnost určit jejich prioritu. Později vložený záznam s větší prioritou má přednost při zpracování před dříve vloženým.

V nastavení určím, jaké priority budou využívány. Parametr je nepovinný s výchozí hodnotou [1].

$backgroundQueue = new \ADT\BackgroundQueue\BackgroundQueue([
	...
	'priorities' => [10, 15, 20, 25, 30, 35, 40, 45, 50],
	...
]);

Přímo u jednotlivých callbacků pro zpracování záznamů lze určit jejich priority. Pokud callback nemá určenou prioritu, použije se nejvyšší dostupná priorita.

Pro příklad mějme následující typy úloh:

  • Přepočet ACL
  • Stahování dat z API třetích sran
  • Rozesílání emailů (např. registrační emaily)

Běžně stahujeme data z API, což může být dlouhotrvající úloha. Pokud je však potřeba přepočítat ACL, nechceme, aby to bylo blokováno stahováním dat z API. A ještě přednostněji chceme odbavit občasné emaily při registraci.

$backgroundQueue = new \ADT\BackgroundQueue\BackgroundQueue([
	...
	'priorities' => [10, 15, 20, 25, 30, 35, 40, 45, 50],
	'callbacks' => [
		'email' => [$mailer, 'process'], // záznamy budou mít prioritu 10
		'aclRecalculation' => [
			'callback' => [$aclService, 'process'],
			'priority' => 20,
		],
		'dataImporting' => [
			'callback' => [$apiService, 'process'],
			'priority' => 30,
		],
	],
	...
]);

Příkazu background-queue:consume máme možnost nastavit pomocí parametru -p, jaký rozsah priorit má zpracovávat (např. -p 10 , -p 20-40 , -p 25- , -p"-20", ...). Můžeme tedy jednoho konzumera vyhradit například na rozesílání registračním emailů (background-queue:consume -p 10) a ostatní pro všechny úlohy (background-queue:consume). Tím zajistíme, že rychlé odeslání registračního emailu nebude čekat na dlouho trvající úlohy, protože je odbaví první konzumer. Ale pokud by se vyskytlo více požadavků na zasílání emailů, po nějaké době je začnou odbavovat všichni konzumeři.

Dále máme možnost prioritu nastavenou pro callback přetížit při vkládání záznamu v metodě publish. Například víme, že se jedná o rozesílání newsletterů. Tedy se jedná o zasílání emailů, ale s nízkou prioritou zpracování.

$priority = null; // aplikuje se priorita 10 z nastavení pro callback
if ($isNewsletter) {
	$priority = 25;
}
$this->backgroundQueue->publish('email', $parameters, $serialGroup, $identifier, $isUnique, $availableAt, $priority);

6 Integrace do frameworků

7 Upgrade

Zrušení interního jobu _processWaitingJobs. Joby čekající ve stavu WAITING dřív vracel do hry periodický interní job _processWaitingJobs. Byl to jediný bod selhání celého sériového zpracování: když skončil ve stavu PERMANENTLY_FAILED, náhrada se už nikdy nepublikovala a všechny serialGroup skupiny zamrzly natrvalo. Nahradilo ho probuzení nástupce přímo konzumentem předchůdce, se záchrannou sítí v background-queue:process.

Knihovna už tenhle callback nikde neregistruje a po zbylých řádcích neuklízí. Po nasazení je smažte ručně:

DELETE FROM background_job WHERE callback_name = '_processWaitingJobs';

Dokud tam zůstanou, nic se nerozbije - jen se do logu můžou dostat chyby Callback "_processWaitingJobs" does not exist. ze starých zpráv v brokeru.

About

No description, website, or topics provided.

Resources

Stars

4 stars

Watchers

14 watching

Forks

Releases

Packages

Used by

Contributors

Languages