|
@@ -4,14 +4,24 @@ declare(strict_types=1);
|
|
|
|
|
|
|
|
namespace SixShop\System\Command;
|
|
namespace SixShop\System\Command;
|
|
|
|
|
|
|
|
|
|
+use SixShop\System\Event\CrontabWorkerStartEvent;
|
|
|
|
|
+use think\console\Command;
|
|
|
|
|
+use think\console\Input;
|
|
|
use think\console\input\Argument;
|
|
use think\console\input\Argument;
|
|
|
use think\console\input\Option;
|
|
use think\console\input\Option;
|
|
|
-use think\worker\command\Server;
|
|
|
|
|
-use think\worker\Manager;
|
|
|
|
|
|
|
+use think\console\Output;
|
|
|
|
|
+use SixShop\System\Workerman\Worker as SixShopWorker;
|
|
|
|
|
+use Workerman\Crontab\Crontab;
|
|
|
|
|
+use Workerman\Timer;
|
|
|
use Workerman\Worker;
|
|
use Workerman\Worker;
|
|
|
|
|
|
|
|
-class WorkerCommand extends Server
|
|
|
|
|
|
|
+class WorkerCommand extends Command
|
|
|
{
|
|
{
|
|
|
|
|
+ public function __construct(private \think\queue\Worker $queueWorker)
|
|
|
|
|
+ {
|
|
|
|
|
+ parent::__construct();
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
public function configure(): void
|
|
public function configure(): void
|
|
|
{
|
|
{
|
|
|
$this->setName('system:worker')
|
|
$this->setName('system:worker')
|
|
@@ -21,67 +31,106 @@ class WorkerCommand extends Server
|
|
|
->setDescription('Workerman Server for ThinkPHP (custom entry)');
|
|
->setDescription('Workerman Server for ThinkPHP (custom entry)');
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- public function handle(Manager $manager)
|
|
|
|
|
|
|
+
|
|
|
|
|
+ protected function execute(Input $input, Output $output): void
|
|
|
{
|
|
{
|
|
|
- $action = $this->input->getArgument('action') ?: 'start';
|
|
|
|
|
- $grace = $this->input->getOption('grace');
|
|
|
|
|
|
|
+ SixShopWorker::setArgv($this->inputToArgv($input));
|
|
|
|
|
+ SixShopWorker::setApp($this->app);
|
|
|
|
|
+ $this->addCronWorker();
|
|
|
|
|
+ $this->addQueueWorker();
|
|
|
|
|
+ SixShopWorker::runAll();
|
|
|
|
|
+ }
|
|
|
|
|
|
|
|
- // 启动流程,直接使用框架原生启动逻辑
|
|
|
|
|
- if ($action === 'start') {
|
|
|
|
|
- if ($grace) {
|
|
|
|
|
- putenv('WORKER_GRACE=1');
|
|
|
|
|
- }
|
|
|
|
|
- $manager->start();
|
|
|
|
|
- return;
|
|
|
|
|
|
|
+ /**
|
|
|
|
|
+ * @param Input $input
|
|
|
|
|
+ * @return array
|
|
|
|
|
+ */
|
|
|
|
|
+ private function inputToArgv(Input $input): array
|
|
|
|
|
+ {
|
|
|
|
|
+ $argv = [PHP_BINARY, 'worker', $input->getArgument('action')];
|
|
|
|
|
+ $daemon = $input->getOption('daemon');
|
|
|
|
|
+ if ($daemon) {
|
|
|
|
|
+ $argv[] = '-d';
|
|
|
}
|
|
}
|
|
|
-
|
|
|
|
|
- // 1. 重置全局$argv,绕过父类私有static::$command,底层getArgv正常识别指令
|
|
|
|
|
- global $argv;
|
|
|
|
|
- $argv = [
|
|
|
|
|
- PHP_BINARY,
|
|
|
|
|
- 'think',
|
|
|
|
|
- 'worker',
|
|
|
|
|
- $action
|
|
|
|
|
- ];
|
|
|
|
|
|
|
+ $grace = $input->getOption('grace');
|
|
|
if ($grace) {
|
|
if ($grace) {
|
|
|
$argv[] = '-g';
|
|
$argv[] = '-g';
|
|
|
}
|
|
}
|
|
|
|
|
+ return $argv;
|
|
|
|
|
+ }
|
|
|
|
|
|
|
|
- $pidFile = runtime_path() . 'worker.pid';
|
|
|
|
|
- $masterPid = file_exists($pidFile) ? (int)file_get_contents($pidFile) : 0;
|
|
|
|
|
- // 检测进程是否运行
|
|
|
|
|
- if ($masterPid <= 0 || !posix_kill($masterPid, 0)) {
|
|
|
|
|
- $this->output->error('Worker is not running');
|
|
|
|
|
- return;
|
|
|
|
|
- }
|
|
|
|
|
|
|
+ /**
|
|
|
|
|
+ * @return void
|
|
|
|
|
+ */
|
|
|
|
|
+ private function addCronWorker(): void
|
|
|
|
|
+ {
|
|
|
|
|
+ SixShopWorker::addWorker(
|
|
|
|
|
+ onWorkerStart: function (Worker $worker) {
|
|
|
|
|
+ $this->app->event->trigger(CrontabWorkerStartEvent::class, $worker);
|
|
|
|
|
+ },
|
|
|
|
|
+ name: 'crontab',
|
|
|
|
|
+ count: 1,
|
|
|
|
|
+ onWorkerStop: function (Worker $worker) {
|
|
|
|
|
+ foreach (Crontab::getAll() as $crontab) {
|
|
|
|
|
+ /** @var Crontab $crontab */
|
|
|
|
|
+ $crontab->destroy();
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ );
|
|
|
|
|
+ }
|
|
|
|
|
+ /**
|
|
|
|
|
+ * @return void
|
|
|
|
|
+ */
|
|
|
|
|
+ private function addQueueWorker(): void
|
|
|
|
|
+ {
|
|
|
|
|
+ SixShopWorker::addWorker(
|
|
|
|
|
+ onWorkerStart: function (Worker $worker) {
|
|
|
|
|
+ $shouldExit = false;
|
|
|
|
|
+ $logger = $this->app->log;
|
|
|
|
|
+ $connection = $this->app->config->get('queue.default');
|
|
|
|
|
+ $queue = $this->app->config->get("queue.connections.{$connection}.queue", 'default');
|
|
|
|
|
+ $delay = 0;
|
|
|
|
|
+ $sleep = 2;
|
|
|
|
|
+ $tries = 3;
|
|
|
|
|
+ $memoryLimit = 128; // 内存上限MB
|
|
|
|
|
+ $queueWorker = $this->queueWorker;
|
|
|
|
|
|
|
|
- match ($action) {
|
|
|
|
|
- // 查询类:调用Workerman底层原生打印表格逻辑
|
|
|
|
|
- 'status', 'connections' => $this->renderWorkerStatus(),
|
|
|
|
|
|
|
+ $run = function () use (&$run, &$shouldExit, $queueWorker, $connection, $queue, $delay, $sleep, $tries, $memoryLimit, $logger, $worker) {
|
|
|
|
|
+ if ($shouldExit) {
|
|
|
|
|
+ $logger->info(sprintf("[pid:%d] 收到退出标记,停止消费调度", $worker->pid));
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
|
|
|
- // 信号控制类
|
|
|
|
|
- 'stop' => posix_kill($masterPid, $grace ? SIGQUIT : SIGINT),
|
|
|
|
|
- 'reload' => posix_kill($masterPid, $grace ? SIGUSR2 : SIGUSR1),
|
|
|
|
|
- 'restart' => posix_kill($masterPid, $grace ? SIGQUIT : SIGINT),
|
|
|
|
|
|
|
+ try {
|
|
|
|
|
+ $queueWorker->runNextJob($connection, $queue, $delay, $sleep, $tries);
|
|
|
|
|
+ } catch (\Throwable $e) {
|
|
|
|
|
+ $logger->error(sprintf(
|
|
|
|
|
+ "[pid:%d] 队列底层异常:%s \n堆栈:%s",
|
|
|
|
|
+ $worker->pid,
|
|
|
|
|
+ $e->getMessage(),
|
|
|
|
|
+ $e->getTraceAsString()
|
|
|
|
|
+ ));
|
|
|
|
|
+ }
|
|
|
|
|
|
|
|
- default => $this->output->error("Unknown action: {$action}"),
|
|
|
|
|
- };
|
|
|
|
|
|
|
+ $usedMb = memory_get_usage(true) / 1024 / 1024;
|
|
|
|
|
+ if ($usedMb > $memoryLimit) {
|
|
|
|
|
+ $logger->info(sprintf(
|
|
|
|
|
+ "[pid:%d] 内存超限 %.2fMB > %dMB,进程标记退出自愈",
|
|
|
|
|
+ $worker->pid,
|
|
|
|
|
+ $usedMb,
|
|
|
|
|
+ $memoryLimit
|
|
|
|
|
+ ));
|
|
|
|
|
+ $shouldExit = true;
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
|
|
|
- // 非查询指令输出提示
|
|
|
|
|
- if (!in_array($action, ['status', 'connections'])) {
|
|
|
|
|
- $this->output->info("{$action} signal sent to master PID: {$masterPid}");
|
|
|
|
|
- }
|
|
|
|
|
- }
|
|
|
|
|
|
|
+ Timer::add(1, $run, [], false);
|
|
|
|
|
+ };
|
|
|
|
|
|
|
|
- /**
|
|
|
|
|
- * 复用Workerman底层 status / connections 原生输出逻辑
|
|
|
|
|
- */
|
|
|
|
|
- protected function renderWorkerStatus(): void
|
|
|
|
|
- {
|
|
|
|
|
- Worker::$pidFile = runtime_path() . 'worker.pid';
|
|
|
|
|
- Worker::$statusFile = runtime_path() . 'worker.status';
|
|
|
|
|
- Worker::$logFile = runtime_path() . 'worker.log';
|
|
|
|
|
- // 底层自动读取重写后的$argv,输出进程/连接表格
|
|
|
|
|
- Worker::runAll();
|
|
|
|
|
|
|
+ Timer::add(1, $run, [], false);
|
|
|
|
|
+ },
|
|
|
|
|
+ name: 'queue',
|
|
|
|
|
+ count: 1,
|
|
|
|
|
+ );
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|