As of version 0.4.0 the worker has been rewritten to a flexible event driven approach. The processing logic is now a very minimalistic method. In pseudocode it looks like this;
processQueue
trigger event 'bootstrap'
while event says continue processing
trigger event 'process.queue'
trigger event 'process.job'
trigger event 'finish'
trigger event 'state'
To get some useful results it is required to register so called 'worker strategies' to the worker. SlmQueue makes this trivial via configuration.
Worker strategies are aggregate listeners which are created via a plugin manager.
At least one worker strategy listening to the bootstrap event must be registered to the worker. The Worker Factory will
throw an exception if its not. SlmQueue attaches the provided AttachQueueListenersStrategy to do just that.
It is worth noting that events will be dispatched from the worker (obviously) but can also be dispatch from within worker strategies.
The plugin manager ensures they extend SlmQueue\Listener\Strategy\AbstractStrategy and each worker strategy therefore
gains the following capabilities;
Configuration options are passed by the plugin manager to the constructor of an worker strategy. Setter methods will be called for each option. If a setter does not exist an exception will be thrown.
'SlmQueue\Strategy\MaxRunStrategy' => array('max_runs' => 10);Such a config will result in an MaxRunStrategy instance of which the setMaxRuns method is called with '10'.
The optional 'priority' option is used when the aggregates listeners are are registered with event manager and is thereafter removed from the passed options. This means a Worker Strategy cannot have this option.
Worker strategies may inform the worker to stop processing the queue. Or more concrete; invalidate the condition of the while loop.
public function onSomeListener(WorkerEvent $event)
{
$event->exitWorkerLoop();
...
}While processing a queue it might be required to execute some setup- or teardown logic. A worker strategy may listen to
the bootstrap and/or finish event to do just this.
/**
* @param EventManagerInterface $events
*/
public function attach(EventManagerInterface $events)
{
$this->listeners[] = $events->attach(
WorkerEvent::EVENT_BOOTSTRAP,
array($this, 'onBootstrap')
);
$this->listeners[] = $events->attach(
WorkerEvent::EVENT_FINISH,
array($this, 'onFinish')
);
}
/**
* @param WorkerEvent $e
*/
public function onBootstrap(WorkerEvent $e)
{
// setup code
}
/**
* @param WorkerEvent $e
*/
public function onFinish(WorkerEvent $e)
{
// teardown code
}For some types of jobs it might be required to do something before or after the execution of an individual job.
This can be done by listening to the process event at different priorities.
/**
* @param EventManagerInterface $events
*/
public function attach(EventManagerInterface $events)
{
$this->listeners[] = $events->attach(
WorkerEvent::EVENT_PROCESS_JOB,
array($this, 'onPreProcess'),
100
);
$this->listeners[] = $events->attach(
WorkerEvent::EVENT_PROCESS_JOB,
array($this, 'onPostProcess'),
-100
);
}
/**
* @param WorkerEvent $e
*/
public function onPreProcess(WorkerEvent $e)
{
// pre job execution code
}
/**
* @param WorkerEvent $e
*/
public function onPostProcess(WorkerEvent $e)
{
// post job execution code
}A worker strategy may report a state. A state is simply a string which will be returned when
the WorkerEvent::EVENT_PROCESS_STATE is received.
From the MaxRunStrategy;
public function onSomeListener(WorkerEvent $event)
{
$this->runCount++;
if ($this->maxRuns && $this->runCount >= $this->maxRuns) {
$event->exitWorkerLoop();
$this->state = sprintf('maximum of %s jobs processed', $this->runCount);
} else {
$this->state = sprintf('%s jobs processed', $this->runCount);
}
}A worker strategy may ask the worker to dispatch events.
From the ProcessQueueStrategy
public function onJobPop(WorkerEvent $e)
{
$queue = $e->getQueue();
$options = $e->getOptions();
$job = $queue->pop($options);
// The queue may return null, for instance if a timeout was set
if (!$job instanceof JobInterface) {
/** @var AbstractWorker $worker */
$worker = $e->getTarget();
$worker->getEventManager()->trigger(WorkerEvent::EVENT_PROCESS_IDLE, $e);
// make sure the event doesn't propagate
$e->stopPropagation();
return;
}
$e->setJob($job);
}Worker strategies are regular ZF2 services that are instanciated via a plugin manager. If a worker strategy has dependancies on other services it should be created it via service factory.
The plugin manager is configured to not share services.
Events the worker and worker strategies may dispatch;
WorkerEvent::EVENT_BOOTSTRAPjust before loop is enteredWorkerEvent::EVENT_FINISHjust after the loop has exitedWorkerEvent::EVENT_PROCESS_QUEUEfetch job(s)WorkerEvent::EVENT_PROCESS_JOBprocesses job(s)WorkerEvent::EVENT_PROCESS_IDLEwhen the queue is emptyWorkerEvent::EVENT_PROCESS_STATEcollects 'states' from strategies.
Any listener waiting for above events will be passed a WorkerEvent class which contains a reference to the queue.
The EVENT_PROCESS_JOB event also has a reference to the job instance.
$em->attach(WorkerEvent::EVENT_PROCESS_JOB, function(WorkerEvent $e) {
$queue = $e->getQueue();
$job = $e->getJob();
});In above example, $em refers to the event manager inside the worker object: $em = $worker->getEventManager();.
When a job is processed, the job or worker returns a status code. You can use a listener to act upon this status, for example to log any failed jobs:
$logger = $sm->get('logger');
$em->attach(WorkerEvent::EVENT_PROCESS_JOB, function(WorkerEvent $e) use ($logger) {
$result = $e->getResult();
if (WorkerEvent::JOB_STATUS_FAILURE === $result) {
$job = $e->getJob();
$logger->warn(sprintf(
'Job #%s (%s) failed executing', $job->getId(), get_class($job)
));
}
}, -1000);The purpose of this strategy is to register additional strategies that are specific to the queue that is being processed.
After registering any additional worker strategies it will unregister itself as a listener. Finally it halts the event
propagation and re-triggers the bootstrap event.
A new cycle of bootstraping will occure but now with additional queue specific strategies.
listens to:
bootstrapevent at priority PHP_MAX_INT
triggers:
bootstrap
throws:
- RunTimeException if the
process.queueevent isn't listened to by any registered strategy.
This strategy is enabled by default for all queue's.
This strategy is able to 'watch' files by creating a hash of their contents. If it detects a change it will request to stop processing the queue. This is useful if you have something like supervisor automatically restarting the worker process.
The strategy builds a list of files it needs to watch via a preg_match on the filenames within the application.
listens to:
process.idleevent at priority 1process.jobevent at priority 1000process.stateevent at priority 1
options:
- pattern defaults to '/^./(config|module).*.(php|phtml)$/'
- idle_throttle_time defaults to 300 seconds
This strategy is not enabled by default. It can be slow and is recommended for development only. In production you may
watch a single file. It will run the check before each job and while idling after idle_throttle_time seconds
have passed.
The InterruptStrategy is able to catch a stop condition under Linux-like systems (as well as OS X). If a worker is started from the command line interface (CLI), it is possible to send a SIGTERM or SiGINT call to the worker. SlmQueue is smart enough not to quit the script directly, but let the job finish its work first and then break out of the loop. On Windows systems this strategy does nothing.
listens to:
process.idleevent at priority 1process.queueevent at priority -1000process.stateevent at priority 1
This strategy is enabled by default for all queue's.
Simple outputs to the console a when it begin executing a job and when it is done.
listens to:
process.jobevent at priority 1000process.jobevent at priority -1000
The MaxMemoryStrategy will measure the amount of memory allocated to PHP after each processed job. It will request to exit when a threshold is exceeded.
Note that an individual job may exceed this threshold during it's live time. But if you have a memory leak this strategy can make sure the script aborts eventually.
listens to:
process.idleevent at priority 1process.queueevent at priority -1000process.stateevent at priority 1
options:
- max_memory defaults to 100*1024*1024
This strategy is enabled by default for all queue's.
The MaxRunStrategy will request to exit after a set number of jobs have been processed.
listens to:
process.idleevent at priority 1process.jobevent at priority -1000process.stateevent at priority 1
options:
- max_runs defaults to 100000
This strategy is enabled by default for all queue's.
Responsible for quering the queue for jobs and executing them.
listens to:
process.queueevent at priority 1process.jobevent at priority 1
triggers:
process.jobfor each job pop'ed from the queueprocess.idleif the queue returns null (it might be empty or timed out)
The MaxPollingFrequencyStrategy ensure the polling frequency don't exceed a defined value. This can be useful in the case where you are using a system like SQS which makes you pay the service per request.
listens to:
process.queueevent at priority 1000
options:
- max_frequency
This strategy is not enabled by default.
| Frequency | x / sec | x / min | x / hour | x / day | x / month |
|---|---|---|---|---|---|
| 0.1 | 0.1 | 6 | 360 | 8640 | 259200 |
| 0.2 | 0.2 | 12 | 720 | 17280 | 518400 |
| 0.5 | 0.5 | 30 | 1800 | 43200 | 1296000 |
| 1 | 1 | 60 | 3600 | 86400 | 2592000 |
| 2 | 2 | 120 | 7200 | 172800 | 5184000 |
| 5 | 5 | 300 | 18000 | 432000 | 12960000 |
| 10 | 10 | 600 | 36000 | 864000 | 25920000 |
| Frequency | Delay |
|---|---|
| 0.000278 | 1 hour |
| 0.0167 | 1 min |
| 0.1 | 10 s |
| 0.2 | 5 s |
| 0.5 | 2 s |
| 1 | 1 s |
| 2 | 500 ms |
| 5 | 200 ms |
| 10 | 100 ms |
In the SlmQueue config, find the part named worker_strategies and add the
following line:
'SlmQueue\Strategy\MaxPollingFrequencyStrategy' => array('max_frequency' => 1)Replace the max_frequency value helping you with the tables above.
Instead of direct access to the worker's event manager, the shared manager is available to register events too:
<?php
namespace MyModule;
use SlmQueue\Worker\WorkerEvent;
use Zend\Mvc\MvcEvent;
class Module
{
public function onBootstrap(MvcEvent $e)
{
$em = $e->getApplication()->getEventManager();
$sharedEm = $em->getSharedManager();
$sharedEm->attach('SlmQueue\Worker\WorkerInterface', WorkerEvent::EVENT_PROCESS_JOB, function() {
// some thing just before a job starts.
}, 1000);
}
}A good example is i18n: a job is given a locale if the job performs localized actions. This locale is set to the translator just before processing starts. The original locale is reverted when the job has finished processing.
In this case, all jobs which require a locale set are implementing a LocaleAwareInterface:
<?php
namespace MyModule\Job;
interface LocaleAwareInterface
{
/**
* @param string $locale
*/
public function setLocale($locale);
/**
* @return string
*/
public function getLocale();
}An worker strategy will listen for two events to set and revert the locale:
<?php
namespace MyModule\Strategy;
use MyModule\Job\LocaleAwareInterface;
use SlmQueue\Listener\Strategy\AbstractStrategy;
use SlmQueue\Worker\WorkerEvent;
use Zend\EventManager\EventManagerInterface;
use Zend\I18n\Translator\Translator;
class JobTranslatorStrategy extends AbstractStrategy
{
/**
* @var Stores original locale while processing a Job
*/
protected $locale;
/**
* @var Instance of Translator
*/
protected $translator;
/**
* @param Translator $translator
*/
public function __construct(Translator $translator)
{
$this->translator = $translator;
}
/**
* {@inheritDoc}
*/
public function attach(EventManagerInterface $events)
{
$this->listeners[] = $events->attach(WorkerEvent::EVENT_PROCESS_JOB, array($this, 'onPreJobProc'), 1000);
$this->listeners[] = $events->attach(WorkerEvent::EVENT_PROCESS_JOB, array($this, 'onPostJobProc'), -1000);
}
public function onPreJobProcessing(WorkerEvent $e)
{
$job = $e->getJob();
if (!$job instanceof LocaleAwareInterface) {
return;
}
$this->locale = $this->translator->getLocale();
$this->translator->setLocale($job->getLocale());
}
public function onPostJobProcessing(WorkerEvent $e)
{
$job = $e->getJob();
if (!$job instanceof LocaleAwareInterface) {
return;
}
$this->translator->setLocale($this->locale);
}
}Since this worker strategy has a dependency that needs to be injected we should create a factory for it.
<?php
namespace MyModule\Strategy\Factory;
use MyModule\Strategy\JobTranslatorStrategy;
use Zend\ServiceManager\FactoryInterface;
use Zend\ServiceManager\ServiceLocatorInterface;
class JobTranslatorStrategyFactory implements FactoryInterface
{
/**
* Create service
*
* @param ServiceLocatorInterface $serviceLocator
* @return JobTranslatorStrategy
*/
public function createService(ServiceLocatorInterface $serviceLocator)
{
$sm = $serviceLocator->getServiceLocator();
/** @var $sm \Zend\Mvc\I18n\Translator */
$translator = $sm->get('MvcTranslator');
$strategy = new JobTranslatorStrategy($translator);
return $strategy;
}
}Finally add two configuration settings;
- Register the factory to the plugin manager to the Strategy Manager.
- Add the strategy by name to the worker strategies. Note we can do this for all queue's or for specific ones.
<?php
return array(
'slm_queue' => array(
/**
* Worker Strategies
*/
'worker_strategies' => array(
'default' => array( // per worker
// add it here to enable the
),
'queues' => array( // per queue
'my-queue' => array(
'MyModule\Strategy\JobTranslatorStrategy',
)
),
),
/**
* Strategy manager
*/
'strategy_manager' => array(
'factories' => array(
'MyModule\Strategy\JobTranslatorStrategy' => 'MyModule\Strategy\Factory\JobTranslatorStrategyFactory',
)
),
)
);Previous page: Workers Next page: Worker Management