Installation
composer require cvek/messenger-kinesis
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,
],
],
],
First Message Dispatch
use Symfony\Component\Messenger\MessageBusInterface;
public function __construct(private MessageBusInterface $bus) {}
public function sendMessage()
{
$this->bus->dispatch(new YourMessageClass());
}
Verify Stream Check AWS Kinesis console for the stream and messages.
$this->bus->dispatch(new OrderCreatedEvent($orderId));
$this->bus->dispatch(new OrderCreatedEvent($orderId), ['transport' => 'kinesis']);
retry_after in options to control backoff.// 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
}
}
batch_size to optimize throughput:
'options' => [
'batch_size' => 500, // Adjust based on AWS limits
],
'options' => [
'serializer' => new CustomSerializer(), // Implement Symfony\Component\Serializer\SerializerInterface
],
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;
});
}
$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);
}
}
$client = new KinesisClient([
'region' => 'us-west-2',
'endpoint' => 'https://kinesis.us-west-2.amazonaws.com',
]);
$transport = new KinesisTransport($client);
// 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')
PutRecord limits (1MB/s or 1,000 records/s) causes ProvisionedThroughputExceededException.
// In your message class
public function getPartitionKey(): string
{
return 'user_' . $this->userId; // Consistent key for ordering
}
Symfony\Component\Serializer\NormalizerInterface or use ignore_errors:
'options' => [
'serializer' => new Serializer([], [new IgnoreExceptionNormalizer()]),
],
max_retries in middleware or use a custom transport:
'options' => [
'max_retries' => 3, // Add this to OptionsResolver
],
$client = new KinesisClient([
'region' => 'us-east-1',
'debug' => true, // Enable debug logging
]);
IncomingBytes, IncomingRecords, and IteratorAgeMilliseconds.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>
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'
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';
}
}
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(),
];
});
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';
});
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');
}
});
How can I help you explore Laravel packages today?