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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.23 06:54:23