Symfony bundle for convenient work with queues. Currently it supports RabbitMQ.
Install bundle
composer require lamoda/queue-bundle
Extend
Lamoda\QueueBundle\Entity\QueueEntityMappedSuperclassuseDoctrine\ORM\MappingasORM; useLamoda\QueueBundle\Entity\QueueEntityMappedSuperclass; /** * @ORM\Entity(repositoryClass="Lamoda\QueueBundle\Entity\QueueRepository") */class Queue extends QueueEntityMappedSuperclass { }
Configure bundle parameters
lamoda_queue: ## requiredentity_class: App\Entity\Queuemax_attempts: 5batch_size_per_requeue: 5batch_size_per_republish: 5## optional (will use for default delay Geometric Progression Strategy)strategy_delay_geometric_progression_start_interval_sec: 60strategy_delay_geometric_progression_multiplier: 2
Register bundle
class AppKernel extends Kernel { // ...publicfunctionregisterBundles() { $bundles = [ // ...newLamoda\QueueBundle\LamodaQueueBundle(), // ... ]; return$bundles; } // ... }
or add to
config/bundles.phpreturn [ // ...Lamoda\QueueBundle\LamodaQueueBundle::class => ['all' => true], // ... ];
Migrate schema
doctrine:migrations:diffto create migration forqueuetabledoctrine:migrations:migrate- apply the migration
Define new exchange constant
namespaceApp\Constant; class Exchanges { publicconstDEFAULT = 'default'; }
Add new node to
old_sound_rabbit_mq.producerswith previous defined constant name, example:old_sound_rabbit_mq: producers: default: connection: defaultexchange_options: name: !php/const App\Constant\Exchanges::DEFAULTtype: "direct"
Define new queue constant
namespaceApp\Constant; class Queues { publicconstNOTIFICATION = 'notification'; }
Register consumer for queue in
old_sound_rabbit_mq.consumerswith previous defined constant name, example:old_sound_rabbit_mq: consumers: notification: connection: defaultexchange_options: name: !php/const App\Constant\Exchanges::DEFAULTtype: "direct"queue_options: name: !php/const App\Constant\Queues::NOTIFICATIONrouting_keys: - !php/constApp\Constant\Queues::NOTIFICATIONcallback: "lamoda_queue.consumer"
Create job class, extend
AbstractJobby example:namespaceApp\Job; useApp\Constant\Exchanges; useApp\Constant\Queues; useLamoda\QueueBundle\Job\AbstractJob; useJMS\Serializer\AnnotationasJMS; class SendNotificationJob extends AbstractJob { /** * @var string * * @JMS\Type("int") */private$message; publicfunction__construct(string$message) { $this->message = $message; } publicfunctiongetDefaultQueue(): string { return Queues::NOTIFICATION; } publicfunctiongetDefaultExchange(): string { return Exchanges::DEFAULT; } }
Create job handler, implement HandlerInterface by example:
namespaceApp\Handler; useLamoda\QueueBundle\Handler\HandlerInterface; useLamoda\QueueBundle\QueueInterface; class SendNotificationHandler implements HandlerInterface { publicfunctionhandle(QueueInterface$job): void { // implement service logic here } }
Tag handler at service container
services: App\Handler\SendNotificationHandler: public: truetags: - { name: queue.handler, handle: App\Job\SendNotificationJob }
Configure delay strategy parameters (optional)
lamoda_queue: ... queues: queue_one: 'delay_arithmetic_progression'queue_two: 'delay_special_geometric_progression'# Settings of special behaviors for Delay strategies (optional)services: lamoda_queue.strategy.delay.arithmetic_progression: class: Lamoda\QueueBundle\Strategy\Delay\ArithmeticProgressionStrategytags: - { name: 'lamoda_queue_strategy', key: 'delay_arithmetic_progression' }arguments: - 60# start_interval_sec parameter - 1700# multiplier parameterlamoda_queue.strategy.delay.geometric_progression: class: Lamoda\QueueBundle\Strategy\Delay\GeometricProgressionStrategytags: - { name: 'lamoda_queue_strategy', key: 'delay_special_geometric_progression' }arguments: - 70# start_interval_sec parameter - 4# multiplier parameter
In this block, you can config special delay behaviors for each queue. For this, you have to register new services that use one of several base strategies (ArithmeticProgressionStrategy, GeometricProgressionStrategy) or yours (for this you have to make Service Class that implements DelayStrategyInterface).
Each strategy service has to have a tag with name
lamoda_queue_strategyand uniquekey. After this, you can use thesekeysfor matching with queues inlamoda_queue.queuessection.By default, use GeometricProgressionStrategy with params (their you can customize in
lamoda_queueconfig section):strategy_delay_geometric_progression_start_interval_sec: 60 strategy_delay_geometric_progression_multiplier: 2Add queue name in "codeception.yml" at
modules.config.AMQP.queuesExecute
./bin/console queue:initcommand
./bin/console queue:init
$job = newSendNotificationJob($id);
$container->get(Lamoda\QueueBundle\Factory\PublisherFactory::class)->publish($job);./bin/console queue:consume notification
./bin/console queue:requeue
You can queue any primitive class, just implement QueueInterface:
namespaceApp\Process;
useLamoda\QueueBundle\Entity\QueueInterface;
class MyProcess implements QueueInterface
{
// implement interface functions
}services:
App\Handler\MyProcessHandler:
public: truetags:
- { name: queue.handler, handle: App\Process\MyProcess }$process = newMyProcess();
$container->get('queue.publisher')->publish($process);If you want to rerun queue, throw Lamoda\QueueBundle\Exception\RuntimeException.
If you want mark queue as failed, throw any another kind of exception.
namespaceApp\Handler;
useLamoda\QueueBundle\Handler\HandlerInterface;
useLamoda\QueueBundle\QueueInterface;
class SendNotificationHandler implements HandlerInterface
{
publicfunctionhandle(QueueInterface$job): void
{
// implement service logic here// Rerun queueif ($rerun === true) {
thrownewLamoda\QueueBundle\Exception\RuntimeException('Error message');
}
// Mark queue as failedif ($failed === true) {
thrownew \Exception();
}
}
}By default delay time is calculated exponentially. You can affect it through configuration.
lamoda_queue:
## required## ...max_attempts: 5## optionalstrategy_delay_geometric_progression_start_interval_sec: 60strategy_delay_geometric_progression_multiplier: 2When consumer wants to execute reached maximum attempts queue.
Properties:
- Queue Entity
QueueAttemptsReachedEvent::getQueue()
make php-cs-check
make php-cs-fixUnit
make test-unit

