SEPEHR.SYS

Starting SEPEHR.SYS…

Loading

سپهر محسنی

مهندس نرم‌افزار فول‌استک و هوش مصنوعی

NOTES.TXT ×

لاراول توزیع‌شده: معماری‌های رویدادمحور با کارایی بالا

3′ · ۱۷ دی ۱۴۰۴

در محیط‌های توزیع‌شده، رویکرد «Fire and Forget» دستورالعملی تضمینی برای خرابی داده‌هاست. معماری پیشرفتهٔ لاراول مستلزم حل مسئلهٔ Atomic Commitment است: یعنی اطمینان از اینکه پایگاه داده و Message Broker شما بر سر وضعیت دنیا هم‌نظر باشند، بدون آنکه کارایی یا سازگاری (Consistency) قربانی شود.

۱. تعریف قرارداد: الگوی Envelope

استفاده از payload به‌شکل «آرایهٔ خام»، یک نقطه‌ضعف معماری است. سیستم‌های پیشرفته از یک Message Envelope استفاده می‌کنند که متادیتای مربوط به ردیابی توزیع‌شده (OpenTelemetry)، نسخه‌بندی Schema و Idempotency را در خود کپسوله می‌کند.

<?php

namespace App\Core\Messaging;

use JsonSerializable;
use Illuminate\Support\Str;

readonly class MessageEnvelope implements JsonSerializable
{
    public function __construct(
        public string $id,
        public string $topic,
        public array $payload,
        public string $version = '1.0',
        public array $metadata = [],
        public float $occurredAt = 0.0,
    ) {
        $this->occurredAt = $this->occurredAt ?: microtime(true);
    }

    public static function create(string $topic, array $payload): self
    {
        return new self(
            id: Str::uuid()->toString(),
            topic: $topic,
            payload: $payload,
            metadata: [
                'correlation_id' => request()->header('X-Correlation-ID') ?? Str::uuid()->toString(),
            ]
        );
    }

    public function jsonSerialize(): array
    {
        return get_object_vars($this);
    }
}

۲. پیاده‌سازی درایور در سطح پایین

برای دستیابی به دقت ۱۰۰٪، از انتزاع پیش‌فرض Queue در لاراول فراتر می‌رویم و مستقیماً با قابلیت‌های سطح پروتکل Brokerهای خود کار می‌کنیم.

Redis Streams (مسیر توان عملیاتی بالا)

به‌جای لیست‌های استاندارد Redis، از Redis Streams (XADD) استفاده می‌کنیم. این قابلیت، ماندگاری بومی پیام‌ها و Consumer Groupها را فراهم می‌کند و به لاراول امکان می‌دهد جهش‌های عظیم رویدادها را بدون از دست رفتن داده مدیریت کند.

namespace App\Infrastructure\Messaging\Drivers;

use App\Core\Messaging\MessageEnvelope;
use Illuminate\Support\Facades\Redis;

class RedisStreamDriver
{
    public function publish(MessageEnvelope $envelope): void
    {
        // Use MAXLEN ~ 10000 to implement an eviction policy
        Redis::connection('messaging')->xadd(
            $envelope->topic, 
            'MAXLEN', '~', 10000, 
            '*', 
            ['data' => json_encode($envelope)]
        );
    }
}

RabbitMQ (مسیر قابل‌اعتماد)

برای RabbitMQ، مکانیزم Publisher Confirms را پیاده‌سازی می‌کنیم. برخلاف Redis، این پیاده‌سازی فرایند PHP را وادار می‌کند پیش از علامت‌گذاری پیام به‌عنوان «ارسال‌شده»، منتظر دریافت تأییدیه (ACK) از سمت Broker مربوط به RabbitMQ بماند.

namespace App\Infrastructure\Messaging\Drivers;

use PhpAmqpLib\Message\AMQPMessage;
use PhpAmqpLib\Exchange\AMQPExchangeType;

class RabbitMQDriver
{
    public function publish(MessageEnvelope $envelope): void
    {
        $channel = app('amqp.connection')->channel();
        
        $channel->exchange_declare($envelope->topic, AMQPExchangeType::TOPIC, false, true, false);

        $msg = new AMQPMessage(json_encode($envelope), [
            'delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT,
            'content_type' => 'application/json',
        ]);

        $channel->confirm_select(); 
        $channel->basic_publish($msg, $envelope->topic);
        
        // Wait for ACK with a 5s timeout
        $channel->wait_for_pending_acks(5.0); 
        
        $channel->close();
    }
}

۳. Outbox تراکنشی: ارسال بدون نشتی

رایج‌ترین نقطهٔ شکست، Dispatch کردن یک رویداد در داخل تراکنش پایگاه داده است. اگر تراکنش Rollback شود اما رویداد از قبل ارسال شده باشد، سیستم شما وارد وضعیتی ناسازگار می‌شود. ما این مشکل را با نوشتن در یک Outbox درون همان تراکنش ACID حل می‌کنیم.

Relay Worker با کارایی بالا

یک Relay خام، مصرف CPU را به اوج می‌رساند. ما از Row-Level Locking (lockForUpdate) استفاده می‌کنیم تا چندین نمونهٔ Relay بتوانند بدون پردازش تکراری، به‌صورت افقی مقیاس بگیرند.

namespace App\Console\Commands;

use App\Models\OutboxMessage;
use App\Infrastructure\Messaging\MessagingMapper;
use Illuminate\Console\Command;

class OutboxRelay extends Command
{
    protected $signature = 'messaging:relay {--driver=redis}';

    public function handle(): void
    {
        while (true) {
            $processed = OutboxMessage::whereNull('dispatched_at')
                ->where('retry_count', '<', 5)
                ->orderBy('created_at', 'asc')
                ->limit(100)
                ->lockForUpdate() 
                ->get()
                ->each(function ($msg) {
                    try {
                        $driver = MessagingMapper::get($this->option('driver'));
                        $driver->publish($msg->toEnvelope());
                        
                        $msg->update(['dispatched_at' => now()]);
                    } catch (\Throwable $e) {
                        $msg->increment('retry_count');
                        report($e);
                    }
                });

            $processed->isEmpty() ? sleep(1) : usleep(100000);
        }
    }
}

۴. Idempotency در Consumer: توهم «دقیقاً یک بار»

در سیستم‌های توزیع‌شده نمی‌توانید تضمین کنید که یک پیام فقط یک بار ارسال شود، اما می‌توانید تضمین کنید که فقط یک بار پردازش شود.

namespace App\Core\Messaging;

use Illuminate\Support\Facades\DB;

trait InteractsWithIdempotency
{
    public function executeOnce(string $messageId, callable $action): void
    {
        try {
            DB::table('processed_events')->insert([
                'id' => $messageId,
                'processed_at' => now(),
            ]);
            
            $action();
        } catch (\Illuminate\Database\QueryException $e) {
            if ($e->getCode() === '23000') {
                return;
            }
            throw $e;
        }
    }
}

مقایسهٔ نهایی

| ویژگی | Redis Streams | RabbitMQ (AMQP) | |--------------|--------------------------------|----------------------------------| | توان عملیاتی (Throughput) | ~150k msg/s | ~40k msg/s | | تأیید (Acknowledge) | سطح Consumer | سطح Broker و Consumer | | ماندگاری (Persistence) | درون‌حافظه‌ای / AOF | مبتنی بر دیسک | | بهترین کاربرد | فیدهای بلادرنگ، متریک‌ها | تراکنش‌های مالی، سفارش‌ها |

این یادداشت ترجمهٔ فارسی نوشتهٔ خودم است — نسخهٔ اصلی (انگلیسی) در dev.to