如何在PyFlink Table API中读取多个目录下的CSV数据?
PyFlink Table API读取多个目录的解决方案
错误原因
Flink Filesystem Connector的path参数默认仅支持传入单个路径,你传入逗号分隔的多路径字符串时,Flink会将整个字符串识别为单个目录路径去访问,因此触发路径不存在的报错。
解决方案
方案1:通配符匹配(适合路径有统一规则的场景)
你的路径符合固定前缀+连续后缀的规则,直接用通配符匹配即可,修改DDL中的path配置:
ddl = """ CREATE TABLE test ( a INT, b STRING ) WITH ( 'connector' = 'filesystem', 'path' = '/opt/data/day=2021-11-1[4-6]', 'format' = 'csv', 'csv.ignore-first-line' = 'true', 'csv.ignore-parse-errors' = 'true', 'csv.array-element-delimiter' = ';' ) """
也可以用大括号匹配离散值:'/opt/data/day=2021-11-1{4,5,6}'
方案2:分区识别(适合按分区规则命名的目录场景)
你的目录本身就是{分区名}={分区值}的标准分区命名格式,可以通过配置分区识别直接读取父目录,查询时过滤需要的分区即可:
ddl = """ CREATE TABLE test ( a INT, b STRING, `day` STRING ) WITH ( 'connector' = 'filesystem', 'path' = '/opt/data', 'format' = 'csv', 'csv.ignore-first-line' = 'true', 'csv.ignore-parse-errors' = 'true', 'csv.array-element-delimiter' = ';', 'scan.partition.discovery.enabled' = 'true' ) """ # 执行查询时过滤需要的分区 table_env.execute_sql("SELECT * FROM test WHERE `day` IN ('2021-11-14','2021-11-15','2021-11-16')")
该方案还可以直接在查询结果中获取day分区字段值,不需要额外解析路径。
方案3:指定多路径参数(适合无规律离散路径场景,Flink 1.14+支持)
对于完全无规律的多个离散路径,可以通过scan.file.paths参数传入逗号分隔的多路径,该参数会覆盖path配置的读取路径:
ddl = """ CREATE TABLE test ( a INT, b STRING ) WITH ( 'connector' = 'filesystem', 'path' = '/opt/data', 'scan.file.paths' = '/opt/data/day=2021-11-14,/opt/data/day=2021-11-15,/opt/data/day=2021-11-16', 'format' = 'csv', 'csv.ignore-first-line' = 'true', 'csv.ignore-parse-errors' = 'true', 'csv.array-element-delimiter' = ';' ) """
以上三种方案都不需要创建多张表再做union操作,不会产生冗余代码。
内容的提问来源于stack exchange,提问作者chenyf
相关产品推荐
相关产品推荐

