Skip to content
advanced Phase 88 · Integration Advanced

Queue-Based Integration

Queue-based integration - async message processing, queue patterns, integration buses

45m
0 problems
Topic Progress 0%

Queue Architecture

Magento 2 Message Queue

┌─────────────┐     ┌──────────────────┐     ┌─────────────┐
│   Producer  │────►│   Message Queue  │────►│   Consumer  │
│   (Event)   │     │   (RabbitMQ)     │     │   (Worker)  │
└─────────────┘     └──────────────────┘     └─────────────┘

Queue Configuration

<!-- app/code/Vendor/Integration/etc/queue.xml -->
<config xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
        xsi:noNamespaceSchemaLocation="urn:magento:framework-message-queue:etc/etc.xsd">

    <topic name="erp.product.sync" publisher="amqp">
        <subscription name="erp.product.sync.default"
                       queue="erp.product.sync"
                       handler="Vendor\Integration\Model\Queue\Consumer\ProductSyncHandler::process"/>
    </topic>

    <topic name="erp.order.push" publisher="amqp">
        <subscription name="erp.order.push.default"
                       queue="erp.order.push"
                       handler="Vendor\Integration\Model\Queue\Consumer\OrderPushHandler::process"/>
    </topic>
</config>

Producer

namespace Vendor\Integration\Model\Queue\Producer;

class ProductSyncProducer
{
    private TopicPublisherInterface $publisher;

    public function publish(array $skuList): void
    {
        foreach ($skuList as $sku) {
            $this->publisher->publish(
                'erp.product.sync',
                json_encode([
                    'sku' => $sku,
                    'timestamp' => time(),
                ])
            );
        }
    }
}

Message Processing

Consumer Handler

namespace Vendor\Integration\Model\Queue\Consumer;

use Magento\Framework\MessageQueue\Consumer\HandlerInterface;

class ProductSyncHandler implements HandlerInterface
{
    private ErpAdapterInterface $erpAdapter;
    private ProductRepositoryInterface $productRepo;
    private LoggerInterface $logger;

    public function process(string $messageBody): void
    {
        $data = json_decode($messageBody, true);
        $sku = $data['sku'];

        try {
            // Pull from ERP
            $erpData = $this->erpAdapter->syncProduct($sku);

            // Update Magento product
            $product = $this->productRepo->get($sku);
            $product->setName($erpData->getName());
            $product->setPrice($erpData->getPrice());
            $this->productRepo->save($product);

            $this->logger->info('Product synced', ['sku' => $sku]);
        } catch (\Exception $e) {
            $this->logger->error('Product sync failed', [
                'sku' => $sku,
                'error' => $e->getMessage(),
            ]);
            throw $e; // Re-throw for retry
        }
    }
}

Consumer Configuration

<!-- app/code/Vendor/Integration/etc/queue_consumer.xml -->
<config xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
        xsi:noNamespaceSchemaLocation="urn:magento:framework-message-queue:etc/consumer.xsd">

    <consumer name="erp.product.sync"
              queue="erp.product.sync"
              connection="amqp"
              handler="Vendor\Integration\Model\Queue\Consumer\ProductSyncHandler::process"
              maxMessages="100"
              maxIdleTime="0"
              sleep="0"/>
</config>

Queue Patterns

Priority Queue

namespace Vendor\Integration\Model\Queue\Priority;

class PriorityProducer
{
    public function publishWithPriority(string $topic, array $payload, int $priority): void
    {
        $this->publisher->publish($topic, json_encode([
            'payload' => $payload,
            'priority' => $priority,
            'timestamp' => time(),
        ]), [
            'priority' => $priority,
        ]);
    }
}

Fan-out Pattern

<!-- Fan-out: one topic, multiple queues -->
<config xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
        xsi:noNamespaceSchemaLocation="urn:magento:framework-message-queue:etc/etc.xsd">

    <topic name="order.created" publisher="amqp">
        <subscription name="erp.order.sync"
                       queue="erp.order.sync"
                       handler="...\ErpOrderHandler::process"/>
        <subscription name="crm.order.sync"
                       queue="crm.order.sync"
                       handler="...\CrmOrderHandler::process"/>
        <subscription name="notification.order.created"
                       queue="notification.order.created"
                       handler="...\NotificationHandler::process"/>
    </topic>
</config>

Delayed Message

namespace Vendor\Integration\Model\Queue\Delay;

class DelayedProducer
{
    public function publishDelayed(string $topic, array $payload, int $delaySeconds): void
    {
        $this->publisher->publish($topic, json_encode($payload), [
            'delay' => $delaySeconds,
        ]);
    }
}

Queue Monitoring

Queue Monitor

namespace Vendor\Integration\Model\Queue\Monitor;

class QueueMonitor
{
    private ManagementInterface $management;
    private AlertInterface $alert;

    public function checkQueues(): void
    {
        $queues = $this->management->getQueues();

        foreach ($queues as $queue) {
            $size = $queue->getSize();
            $consumerCount = $queue->getConsumers()->count();

            // Alert on queue backlog
            if ($size > 1000 && $consumerCount === 0) {
                $this->alert->send(
                    'Queue backlog warning',
                    sprintf(
                        'Queue %s has %d messages with no consumers',
                        $queue->getName(),
                        $size
                    )
                );
            }
        }
    }
}

Queue Statistics

namespace Vendor\Integration\Model\Queue\Stats;

class QueueStats
{
    private array $stats = [];

    public function recordProcessed(string $queueName, int $processingTime): void
    {
        if (!isset($this->stats[$queueName])) {
            $this->stats[$queueName] = [
                'total' => 0,
                'failed' => 0,
                'avg_time' => 0,
            ];
        }

        $this->stats[$queueName]['total']++;
        $this->stats[$queueName]['avg_time'] =
            ($this->stats[$queueName]['avg_time'] * ($this->stats[$queueName]['total'] - 1) + $processingTime)
            / $this->stats[$queueName]['total'];
    }
}

Quiz

1. What is the benefit of queue-based integration?

Question 1 options

2. What is a fan-out pattern?

Question 2 options

3. What happens when a consumer fails?

Question 3 options

Flashcards

Question

What is a message queue?

Answer

Buffer between producers and consumers for async processing

Question

What is fan-out?

Answer

One message delivered to multiple queues/consumers

Question

What is a dead letter queue?

Answer

Queue for messages that failed all retry attempts

Question

What is consumer lag?

Answer

When queue grows faster than consumers can process

Revision Notes

Key Takeaways

  • 1. Queues decouple producers and consumers for async processing
  • 2. Fan-out pattern routes messages to multiple consumers
  • 3. Dead letter queues handle permanently failed messages
  • 4. Queue monitoring detects backlogs and consumer failures
  • 5. Priority queues ensure critical messages are processed first

Interview Tips

  • Explain when to use sync vs async processing
  • Discuss fan-out vs point-to-point patterns
  • Describe handling consumer failures gracefully
  • Talk about queue monitoring and alerting

Cheat Sheet

Queue Patterns:
  Point-to-Point → one consumer per message
  Fan-out → multiple consumers per message
  Priority → order by importance
  Delayed → schedule future delivery

Components:
  Producer → creates messages
  Queue → stores messages
  Consumer → processes messages
  Dead Letter → failed messages

Monitoring:
  Queue size → detect backlogs
  Consumer count → detect idle
  Processing time → detect slow consumers