Logstash通过IAM角色连接AWS MSK时遇NoClassDefFoundError问题
Logstash通过IAM角色连接AWS MSK报错排查与解决
问题场景
使用aws-msk-iam-auth插件配置Logstash连接AWS MSK,配置如下:
kafka { codec => "json" topic_id => "my_topic" bootstrap_servers => "my_iamBroder:9098" security_protocol => "SASL_SSL" sasl_mechanism => "AWS_MSK_IAM" sasl_jaas_config => "software.amazon.msk.auth.iam.IAMLoginModule required awsRoleArn='my_role_arn' sasl.client.callback.handler.class='software.amazon.msk.auth.iam.IAMClientCallbackHandler';" }
运行时抛出如下错误:
[2022-07-26T07:29:29,927][ERROR][org.apache.kafka.common.utils.KafkaThread] Uncaught exception in thread 'kafka-producer-network-thread | producer-1': java.lang.NoClassDefFoundError: org/apache/kafka/common/errors/IllegalSaslStateException at software.amazon.msk.auth.iam.internals.IAMSaslClient$IAMSaslClientFactory.createSaslClient(IAMSaslClient.java:216) ~[aws-msk-iam-auth-1.1.4.jar:?] at javax.security.sasl.Sasl.createSaslClient(Sasl.java:433) ~[?:?] at org.apache.kafka.common.security.authenticator.SaslClientAuthenticator.lambda$createSaslClient$0(SaslClientAuthenticator.java:217) ~[kafka-clients-2.5.1.jar:?] at java.security.AccessController.doPrivileged(Native Method) ~[?:?] at javax.security.auth.Subject.doAs(Subject.java:423) ~[?:?] at org.apache.kafka.common.security.authenticator.SaslClientAuthenticator.createSaslClient(SaslClientAuthenticator.java:213) ~[kafka-clients-2.5.1.jar:?] at org.apache.kafka.common.security.authenticator.SaslClientAuthenticator.<init>(SaslClientAuthenticator.java:204) ~[kafka-clients-2.5.1.jar:?] at org.apache.kafka.common.network.SaslChannelBuilder.buildClientAuthenticator(SaslChannelBuilder.java:274) ~[kafka-clients-2.5.1.jar:?] at org.apache.kafka.common.network.SaslChannelBuilder.lambda$buildChannel$1(SaslChannelBuilder.java:216) ~[kafka-clients-2.5.1.jar:?] at org.apache.kafka.common.network.KafkaChannel.<init>(KafkaChannel.java:142) ~[kafka-clients-2.5.1.jar:?] at org.apache.kafka.common.network.SaslChannelBuilder.buildChannel(SaslChannelBuilder.java:224) ~[kafka-clients-2.5.1.jar:?] at org.apache.kafka.common.network.Selector.buildAndAttachKafkaChannel(Selector.java:338) ~[kafka-clients-2.5.1.jar:?] at org.apache.kafka.common.network.Selector.registerChannel(Selector.java:329) ~[kafka-clients-2.5.1.jar:?]
错误原因
IllegalSaslStateException是Apache Kafka客户端(kafka-clients)2.6.0版本新增的异常类,当前环境中使用的kafka-clients版本为2.5.1,不包含该类。- aws-msk-iam-auth 1.1.4版本依赖的kafka-clients版本高于2.5.1,版本不匹配导致类找不到的错误。
解决方法
方法1:升级kafka-clients版本
找到Logstash安装目录下的kafka-clients jar包(路径通常为vendor/bundle/jruby/2.5.0/gems/logstash-integration-kafka-*/vendor/kafka-clients-*.jar),替换为2.6.0及以上兼容版本的kafka-clients jar包。
方法2:降级aws-msk-iam-auth插件
选择支持kafka-clients 2.5.x的aws-msk-iam-auth版本(比如1.0.0版本),重新安装该版本插件:
logstash-plugin install --version 1.0.0 aws-msk-iam-auth
方法3:验证IAM角色权限(辅助排查)
确保Logstash所在实例使用的IAM角色(或配置中指定的awsRoleArn)拥有kafka:DescribeCluster、kafka:Connect等MSK连接所需的权限,且角色的信任关系配置正确,允许Logstash所在实体(EC2实例、ECS任务等)扮演该角色。
内容的提问来源于stack exchange,提问作者Shihabudheen K M
相关产品推荐
相关产品推荐

