Symfony Messenger component loses connection and duplicates saving records

Viewed 421

The scheme of work is as follows: I send json to the controller (then I convert json into an object) -> then I send the object to the queue (I use doctrines as a transport) -> then the handler reads the message and saves the object to the database. Problem: when processing messages asynchronously, several duplicates of the object are saved in the database + the handler writes about errors in the connection (see screenshot).

Found a description of this problem on symfony casts, they suggest using the flush function when saving an object to the database, without first using the persist function. With this approach, saving to the database does not occur.

With a synchronous save, this error does not occur.

Does anyone know a solution to this problem?

ApiController:

 /**
 * @Route("/api/task")
 * @param Request $request
 * @param MessageBusInterface $messageBus
 * @return Response
 */
public function createTask(Request $request, MessageBusInterface $messageBus): Response
{
    try {
        $serializer = SerializerBuilder::create()->build();

        $this->inputCalculatorTaskDto = $serializer->deserialize($request->getContent(),
            InputCalculatorTaskDto::class,
            'json');

        $this->inputCalculatorTaskDto
            ->setTaskGuid(Uuid::uuid4())
            ->setServiceName(ServicesListEnum::SERVICE_CALULATOR)
            ->setTaskName(TaskListEnum::CALCULATE_ORDER_PRICE)
            ->setCreatedAt(DateTime::createFromFormat('U.u', microtime(TRUE)));

        $calculatorTaskMessage = new CalculatorTaskMessage($this->inputCalculatorTaskDto);
        $messageBus->dispatch($calculatorTaskMessage);

        $this->logger->info('done ' . $this->inputCalculatorTaskDto->getTaskGuid());

        return new Response('done');

    } catch (\Throwable $t) {
        $this->logger->error($t->getMessage());

        return new Response('error ' . $t->getMessage());
    }
}

Message:

class CalculatorTaskMessage
{
    /**
     * @var InputCalculatorTaskDto
     */
    private InputCalculatorTaskDto $inputCalculatorTaskDto;

    public function __construct(InputCalculatorTaskDto $inputCalculatorTaskDto)
    {
        $this->inputCalculatorTaskDto = $inputCalculatorTaskDto;
    }

    /**
     * @return InputCalculatorTaskDto
     */
    public function getInputCalculatorTaskDto(): InputCalculatorTaskDto
    {
        return $this->inputCalculatorTaskDto;
    }
}

MessageHandler

class CalculatorNewTaskCreatorHandler implements MessageHandlerInterface
{
    public InputCalculatorTaskDtoToTaskEntityImpl $inputCalculatorTaskDtoToTaskEntity;
    public EntityManagerInterface $entityManager;
    public TasksRepository $tasksRepository;
    private CalculatorTaskProducer $calculatorTaskProducer;
    
    /**
        * CalculatorTaskCreatorHandler constructor.
        * @param EntityManagerInterface $entityManager
        * @param TasksRepository $tasksRepository
        * @param CalculatorTaskProducer $calculatorTaskProducer
    */
    public function __construct(EntityManagerInterface $entityManager,
    TasksRepository $tasksRepository,
    CalculatorTaskProducer $calculatorTaskProducer)
    {
        $this->inputCalculatorTaskDtoToTaskEntity = new InputCalculatorTaskDtoToTaskEntityImpl();
        $this->entityManager = $entityManager;
        $this->tasksRepository = $tasksRepository;
        $this->calculatorTaskProducer = $calculatorTaskProducer;
    }
    
    public function __invoke(CalculatorTaskMessage $calculatorTaskMessage)
    {
        $inputCalculatorTaskDto = $calculatorTaskMessage->getInputCalculatorTaskDto();
        $tasks = $this->inputCalculatorTaskDtoToTaskEntity->map($inputCalculatorTaskDto);
        
        $this->entityManager->persist($tasks);
        $this->entityManager->flush();
        
        $this->calculatorTaskProducer->sendMessage($inputCalculatorTaskDto);
    }
}   

enter image description here

1 Answers

This is indeed a frustrating problem. Luckily, the solution to your duplication problem should be simple.

The problem that occurs is that when you publish your message to your queue and consuming it again from the queue, you are not using the same entity manager anymore. In your consumer part, Doctrine does not recognise the object anymore, and treats it like an unmanaged entity. That is why Doctrine is inserting multiple instances of the same thing.

How you can fix this:

In your consumer, do a lookup before you try to persist your object (from your message). If the entity is not found, you know it is a new entity and you can safely insert it.

If an existing entity is found, Doctrine understands that it already exists and will not insert it again. At this point you can update the values of the entity (if needed) and then persist & flush it. At this point, your entity should be updated and you should not have any duplications anymore.

Related