Spark SQL:如何批量读取Parquet文件并执行全外连接查询
问题解决方法
你直接用逗号分隔多个Parquet路径的写法是错误的——这种写法会被Spark解析为多表笛卡尔积,且别名target_table仅绑定了第二个文件,导致查询逻辑混乱,无法正常执行。以下是三种可行的解决方案:
方案1:使用通配符或目录路径读取多个文件(性能最优)
直接通过通配符匹配目标文件,或读取整个分区目录(Spark会自动加载目录下所有Parquet文件),无需逐个指定文件路径:
-- 方法A:用通配符匹配指定范围的part文件 spark.sql("select target_table.*,lkp.account FROM parquet.`hdfs://gcgprod/data/work/hive/db/table/process_date=20221220/part-0088[6-7]*` as target_table Full OUTER JOIN db.table_ref lkp ON target_table.chd_account_number = lkp.account") -- 方法B:直接读取整个分区目录(加载目录下所有Parquet文件) spark.sql("select target_table.*,lkp.account FROM parquet.`hdfs://gcgprod/data/work/hive/db/table/process_date=20221220/` as target_table Full OUTER JOIN db.table_ref lkp ON target_table.chd_account_number = lkp.account")
方案2:用UNION ALL合并指定文件后再连接
如果需要精确指定某几个文件,可先通过UNION ALL合并数据,再执行全外连接:
spark.sql(""" select combined.*, lkp.account FROM ( select * from parquet.`hdfs://gcgprod/data/work/hive/db/table/process_date=20221220/part-00886-6ba5c0ad97ba.c000` union all select * from parquet.`hdfs://gcgprod/data/work/hive/db/table/process_date=20221220/part-00887-6ba5c0ad97ba.c000` ) as combined Full OUTER JOIN db.table_ref lkp ON combined.chd_account_number = lkp.account """)
方案3:通过DataFrame API读取文件并注册临时视图
先利用DataFrame API加载多个文件,再注册为临时视图执行SQL查询:
# 读取指定的多个Parquet文件 target_df = spark.read.parquet( "hdfs://gcgprod/data/work/hive/db/table/process_date=20221220/part-00886-6ba5c0ad97ba.c000", "hdfs://gcgprod/data/work/hive/db/table/process_date=20221220/part-00887-6ba5c0ad97ba.c000" ) # 注册为临时视图 target_df.createOrReplaceTempView("target_table") # 执行全外连接查询 spark.sql("select target_table.*,lkp.account FROM target_table Full OUTER JOIN db.table_ref lkp ON target_table.chd_account_number = lkp.account")
内容的提问来源于stack exchange,提问作者Arvinth
相关产品推荐
相关产品推荐

