PyFlink导入错误:无法从pyflink.table.descriptors导入OldCsv如何解决
问题原因
该报错是因为你通过pip install apache-flink默认安装的是最新版Flink,而你参考的是1.13版本的教程:Flink从1.14版本开始就标记OldCsv类为废弃状态,1.15及后续版本直接从pyflink.table.descriptors模块中移除了该类,同时原来的descriptor链式建表API也已经被弃用。
解决方法
方法一:降级Flink版本匹配教程
如果你希望直接沿用1.13版本的教程代码,直接将pyflink版本降到1.13.x系列即可:
pip install apache-flink==1.13.6
降级后原来的导入语句和业务代码无需修改,可直接运行。
方法二:适配新版本Flink的API
如果你希望使用最新版Flink,需要替换废弃的descriptor API,改用DDL语句创建源表:
修改导入语句
删除descriptors相关的导入,修改后导入如下:
from pyflink.table import DataTypes, TableEnvironment, EnvironmentSettings from pyflink.table.expressions import lit
修改建表逻辑
原来的链式建表代码替换为DDL执行逻辑:
t_env.get_config().get_configuration().set_string("parallelism.default", "1") # 用DDL声明文件数据源 t_env.execute_sql(f""" CREATE TABLE Source ( word STRING ) WITH ( 'connector' = 'filesystem', 'path' = '{input_file}', 'format' = 'csv' ) """)
新版本中csv格式直接替代了原来的OldCsv,不需要额外导入依赖,pyflink默认已内置csv格式支持。
内容的提问来源于stack exchange,提问作者Samarth BM
相关产品推荐
相关产品推荐

