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

PyFlink 1.14 Table API连接带认证Kafka无任务生成问题求助

问题排查及修正点

  • 依赖冲突:同时加载了Scala 2.11、2.12版本的Kafka连接器,还额外引入了独立的kafka-clients包,版本冲突会导致作业初始化卡住。flink-sql-connector-kafka 已经内置了兼容的kafka客户端依赖,无需额外引入其他Kafka相关jar包,且连接器的Scala版本需要和你的Flink集群使用的Scala版本完全一致。
  • 执行逻辑重复:execute_insert("sink_table").wait() 已经会触发作业提交并等待执行,后续的 t_env.execute("kafka-table") 属于冗余执行逻辑,会导致提交流程阻塞。
  • Sink配置冗余:Kafka Sink不需要配置scan.startup.mode、properties.group.id这两个仅Source生效的参数。
  • 认证配置转义风险:Python f-string中转义双引号可能会导致jaas配置解析异常,建议改用单引号包裹配置内容避免转义问题。

修正后完整代码

from pyflink.datastream.stream_execution_environment import StreamExecutionEnvironment
from pyflink.table import EnvironmentSettings
from pyflink.table.table_environment import StreamTableEnvironment

KAFKA_SERVERS = 'localhost:9092'
KAFKA_USERNAME = "user"
KAFKA_PASSWORD = "XXX"
KAFKA_SOURCE_TOPIC = 'source'
KAFKA_SINK_TOPIC = 'dest'


def log_processing():
    env = StreamExecutionEnvironment.get_execution_environment()
    # 仅保留和集群Scala版本匹配的flink-sql-connector-kafka即可,无需额外加载其他Kafka相关jar
    env.add_jars("file:///opt/flink/lib_py/flink-sql-connector-kafka_2.12-1.14.0.jar")
    settings = EnvironmentSettings.new_instance()\
                      .in_streaming_mode()\
                      .use_blink_planner()\
                      .build()

    t_env = StreamTableEnvironment.create(stream_execution_environment=env, environment_settings=settings)
    
    source_ddl = f"""
            CREATE TABLE source_table(
                Cylinders INT,
                Displacement INT,
                Horsepower INT,
                Weight INT,
                Acceleration INT,
                Model_Year INT,
                USA INT,
                Europe INT,
                Japan INT
            ) WITH (
              'connector' = 'kafka',
              'topic' = '{KAFKA_SOURCE_TOPIC}',
              'properties.bootstrap.servers' = '{KAFKA_SERVERS}',
              'properties.group.id' = 'testgroup12',
              'properties.sasl.mechanism' = 'PLAIN',
              'properties.security.protocol' = 'SASL_PLAINTEXT',
              'properties.sasl.jaas.config' = 'org.apache.kafka.common.security.plain.PlainLoginModule required username=\'{KAFKA_USERNAME}\' password=\'{KAFKA_PASSWORD}\';',
              'scan.startup.mode' = 'latest-offset',
              'format' = 'json'
            )
            """

    sink_ddl = f"""
            CREATE TABLE sink_table(
                Cylinders INT,
                Displacement INT,
                Horsepower INT,
                Weight INT,
                Acceleration INT,
                Model_Year INT,
                USA INT,
                Europe INT,
                Japan INT
            ) WITH (
              'connector' = 'kafka',
              'topic' = '{KAFKA_SINK_TOPIC}',
              'properties.bootstrap.servers' = '{KAFKA_SERVERS}',
              'properties.sasl.mechanism' = 'PLAIN',
              'properties.security.protocol' = 'SASL_PLAINTEXT',
              'properties.sasl.jaas.config' = 'org.apache.kafka.common.security.plain.PlainLoginModule required username=\'{KAFKA_USERNAME}\' password=\'{KAFKA_PASSWORD}\';',
              'format' = 'json'
            )
            """

    t_env.execute_sql(source_ddl)
    t_env.execute_sql(sink_ddl)

    # 移除冗余的execute调用,仅保留insert执行即可
    t_env.sql_query("SELECT * FROM source_table").execute_insert("sink_table").wait()


if __name__ == '__main__':
    log_processing()

提交验证建议

  • 提交前先确认Flink集群使用的Scala版本,若为2.11版本则替换connector包为flink-sql-connector-kafka_2.11-1.14.0.jar即可
  • 提交作业时可增加-pyclientexec python3参数指定Python解释器版本,避免本地和集群Python环境不一致导致的初始化失败
  • 若仍提交失败,可查看Flink客户端提交日志、集群JobManager日志排查具体报错信息

内容的提问来源于stack exchange,提问作者Paul

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 00:06:01