如何让本地Laravel应用连接AWS MSK Kafka集群实现生产消费?
连接本地Laravel应用到AWS MSK集群(使用mateusjunges/laravel-kafka)
前置准备
- 确保本地应用可访问MSK集群:若为MSK Serverless,需配置VPC端点或开启公网访问权限,同时安全组要允许本地IP访问9098(TLS)端口
- 安装依赖:
- AWS SDK for PHP:
composer require aws/aws-sdk-php - mateusjunges/laravel-kafka包:
composer require mateusjunges/laravel-kafka
- AWS SDK for PHP:
配置连接信息
1. 设置AWS凭证
在Laravel项目的.env文件中添加AWS认证信息(本地开发优先使用凭证,生产环境可改用IAM角色):
AWS_ACCESS_KEY_ID=你的访问密钥ID AWS_SECRET_ACCESS_KEY=你的秘密访问密钥 AWS_DEFAULT_REGION=MSK集群所在区域
2. 配置Kafka客户端
修改config/kafka.php文件,添加MSK集群专属配置:
return [ 'brokers' => env('KAFKA_BROKERS'), 'ssl' => true, 'sasl' => [ 'mechanism' => 'AWS_MSK_IAM', 'username' => env('AWS_ACCESS_KEY_ID'), 'password' => env('AWS_SECRET_ACCESS_KEY'), 'callback' => function ($args) { $signer = new Aws\Signature\SignatureV4('kafka', env('AWS_DEFAULT_REGION')); $request = new Aws\Http\Request( 'GET', 'https://' . $args['hostname'], [], '' ); $credentials = new Aws\Credentials\Credentials( env('AWS_ACCESS_KEY_ID'), env('AWS_SECRET_ACCESS_KEY') ); $signedRequest = $signer->signRequest($request, $credentials); return $signedRequest->getHeader('Authorization')[0]; }, ], // 保留其他默认配置项 ];
接着在.env中填入从MSK控制台获取的Broker列表:
KAFKA_BROKERS=b-1.xxx.xxx.kafka.你的区域.amazonaws.com:9098,b-2.xxx.xxx.kafka.你的区域.amazonaws.com:9098
生产者配置(应用1)
可直接在控制器或业务类中实现消息发送逻辑:
use Junges\Kafka\Facades\Kafka; // 发送消息到指定Topic Kafka::publishOn('你的Topic名称') ->withHeaders(['type' => 'user_notification']) ->withBodyKey('content', '来自Laravel生产者的消息') ->send();
消费者配置(应用2)
有两种方式实现消息消费:
方式1:使用内置Artisan命令
直接通过命令启动消费者:
php artisan kafka:consume --topic=你的Topic名称 --group=你的消费者组ID
方式2:自定义消费者类
创建app/Jobs/ConsumeKafkaMessages.php文件:
namespace App\Jobs; use Junges\Kafka\Contracts\KafkaConsumerMessage; use Junges\Kafka\Facades\Kafka; class ConsumeKafkaMessages { public function handle() { $consumer = Kafka::createConsumer(['你的Topic名称']) ->withConsumerGroupId('你的消费者组ID') ->withHandler(function (KafkaConsumerMessage $message) { // 自定义消息处理逻辑,例如记录日志或执行业务操作 info('收到Kafka消息:' . $message->getBody()); }) ->build(); $consumer->consume(); } }
然后启动队列任务执行消费:
php artisan queue:work --queue=kafka-consumer
关键注意事项
- MSK Serverless必须使用AWS_MSK_IAM作为SASL认证机制,不能使用普通用户名密码
- 确保本地PHP版本符合mateusjunges/laravel-kafka的要求(推荐PHP 8.0及以上)
- 若出现连接超时,优先检查VPC配置、安全组规则以及Broker地址是否正确
内容的提问来源于stack exchange,提问作者Thomas Chirwa
相关产品推荐
相关产品推荐

