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?
2. What is a fan-out pattern?
3. What happens when a consumer fails?
Flashcards
Question
What is a message queue?
Click to reveal answer
Answer
Buffer between producers and consumers for async processing
Question
What is fan-out?
Click to reveal answer
Answer
One message delivered to multiple queues/consumers
Question
What is a dead letter queue?
Click to reveal answer
Answer
Queue for messages that failed all retry attempts
Question
What is consumer lag?
Click to reveal answer
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