Weave Code
Code Weaver
Helps Laravel developers discover, compare, and choose open-source packages. See popularity, security, maintainers, and scores at a glance to make better decisions.
Feedback
Share your thoughts, report bugs, or suggest improvements.
Subject
Message

Messenger Kafka Laravel Package

den1008/messenger-kafka

View on GitHub
Deep Wiki
Context7

Getting Started

Minimal Setup

  1. Installation:

    composer require koco/messenger-kafka
    

    For non-Flex projects, add Koco\Kafka\KocoKafkaBundle::class to config/bundles.php.

  2. Configure DSN in config/packages/messenger.yaml:

    messenger:
        transports:
            kafka:
                dsn: '%env(KAFKA_DSN)%'  # e.g., kafka://localhost:9092
                options:
                    topic: 'your_topic_name'
                    group_id: 'your_consumer_group'
    
  3. Dispatch a message:

    use Symfony\Component\Messenger\MessageBusInterface;
    
    $bus->dispatch(new YourMessage());
    
  4. Consume messages (via Symfony CLI):

    php bin/console messenger:consume kafka -vv
    

First Use Case

Replace a default transport (e.g., async) with Kafka for high-throughput messaging:

framework:
    messenger:
        transports:
            async: ~
            kafka:
                dsn: '%env(KAFKA_DSN)%'
                options:
                    topic: 'async_messages'
                    group_id: 'async_workers'
        routing:
            'App\Message\AsyncTask': kafka

Implementation Patterns

Core Workflows

  1. Producer Pattern:

    • Use Kafka for asynchronous, high-volume tasks (e.g., notifications, batch jobs).
    • Example: Dispatch a SendEmailMessage to Kafka instead of a database queue.
    routing:
        'App\Message\SendEmail': kafka
    
  2. Consumer Groups:

    • Scale consumers by group_id (e.g., workers_*). Each group processes messages independently.
    • Example: Deploy 3 workers with group_id: workers_backend for parallel processing.
  3. Error Handling:

    • Use failed_transport to redirect failed messages to a fallback (e.g., doctrine):
    transports:
        kafka:
            dsn: '%env(KAFKA_DSN)%'
            failed_transport: doctrine
    
  4. Serialization:

    • Default: JSON. Override via serializer option:
    options:
        serializer: 'App\Serializer\CustomSerializer'
    

Integration Tips

  • Laravel-Specific: Combine with spatie/laravel-messenger for Laravel support:

    $this->bus->dispatch(new YourMessage());
    

    Configure in config/messenger.php (Laravel’s adapter layer).

  • Monitoring: Use Kafka’s built-in tools (kafka-console-consumer) or integrate with Prometheus via custom metrics.

  • Retry Logic: Leverage Symfony’s retry middleware:

    messenger:
        transports:
            kafka:
                dsn: '%env(KAFKA_DSN)%'
                retry_strategy:
                    max_retries: 3
                    delay: 1000
    

Gotchas and Tips

Pitfalls

  1. Connection Issues:

    • Symptom: Messages silently fail to send/receive.
    • Fix: Validate DSN format (kafka://host:port) and ensure brokers are reachable.
    • Debug: Use php bin/console debug:container koco_kafka to inspect the transport.
  2. Consumer Lag:

    • Cause: Slow message processing or insufficient consumers.
    • Solution: Scale consumers or optimize message handlers.
  3. Schema Evolution:

    • Kafka has no built-in schema registry. Use a tool like Confluent Schema Registry or validate messages manually.
  4. SSL/TLS:

    • For kafka+ssl://, ensure certificates are trusted:
    options:
        ssl:
            cafile: '/path/to/ca.pem'
            certfile: '/path/to/client.pem'
            keyfile: '/path/to/client.key'
    

Debugging

  • Enable Verbose Logging:

    php bin/console messenger:consume kafka -vv
    

    Look for KafkaTransport logs in var/log/dev.log.

  • Check Topic Existence: Use kafka-topics.sh to verify topics:

    kafka-topics.sh --list --bootstrap-server localhost:9092
    

Extension Points

  1. Custom Serializers: Implement MessageSerializerInterface for non-JSON formats (e.g., Avro):

    class AvroSerializer implements MessageSerializerInterface
    {
        public function serialize($message): string { ... }
        public function deserialize(string $data): mixed { ... }
    }
    

    Register in config:

    options:
        serializer: 'App\Serializer\AvroSerializer'
    
  2. Middleware: Add Kafka-specific middleware (e.g., logging, transformation):

    $bus->addMiddleware(new KafkaLoggingMiddleware());
    
  3. Dynamic Topics: Use a message attribute to route dynamically:

    #[Message(['topic' => 'dynamic_topic'])]
    class DynamicMessage {}
    

    Requires custom transport extension.

Configuration Quirks

  • Default Topic: If omitted, the transport uses messenger.<transport_name> (e.g., messenger.kafka).
  • Group ID: Must match consumer group names. Avoid hardcoding; use environment variables:
    options:
        group_id: '%env(KAFKA_GROUP_ID)%'
    
Weaver

How can I help you explore Laravel packages today?

Conversation history is not saved when not logged in.
Prompt
Add packages to context
No packages found.
terminal42/code-quality-tools
codifyo/ts-generator-bundle
andydefer/laravel-cluster
testo/fiber
mintobit/jobqueue
a4sex/maintenance-bundle
a4sex/entity-date-update
a4sex/client-identifier
a4sex/base-utilites
a4sex/key-value-storage
a4sex/micro-status
chilldev/dependency-injection-extra
datinglibre/datinglibre-app-api
biberltd/corebundle
bricre/symfony-bundle-test
biberltd/logbundle
dominium/http-adapter-bundle
dominium/google-analytics
a4sex/auto-clean-entity
christhompsontldr/laravel-inky