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 Kinesis Laravel Package

cvek/messenger-kinesis

View on GitHub
Deep Wiki
Context7

Getting Started

Minimal Setup

  1. Installation

    composer require cvek/messenger-kinesis
    
  2. Configure Kinesis Transport Add to your config/packages/symfony_messenger.php:

    'transports' => [
        'kinesis' => [
            'dsn' => 'aws://access_key:secret_key@?region=us-east-1&stream=kinesis-stream-name',
            'options' => [
                'retry_after' => 100, // ms
                'batch_size' => 100,
            ],
        ],
    ],
    
  3. First Message Dispatch

    use Symfony\Component\Messenger\MessageBusInterface;
    
    public function __construct(private MessageBusInterface $bus) {}
    
    public function sendMessage()
    {
        $this->bus->dispatch(new YourMessageClass());
    }
    
  4. Verify Stream Check AWS Kinesis console for the stream and messages.


Implementation Patterns

Common Workflows

1. Dispatching Messages

  • Synchronous Dispatch (default):
    $this->bus->dispatch(new OrderCreatedEvent($orderId));
    
  • Asynchronous Dispatch (with transport-specific options):
    $this->bus->dispatch(new OrderCreatedEvent($orderId), ['transport' => 'kinesis']);
    

2. Handling Failures

  • Retry Logic: Configure retry_after in options to control backoff.
  • Dead Letter Queue (DLQ): Extend the transport to route failed messages to a DLQ stream:
    // In a custom transport class
    public function send(Envelope $envelope): Envelope
    {
        try {
            return parent::send($envelope);
        } catch (Exception $e) {
            $this->sendToDLQ($envelope, $e);
            throw $e; // Re-throw to trigger retry logic
        }
    }
    

3. Batch Processing

  • Leverage batch_size to optimize throughput:
    'options' => [
        'batch_size' => 500, // Adjust based on AWS limits
    ],
    

4. Message Serialization

  • Custom Serializer: Override the default serializer (e.g., for large payloads):
    'options' => [
        'serializer' => new CustomSerializer(), // Implement Symfony\Component\Serializer\SerializerInterface
    ],
    

Integration Tips

With Laravel

  • Service Provider Binding:
    public function register()
    {
        $this->app->extend(MessageBus::class, function ($bus, $app) {
            $bus->addTransport('kinesis', new KinesisTransport(
                new KinesisClient([
                    'region'  => 'us-east-1',
                    'version' => 'latest',
                    'credentials' => [
                        'key'    => env('AWS_ACCESS_KEY_ID'),
                        'secret' => env('AWS_SECRET_ACCESS_KEY'),
                    ],
                ]),
                new Serializer(),
                new OptionsResolver()
            ));
            return $bus;
        });
    }
    

With Symfony Messenger Middleware

  • Logging Middleware:
    $bus->addMiddleware(new HandleMessageBus(new KinesisLogger()));
    
    class KinesisLogger implements MessageBusInterface
    {
        public function handle(Envelope $envelope, callable $next): Envelope
        {
            \Log::debug('Dispatching to Kinesis', ['message' => $envelope->getMessage()]);
            return $next($envelope);
        }
    }
    

With AWS SDK

  • Custom Kinesis Client: Inject a pre-configured AWS SDK client:
    $client = new KinesisClient([
        'region' => 'us-west-2',
        'endpoint' => 'https://kinesis.us-west-2.amazonaws.com',
    ]);
    $transport = new KinesisTransport($client);
    

Gotchas and Tips

Pitfalls

1. Credential Management

  • Hardcoded Credentials: Avoid committing AWS credentials to version control. Use environment variables or AWS IAM roles.
    // Bad: Hardcoded
    'dsn' => 'aws://AKIAEXAMPLE:secret@?region=us-east-1'
    
    // Good: Environment variables
    'dsn' => 'aws://' . env('AWS_ACCESS_KEY_ID') . ':' . env('AWS_SECRET_ACCESS_KEY') . '@?region=' . env('AWS_REGION')
    

2. Stream Provisioned Throughput

  • Throttling: Exceeding PutRecord limits (1MB/s or 1,000 records/s) causes ProvisionedThroughputExceededException.
    • Solution: Monitor CloudWatch metrics and adjust shard count or batch size.

3. Message Ordering

  • Partition Key: Kinesis orders messages by partition key. Ensure consistent keys for ordered processing:
    // In your message class
    public function getPartitionKey(): string
    {
        return 'user_' . $this->userId; // Consistent key for ordering
    }
    

4. Serialization Issues

  • Non-Serializable Objects: Default serializer fails on closures, resources, or circular references.
    • Solution: Implement Symfony\Component\Serializer\NormalizerInterface or use ignore_errors:
      'options' => [
          'serializer' => new Serializer([], [new IgnoreExceptionNormalizer()]),
      ],
      

5. Retry Logic

  • Infinite Retries: Without a DLQ, messages may retry indefinitely.
    • Solution: Configure max_retries in middleware or use a custom transport:
      'options' => [
          'max_retries' => 3, // Add this to OptionsResolver
      ],
      

Debugging Tips

1. Enable AWS SDK Debugging

$client = new KinesisClient([
    'region' => 'us-east-1',
    'debug' => true, // Enable debug logging
]);

2. Check Kinesis Metrics

  • CloudWatch: Monitor IncomingBytes, IncomingRecords, and IteratorAgeMilliseconds.
  • Stream Console: Use AWS CLI to inspect messages:
    aws kinesis get-shard-iterator --shard-id shardId-000000000000 --shard-iterator-type TRIM_HORIZON --stream-name your-stream
    aws kinesis get-records --shard-iterator <iterator-from-above>
    

3. Local Testing

  • LocalStack: Spin up a local Kinesis mock:
    docker run -it -p 4566:4566 localstack/localstack
    
    Configure transport to use LocalStack endpoint:
    'dsn' => 'aws://access_key:secret_key@?region=us-east-1&endpoint=http://localhost:4566&stream=test-stream'
    

Extension Points

1. Custom Transport Class

Extend Cvek\MessengerKinesis\Transport to add features:

class CustomKinesisTransport extends KinesisTransport
{
    public function send(Envelope $envelope): Envelope
    {
        // Add pre-send logic (e.g., logging, validation)
        return parent::send($envelope);
    }

    public function getName(): string
    {
        return 'custom_kinesis';
    }
}

2. Message Metadata

Add custom metadata to Kinesis records:

$transport = new KinesisTransport($client, $serializer, $options);
$transport->setMetadataExtractor(function (Envelope $envelope) {
    return [
        'user_id' => $envelope->getMessage()->getUserId(),
        'timestamp' => now()->toIso8601String(),
    ];
});

3. Dynamic Stream Selection

Route messages to different streams based on conditions:

$transport = new KinesisTransport($client, $serializer, $options);
$transport->setStreamResolver(function (Envelope $envelope) {
    return $envelope->getMessage() instanceof HighPriorityMessage
        ? 'high-priority-stream'
        : 'default-stream';
});

4. Event Listeners for Kinesis

Listen to Kinesis events (e.g., stream creation, throttling):

$client->getStream('your-stream')->subscribe(function (StreamEvent $event) {
    if ($event->isThrottled()) {
        \Log::warning('Kinesis throttling detected');
    }
});
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