Repository navigation
Expand file tree
/
Copy pathConsumerProcess.php
More file actions
59 lines (49 loc) · 1.51 KB
/
Copy pathConsumerProcess.php
File metadata and controls
59 lines (49 loc) · 1.51 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
<?php
declare(strict_types=1);
/**
* This file is part of Hyperf.
*
* @link https://www.hyperf.io
* @document https://hyperf.wiki
* @contact [email protected]
* @license https://github.com/hyperf/hyperf/blob/master/LICENSE
*/
namespace Hyperf\AsyncQueue\Process;
use Hyperf\AsyncQueue\Driver\DriverFactory;
use Hyperf\AsyncQueue\Driver\DriverInterface;
use Hyperf\Contract\StdoutLoggerInterface;
use Hyperf\Process\AbstractProcess;
use Psr\Container\ContainerInterface;
class ConsumerProcess extends AbstractProcess
{
/**
* @var string
*/
protected $queue = 'default';
/**
* @var DriverInterface
*/
protected $driver;
/**
* @var array
*/
protected $config;
public function __construct(ContainerInterface $container)
{
parent::__construct($container);
$factory = $this->container->get(DriverFactory::class);
$this->driver = $factory->get($this->queue);
$this->config = $factory->getConfig($this->queue);
$this->name = "queue.{$this->queue}";
$this->nums = $this->config['processes'] ?? 1;
}
public function handle(): void
{
if (! $this->driver instanceof DriverInterface) {
$logger = $this->container->get(StdoutLoggerInterface::class);
$logger->critical(sprintf('[CRITICAL] process %s is not work as expected, please check the config in [%s]', ConsumerProcess::class, 'config/autoload/queue.php'));
return;
}
$this->driver->consume();
}
}