PHP code example of kafka-bus / laravel-bridge

1. Go to this page and download the library: Download kafka-bus/laravel-bridge 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-bus / laravel-bridge example snippets


'default' => env('KAFKA_CONNECTION', 'kafka'),

'connections' => [
    'kafka' => [
        'driver' => 'kafka',
        'options' => [
            'metadata.broker.list' => env('KAFKA_BROKER_LIST', 'localhost:9092'),
            'security.protocol'    => env('KAFKA_SECURITY_PROTOCOL', 'SASL_PLAINTEXT'),
            'sasl.mechanisms'      => env('KAFKA_SASL_MECHANISMS', 'PLAIN'),
            'sasl.username'        => env('KAFKA_SASL_USERNAME'),
            'sasl.password'        => env('KAFKA_SASL_PASSWORD'),
            'debug'                => env('KAFKA_DEBUG', false),
        ],
    ],

    'testing' => [
        'driver'  => 'null',
        'options' => [],
    ],
],

'topic_prefix' => env('KAFKA_PREFIX', env('APP_ENV', 'local').'.'),

'topics' => [
    'products' => 'fact.products.1',
    'orders'   => 'fact.orders.1',
],

'producers' => [
    'middleware' => [
        // KafkaBus\Commiter\Middleware\ProducerIdempotencyMiddleware::class,
    ],

    'routes' => [
        App\Kafka\Messages\ProductMessage::class => 'products',
    ],

    'flush_timeout' => 5000,
    'flush_retries' => 5,

    'additional_options' => [
        'compression.codec' => env('KAFKA_PRODUCER_COMPRESSION_CODEC', 'snappy'),
    ],
],

'routes' => [
    App\Kafka\Messages\ProductMessage::class => [
        'topic_key'          => 'products',
        'middleware'         => [App\Kafka\Middleware\AuditTrailMiddleware::class],
        'additional_options' => [],
        'flush_timeout'      => 5000,
        'flush_retries'      => 5,
    ],
],

use KafkaBus\Core\Interfaces\Bus\BusInterface;

public function execute(BusInterface $bus): void
{
    $bus->publish(new \App\Kafka\Messages\ProductMessage(/* ... */));
}

use KafkaBus\Laravel\Facades\KafkaBus;

KafkaBus::publish(new \App\Kafka\Messages\ProductMessage(/* ... */));

'consumers' => [
    'middleware' => [
        // KafkaBus\Commiter\Middleware\ConsumerCommiterMiddleware::class,
    ],

    'workers' => [
        // Multi-topic worker with per-worker overrides
        'default' => [
            'middleware'   => [],
            'auto_commit'  => false,
            'consume_timeout' => 20000,
            'topics' => [
                'products' => App\Kafka\Consumers\ProductsTopicConsumer::class,
                'orders'   => [
                    'handler'    => App\Kafka\Consumers\OrdersTopicConsumer::class,
                    'middleware' => [App\Kafka\Middleware\TenantContextMiddleware::class],
                ],
            ],
        ],

        // Single-topic worker, worker name == topic key
        'products' => App\Kafka\Consumers\ProductsTopicConsumer::class,

        // Single-topic worker with overrides, worker name == topic key
        'orders' => [
            'middleware' => [],
            'handler'    => App\Kafka\Consumers\OrdersTopicConsumer::class,
        ],

        // Single-topic worker where worker name != topic key
        'products-secondary' => [
            'topic_key'  => 'products',
            'middleware' => [],
            'handler'    => App\Kafka\Consumers\ProductsTopicConsumer::class,
        ],
    ],

    'auto_commit'     => env('KAFKA_CONSUMER_AUTO_COMMIT', false),
    'consume_timeout' => 5_000,

    'additional_options' => [
        'group.id'              => env('KAFKA_CONSUMER_GROUP_ID', env('APP_NAME')),
        'max.poll.interval.ms'  => env('KAFKA_MAX_POLL_INTERVAL_MS', 300_000),
        'session.timeout.ms'    => env('KAFKA_SESSION_TIMEOUT_MS', 45_000),
        'heartbeat.interval.ms' => env('KAFKA_HEARTBEAT_INTERVAL_MS', 3_000),
        'auto.offset.reset'     => 'beginning',
    ],
],

