PyFlink对接Kafka异常:依赖与任务提交问题咨询
PyFlink读取Kafka Topic问题排查
报错信息
使用PyFlink读取Kafka时出现以下错误:
[ERROR] Could not execute SQL statement. Reason: org.apache.flink.table.api.ValidationException: Could not find any factory for identifier 'kafka' that implements 'org.apache.flink.table.factories.DynamicTableFactory' in the classpath. Available factory identifiers are: blackhole datagen filesystem print python-input-format
运行的PyFlink代码
def main(): env_settings = EnvironmentSettings.in_streaming_mode() t_env = StreamTableEnvironment.create(environment_settings=env_settings) t_env.execute_sql( f""" CREATE TABLE kafka_test_logs ( action2 STRING ) WITH ( 'connector' = 'kafka', 'topic' = 'test-logs', 'properties.bootstrap.servers' = 'kafka:9092', 'properties.group.id' = 'flink_group', 'scan.startup.mode' = 'earliest-offset', 'format' = 'json', 'json.ignore-parse-errors' = 'true' ) """) t_env.execute_sql(""" SELECT COUNT(*) AS message_count FROM kafka_test_logs """).print() if __name__ == "__main__": main()
环境配置
- 使用的Flink镜像:
flink:1.18.1-scala_2.12 - JobManager容器中
/opt/flink/lib目录下的初始依赖:
flink-cep-1.18.1.jar flink-connector-files-1.18.1.jar flink-connector-kafka-1.17.2.jar flink-csv-1.18.1.jar flink-dist-1.18.1.jar flink-json-1.18.1.jar flink-scala_2.12-1.18.1.jar flink-shaded-guava-30.1.1-jre-16.0.jar flink-table-api-java-uber-1.18.1.jar flink-table-planner-loader-1.18.1.jar flink-table-runtime-1.18.1.jar kafka-clients-3.6.1.jar log4j-1.2-api-2.17.1.jar log4j-api-2.17.1.jar log4j-core-2.17.1.jar log4j-slf4j-impl-2.17.1.jar
- docker-compose.yml中的挂载配置:
volumes: - ./libs/flink-connector-kafka-1.17.2.jar:/opt/flink/lib/flink-connector-kafka-1.17.2.jar - ./libs/flink-shaded-guava-30.1.1-jre-16.0.jar:/opt/flink/lib/flink-shaded-guava-30.1.1-jre-16.0.jar - ./libs/kafka-clients-3.6.1.jar:/opt/flink/lib/kafka-clients-3.6.1.jar
初始测试时,Flink SQL可创建表但执行查询仍报缺失Kafka依赖错误,不过JobManager可连接Kafka。
更新后的状态
在docker-compose.yml中添加了flink-sql-connector后,当前JobManager的类路径如下:
/opt/flink/lib/flink-cep-1.18.1.jar:/opt/flink/lib/flink-connector-files-1.18.1.jar:/opt/flink/lib/flink-csv-1.18.1.jar:/opt/flink/lib/flink-json-1.18.1.jar:/opt/flink/lib/flink-scala_2.12-1.18.1.jar:/opt/flink/lib/flink-sql-connector-kafka-3.1.0-1.18.jar:/opt/flink/lib/flink-table-api-java-uber-1.18.1.jar:/opt/flink/lib/flink-table-planner-loader-1.18.1.jar:/opt/flink/lib/flink-table-runtime-1.18.1.jar:/opt/flink/lib/log4j-1.2-api-2.17.1.jar:/opt/flink/lib/log4j-api-2.17.1.jar:/opt/flink/lib/log4j-core-2.17.1.jar:/opt/flink/lib/log4j-slf4j-impl-2.17.1.jar:/opt/flink/lib/flink-dist-1.18.1.jar
现在出现新问题:
- 运行PyFlink客户端容器时,没有任务提交到集群;
- 在JobManager中通过Flink SQL执行
SELECT * FROM kafka_test_logs时,看似提交了任务,但未实际消费数据:执行kafka-consumer-groups.sh查询flink_group时显示该消费组不存在,JobManager的UI stdout为空,TaskManager和JobManager日志仅显示连接相关信息但无消费迹象。
疑问
- 为何Flink SQL可提交任务但PyFlink客户端不行?
- 如何输出消息计数结果以验证Flink是否成功消费Kafka Topic?
- Flink SQL连接器JAR是否与使用的Flink Docker镜像不兼容?
内容的提问来源于stack exchange,提问作者jytu65
相关产品推荐
相关产品推荐

