PyFlink 1.11无法连接Confluent Cloud Kafka集群问题咨询
问题与解决方案:PyFlink连接Confluent Cloud Kafka(SASL/PLAIN认证)
问题描述
配置PyFlink(1.11或1.13版本)连接Confluent Cloud Kafka集群,采用SASL/PLAIN认证方式,使用如下建表SQL:
""" CREATE TABLE {0} ( `transaction_amt` BIGINT NOT NULL, `event_id` VARCHAR(64) NOT NULL, `event_time` TIMESTAMP(6) NOT NULL ) WITH ( 'connector' = 'kafka', 'topic' = '{1}', 'properties.bootstrap.servers' = '{2}', 'properties.group.id' = 'testGroupTFI', 'format' = 'json', 'json.timestamp-format.standard' = 'ISO-8601', 'properties.security.protocol' = 'SASL_SSL', 'properties.sasl.mechanism' = 'PLAIN', 'properties.sasl.jaas.config' = 'org.apache.kafka.common.security.plain.PlainLoginModule required username=\"{3}\" password=\"{4}\";' ) """.format(table_name, stream_name, broker, user, secret)
运行作业后出现失败,核心异常信息为:
Caused by: javax.security.auth.login.LoginException: No LoginModule found for org.apache.kafka.common.security.plain.PlainLoginModule
疑问:PyFlink 1.11或1.13的SQL Connector是否不支持SASL?有何可行解决办法?
解决方案
1. SASL支持判断
你的判断不正确,PyFlink 1.11和1.13版本的Kafka SQL Connector完全支持SASL/PLAIN认证。错误根源是Flink对Kafka依赖包做了shading(包路径重命名)处理,导致配置的LoginModule类名无法被加载。
2. 具体修复步骤
修改JAAS配置中的LoginModule类名为Flink shaded后的完整路径:
将org.apache.kafka.common.security.plain.PlainLoginModule替换为org.apache.flink.kafka.shaded.org.apache.kafka.common.security.plain.PlainLoginModule
修改后的建表SQL示例:
""" CREATE TABLE {0} ( `transaction_amt` BIGINT NOT NULL, `event_id` VARCHAR(64) NOT NULL, `event_time` TIMESTAMP(6) NOT NULL ) WITH ( 'connector' = 'kafka', 'topic' = '{1}', 'properties.bootstrap.servers' = '{2}', 'properties.group.id' = 'testGroupTFI', 'format' = 'json', 'json.timestamp-format.standard' = 'ISO-8601', 'properties.security.protocol' = 'SASL_SSL', 'properties.sasl.mechanism' = 'PLAIN', 'properties.sasl.jaas.config' = 'org.apache.flink.kafka.shaded.org.apache.kafka.common.security.plain.PlainLoginModule required username=\"{3}\" password=\"{4}\";' ) """.format(table_name, stream_name, broker, user, secret)
3. 额外注意事项
- 提交PyFlink作业时,需确保包含对应版本的
flink-sql-connector-kafka_xxx.jar依赖包,可通过--jar参数指定; - 若使用AWS Kinesis Analytics for Flink(从错误日志ARN可判断),需确认托管环境已正确加载包含shaded Kafka依赖的Connector包,部分场景需手动配置依赖加载规则。
内容的提问来源于stack exchange,提问作者Ahmed Zamzam
相关产品推荐
相关产品推荐

