PHP code example of kafka / workflow

1. Go to this page and download the library: Download kafka/workflow library. Choose the download type require.

2. Extract the ZIP file and open the index.php.

3. Add this code to the index.php.
    
        
<?php
require_once('vendor/autoload.php');

/* Start to develop here. Best regards https://php-download.com/ */

    

kafka / workflow example snippets


use Wf\Kafka\Producer\OutboxWriter;
use Wf\Kafka\Producer\ProducerOptions;

DB::transaction(function () use ($order) {
    // Business logic của bạn
    $order = Order::create($orderData);

    // Ghi vào outbox — cùng transaction với business logic
    OutboxWriter::write(
        topic:     'order-events',
        eventType: 'ORDER_CREATED',
        payload:   $order->toArray(),
        options:   ProducerOptions::make()->withKey((string) $order->id)
    );
});

DB::transaction(function () use ($orders) {
    // ...
    $payloads = collect($orders)->map->toArray()->all();

    OutboxWriter::writeBatch(
        topic:     'order-events',
        eventType: 'ORDER_CREATED',
        payloads:  $payloads
    );
    // Trả về: int — số bản ghi đã ghi
});

DB::transaction(function () {
    OutboxWriter::writeMulti([
        ['topic' => 'order-events',    'event_type' => 'ORDER_CREATED',   'payload' => [...]],
        ['topic' => 'payment-events',  'event_type' => 'PAYMENT_CHARGED', 'payload' => [...]],
        ['topic' => 'inventory-events','event_type' => 'STOCK_RESERVED',  'payload' => [...]],
    ]);
    // Trả về: int — số bản ghi đã ghi
});

$options = ProducerOptions::make()
    ->withKey('order-123')               // Partition key — đảm bảo ordering
    ->withHeaders(['source' => 'api'])   // Kafka message headers
    ->withFlushTimeoutMs(3000)           // Override flush timeout (ms)
    ->withPartition(2);                  // Gửi vào partition cụ thể (hiếm khi cần)

// app/Console/Kernel.php
$schedule->command('kafka:outbox-relay')->everyFiveSeconds()->withoutOverlapping();

use Wf\Kafka\Consumer\ConsumerDispatcher;

app(ConsumerDispatcher::class)
    ->onTopic('order-events')
    ->withGroup('order-service-group')
    ->handle(function (array $payload, array $meta) {
        // $payload — dữ liệu gốc đã deserialize hoàn toàn
        // $meta    — ['event_id', 'event_type', 'topic', 'partition', 'offset']

        OrderService::process($payload);
    })
    ->listen(); // blocking loop

use Wf\Kafka\Exceptions\TransientInfraException;
use Wf\Kafka\Exceptions\PoisonPillException;

->handle(function (array $payload, array $meta) {
    // Lỗi tạm thời (mạng, API timeout, DB down)
    // → Package route vào DLQ, KHÔNG commit offset
    // → Message được re-deliver khi consumer restart
    if ($apiIsDown) {
        throw new TransientInfraException("Payment API timeout sau 3s.");
    }

    // Dữ liệu không thể xử lý được (Poison Pill)
    // → Package route vào EDL, commit offset (không block queue)
    if (!in_array($payload['status'], ['PENDING', 'COMPLETED'])) {
        throw new PoisonPillException("Unknown status: {$payload['status']}");
    }

    // Bất kỳ exception nào khác cũng → EDL
    OrderService::process($payload);
})

->onError(function (string $type, array $envelope, \Throwable $e) {
    // $type = 'dlq' | 'edl'
    // Dùng để notify thêm, ghi log riêng, alert Sentry...
    Sentry::captureException($e);
})

// app/Console/Commands/OrderConsumerCommand.php
namespace App\Console\Commands;

use Illuminate\Console\Command;
use Wf\Kafka\Consumer\ConsumerDispatcher;

class OrderConsumerCommand extends Command
{
    protected $signature   = 'consumer:orders';
    protected $description = 'Consume order-events topic';

    public function handle(ConsumerDispatcher $dispatcher): void
    {
        $dispatcher
            ->onTopic('order-events')
            ->withGroup('order-service-group')
            ->handle(function (array $payload, array $meta) {
                OrderService::process($payload);
            })
            ->onError(function (string $type, array $envelope, \Throwable $e) {
                Log::channel('slack')->critical("[$type] {$e->getMessage()}");
            })
            ->listen();
    }
}

// config/kafka.php
'serialization' => [
    'driver' => 'avro',
    'avro' => [
        'schemas' => [
            'order-events'   => base_path('avro/order.avsc'),
            'payment-events' => base_path('avro/payment.avsc'),
        ],
    ],
],
bash
php artisan kafka:install          # Detect OS → cài tự động
php artisan kafka:install --check  # Chỉ kiểm tra, không cài
php artisan kafka:install --force  # Không hỏi confirm
dockerfile
FROM php:8.2-fpm-alpine AS base

RUN apk add --no-cache librdkafka librdkafka-dev $PHPIZE_DEPS \
    && pecl install rdkafka \
    && docker-php-ext-enable rdkafka \
    && apk del --no-cache $PHPIZE_DEPS librdkafka-dev \
    && rm -rf /tmp/pear /var/cache/apk/*
bash
# Publish config vào config/kafka.php của project
php artisan vendor:publish --tag=kafka-config

# Chạy migrations (tạo 3 bảng kafka_*)
php artisan migrate
bash
php artisan vendor:publish --tag=kafka-migrations