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
相关产品推荐
相关产品推荐

