Hudi Sink Connector接入AWS MSK遇Broker断开无数据写入问题
排查Hudi Sink Connector对接AWS MSK无数据写入问题
核心问题分析
日志中的Bootstrap broker disconnected提示:Hudi Sink的任务进程与MSK Broker的连接存在异常。尽管连接器整体状态显示运行中,但单个任务可能已断连,导致无法消费Topic数据,进而无法写入S3。
具体排查与修复步骤
检查MSK Broker的访问配置
AWS MSK支持公有、私有两种端点:- 若Kafka Connect运行在VPC外,必须将MSK的公有端点配置为
bootstrap.servers; - 若在VPC内,需确认Connect所在子网与MSK子网路由可达,且安全组允许Connect的IP访问MSK的9092/9094端口(根据加密方式而定)。
注意:Hudi Sink的任务是独立进程,需单独验证其到MSK Broker的网络连通性,不能仅依赖Connect集群的连通性。
- 若Kafka Connect运行在VPC外,必须将MSK的公有端点配置为
补全SSL/SASL安全配置
若MSK集群启用了加密或身份认证,当前配置缺少对应安全参数:- 启用SSL加密时,需添加:
security.protocol=SSL ssl.truststore.location=/path/to/kafka.client.truststore.jks ssl.truststore.password=your-truststore-password - 使用IAM身份认证(AWS MSK常用方式)时,需添加:
security.protocol=SASL_SSL sasl.mechanism=AWS_MSK_IAM sasl.jaas.config=software.amazon.msk.auth.iam.IAMLoginModule required; sasl.client.callback.handler.class=software.amazon.msk.auth.iam.IAMClientCallbackHandler
其他Sink能正常运行,说明Connect集群可能已有全局安全配置,但Hudi Sink需显式指定这些参数(部分连接器会忽略全局配置)。
- 启用SSL加密时,需添加:
显式配置Hudi内部Kafka Consumer参数
Hudi Sink会启动独立的Kafka Consumer处理数据,该Consumer不会自动继承Connect的全局Kafka配置,需补充:hoodie.kafka.consumer.bootstrap.servers=你的MSK Broker端点 hoodie.kafka.consumer.security.protocol=与MSK匹配的协议(如SASL_SSL) hoodie.kafka.consumer.sasl.mechanism=AWS_MSK_IAM(若使用IAM认证)验证Topic权限与数据状态
确认Hudi Sink使用的IAM角色(或账号)拥有hudi-test-topic的READ权限,同时用kafka-console-consumer.sh测试该Topic是否有未消费的消息。降低任务数排查并发问题
暂时将tasks.max改为1,排除多任务并发导致的连接冲突,待单任务正常运行后再逐步调整任务数。
验证方法
修改配置后重启连接器,查看日志是否仍有Broker断开报错,同时监控S3路径是否有数据写入,或用Hudi的hudi-cli查看表元数据是否更新。
内容的提问来源于stack exchange,提问作者ennox108
相关产品推荐
相关产品推荐

