You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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集群的连通性。
  • 补全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需显式指定这些参数(部分连接器会忽略全局配置)。

  • 显式配置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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.17 02:17:15