// config/kafka-bus-commiter.php
return [
    'connection' => env('KAFKA_COMMITER_CONNECTION'),
    'table'      => 'kafka_bus_commits',
    'repository' => env('KAFKA_COMMITER_REPOSITORY', 'idempotency'),

    'repositories' => [
        'idempotency' => \KafkaBus\Commiter\Repositories\IdempotencyMessageRepository::class,
        'native'      => \KafkaBus\Commiter\Repositories\NativeMessageRepository::class,
    ],
];

'consumers' => [
    'middleware' => [
        \KafkaBus\Commiter\Middleware\ConsumerCommiterMiddleware::class,
    ],

    'workers' => [
        'orders' => [
            'middleware' => [
                \KafkaBus\Commiter\Middleware\ConsumerCommiterMiddleware::class,
            ],
            'handler' => App\Kafka\Consumers\OrdersTopicConsumer::class,
        ],
    ],
],

use KafkaBus\Commiter\Interfaces\HasIdempotency;
use KafkaBus\Core\Messages\ProducerMessage;

final class ProductMessage extends ProducerMessage implements HasIdempotency
{
    public function __construct(private string $productId) {}

    public function getIdempotencyKey(): string
    {
        return $this->productId;
    }
}

'producers' => [
    'middleware' => [
        \KafkaBus\Commiter\Middleware\ProducerIdempotencyMiddleware::class,
    ],
],

use KafkaBus\Laravel\Facades\KafkaBus;

KafkaBus::fake();

use App\Kafka\Messages\ProductMessage;
use KafkaBus\Laravel\Facades\KafkaBus;

it('publishes a product message', function () {
    KafkaBus::fake();

    app(CreateProductAction::class)->execute(productId: 1);

    KafkaBus::assertPublished(ProductMessage::class);
});

KafkaBus::assertPublished(
    ProductMessage::class,
    fn($msg) => str_contains($msg->payload, '"id":1')
        && isset($msg->headers['x-idempotency-key'])
);

// Assert published exactly N times
KafkaBus::assertPublishedTimes(ProductMessage::class, 2);

// Assert a specific message was NOT published
KafkaBus::assertNotPublished(ProductMessage::class);

// Assert no messages were published at all
KafkaBus::assertNothingPublished();

// list<ProducerMessage> — serialised messages including payload, headers, topic
$messages = KafkaBus::getPublished(ProductMessage::class);
$all      = KafkaBus::allPublished();

use KafkaBus\Core\Testing\Consumers\MessageFactory;
use KafkaBus\Laravel\Facades\KafkaBus;

it('handles a product message', function () {
    KafkaBus::fake();

    $message = MessageFactory::for()
        ->withTopicKey('products')
        ->withHeaders(['x-idempotency-key' => 'abc-123'])
        ->make('{"id":1,"name":"Widget"}');

    KafkaBus::addMessage($message);
    KafkaBus::listen('products');

    // Assert side effects produced by the handler
    expect(Product::find(1))->not->toBeNull();
});

$factory = MessageFactory::for()->withTopicKey('products');

KafkaBus::addMessage($factory->make('{"id":1}'));
KafkaBus::addMessage($factory->make('{"id":2}'));

KafkaBus::listen('products');

KafkaBus::addMessage(
    MessageFactory::for()->withTopicKey('products')->make('{"id":1}')
);

KafkaBus::listen('products');

// Assert at least one message was committed on the topic
KafkaBus::assertCommitted('products');

// Assert with a condition on the ConsumerMessageInterface
KafkaBus::assertCommitted(
    'products',
    fn($msg) => $msg->payload() === '{"id":1}'
        && $msg->headers()['x-idempotency-key'] === 'abc-123'
);

// Assert exact count
KafkaBus::assertCommittedTimes('products', 2);

// Assert nothing was committed (e.g. before listen() is called)
KafkaBus::assertNothingCommitted();

$committed = KafkaBus::getCommitted('products'); // list<ConsumerMessageInterface>

expect($committed[0]->payload())->toBe('{"id":1}');
expect($committed[0]->headers())->toHaveKey('x-idempotency-key');
bash
php artisan vendor:publish --tag=kafka-bus
bash
php artisan vendor:publish --tag=kafka-bus-commiter
php artisan migrate
bash
php artisan kafka:consume default
bash
php artisan kafka:worker:list
bash
php artisan kafka:route:list