Scala通过Kafka Java API对接AWS MSK IAM认证消费者报错求助
排查解决步骤
- 验证基础网络与端口配置:确认代码中填写的MSK bootstrap地址、端口正确,IAM认证的MSK默认端口为9098。若代码在本地运行,需确认MSK已开启公网访问权限,且集群安全组已放开你本地出口IP的9098端口访问权限;若代码部署在AWS VPC内,需确认部署资源与MSK在同一VPC或已打通跨VPC网络。额外注意你贴出的Conduktor配置中
sasl.client.callback.handler.class类名拼写少了末尾的r,如果是粘贴失误不影响使用,若实际配置确实写错需要修正。 - 核对AWS凭证与权限:确认代码运行环境使用的AWS凭证,和Conduktor/Kafka CLI使用的凭证一致,且对应身份已配置MSK访问权限,必须包含
kafka-cluster:Connect权限,消费场景还需要额外配置kafka-cluster:DescribeTopic、kafka-cluster:ReadData、kafka-cluster:AlterConsumerGroup等主题、消费组相关权限。 - 排查依赖冲突:你使用的aws-msk-iam-auth 1.1.0依赖与Kafka 2.4.1版本本身兼容,可通过
sbt dependencyTree命令检查项目依赖,确认是否存在AWS SDK版本冲突:如果项目中其他依赖引入了不同版本的AWS SDK包,会导致IAM认证逻辑异常,无法完成集群握手,最终消费者无法正常获取分区分配。 - 调整消费者调用逻辑:你遇到的
java.lang.IllegalStateException: You can only check the position for partitions assigned to this consumer报错,本质是在消费者未完成分区分配时就调用了position()、seek()等需要分区已绑定到当前消费者的方法。可在消费者订阅主题后,先调用一次consumer.poll(Duration.ofSeconds(5))触发服务端的分区分配流程,再执行分区位置相关操作,验证是否为代码逻辑顺序问题。 - 检查遗漏配置项:确认你代码中省略的配置部分包含必填的
group.id、bootstrap.servers配置,且sasl.jaas.config等安全配置末尾没有多余的空格、特殊字符。
内容的提问来源于stack exchange,提问作者l33tHax0r
相关产品推荐
相关产品推荐

