You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Laravel数据库队列执行MQTT发布命令失败,同步队列正常

问题:Laravel队列database模式下MQTT任务仅成功执行一个

我正在开发一个应用,其中包含向MQTT代理发布命令的后台任务。当队列使用sync模式时一切正常,但切换为database模式时,两个独立分发的任务仅能成功发送1个命令。Laravel日志和supervisor日志均无报错信息。我一整天都在搜索并尝试各种解决方案,但均无效。

相关代码

MQTT服务类

<?php

namespace App\Services;

use Illuminate\Contracts\Container\BindingResolutionException;
use PhpMqtt\Client\Exceptions\ConfigurationInvalidException;
use PhpMqtt\Client\Exceptions\ConnectingToBrokerFailedException;
use PhpMqtt\Client\Exceptions\ClientNotConnectedToBrokerException;
use PhpMqtt\Client\Exceptions\RepositoryException;
use PhpMqtt\Client\Exceptions\PendingMessageAlreadyExistsException;
use PhpMqtt\Client\Exceptions\DataTransferException;
use PhpMqtt\Client\MqttClient;
use Psr\Container\NotFoundExceptionInterface;
use Psr\Container\ContainerExceptionInterface;
use App\Models\Devices\Device;

class MqttService {

    public function sendCommand(Device $device, string $command, int $qos = 1)
    {
        $host = config('mqtt.host');
        $port = config('mqtt.port');
        $clientId = config('mqtt.client_id');


        $mqtt = new MqttClient($host, $port, $clientId);
        $mqtt->connect();

        $userId = $device->user->id;
        $deviceId = $device->id;
        $uniqueId = $device->unique_id;
        $topic = "v2/$userId/$deviceId/commands";

        if ($device->recognized_by_unique_id) {
            $topic = "v3/$uniqueId/commands";
        }

        $mqtt->publish($topic, $command, $qos);
        $mqtt->disconnect();
    }

    /**
     * Publishes error message for a specific device and user on the Mqtt broker.
     * 
     * @param string|null $userId the user id - first subtopic
     * @param string|null $deviceId the device id - last sub topic
     * @param string|null $uniqueId the device unique id - if it is using unique id instead of user id and device id
     * @param string|null $errorAsJson the json error message to publish on the mqtt broker
     * 
     * @return void 
     * 
     * @throws BindingResolutionException 
     * @throws NotFoundExceptionInterface 
     * @throws ContainerExceptionInterface 
     * @throws ConfigurationInvalidException 
     * @throws ConnectingToBrokerFailedException 
     * @throws ClientNotConnectedToBrokerException 
     * @throws RepositoryException 
     * @throws PendingMessageAlreadyExistsException 
     * @throws DataTransferException 
     */
    public static function sendError(?string $userId = null, ?string $deviceId = null, ?string $uniqueId = null, string $errorAsJson) 
    {
        $host = config('mqtt.host');
        $port = config('mqtt.port');
        $clientId = config('mqtt.client_id');

        if ($uniqueId) {
            $topic = "errors/$uniqueId";
        } else {
            $topic = "errors/$userId/$deviceId";
        }

        $mqtt = new MqttClient($host, $port, $clientId);
        $mqtt->connect();
        $mqtt->publish($topic, $errorAsJson, 0);
        $mqtt->disconnect();
    }

}

任务类

<?php

namespace App\Jobs;

use App\Services\CommandService;
use Illuminate\Bus\Queueable;
use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Foundation\Bus\Dispatchable;
use Illuminate\Queue\InteractsWithQueue;
use Illuminate\Queue\SerializesModels;
use App\Models\Devices\Device;

class SendCommand implements ShouldQueue
{
    use Dispatchable, InteractsWithQueue, Queueable, SerializesModels;

    /**
     * Create a new job instance.
     */
    public function __construct(public Device $device, public string $command) {}

    /**
     * Execute the job.
     */
    public function handle(): void
    {
        info("Before: {$this->device->unique_id}; {$this->command}");
        CommandService::sendCommands($this->device, $this->command);
        info("After: {$this->device->unique_id}; {$this->command}");
    }
}

CommandService类相关部分

public static function sendCommands(Device $device, string $command)
{
    $mqtt = new MqttService;

    $mqtt->sendCommand($device, $command);

    return true;
}

现象说明

我添加了日志记录,日志显示两个任务的Before和After都有输出,只是先输出两条Before,再输出两条After,考虑到任务是并行执行,这应该是正常的。无法成功发送的设备任务是随机的,有时是设备1,有时是设备2,偶尔两个都能成功,但这种情况很少见。

求解决思路?

内容的提问来源于stack exchange,提问作者Petar Vasilev

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.21 07:12:44