PySpark读取Hudi分区表返回空DataFrame问题排查
问题:Hudi分区表读取返回空DataFrame(写入/流读取/直接读Parquet均正常)
环境信息
- Azure Databricks集群,Runtime版本9.1 LTS(内置Apache Spark 3.1.2、Scala 2.12)
- Python 3.8.10
- 已安装Maven包:
hudi_spark3_1_2_bundle_2_12_0_10_1.jar
写入情况
按照Hudi官方文档编写PySpark代码,将DataFrame写入Azure DataLake Gen2的Hudi分区表,已生成符合预期的多级分区目录结构。写入代码如下:
tableName = "my_hudi_table" basePath = <<table_path>> dataGen = sc._jvm.org.apache.hudi.QuickstartUtils.DataGenerator() inserts = sc._jvm.org.apache.hudi.QuickstartUtils.convertToStringList(dataGen.generateInserts(10)) df_hudi = spark.read.json(spark.sparkContext.parallelize(inserts, 2)) hudi_options = { 'hoodie.table.name': tableName, 'hoodie.datasource.write.recordkey.field': 'uuid', 'hoodie.datasource.write.partitionpath.field': 'partitionpath', 'hoodie.datasource.write.table.name': tableName, 'hoodie.datasource.write.operation': 'upsert', 'hoodie.datasource.write.precombine.field': 'ts', 'hoodie.upsert.shuffle.parallelism': 2, 'hoodie.insert.shuffle.parallelism': 2 } df_hudi.write.format("hudi").options(**hudi_options).mode("overwrite").save(basePath)
读取异常情况
使用Hudi格式读取全表时返回空DataFrame,但Schema结构正确。读取代码如下:
partition_column = "partitionpath" hudi_read_options = { 'hoodie.datasource.read.partitionpath.field': partition_column, 'hoodie.file.index.enable': 'false' } df_hudi = spark.read.format("hudi").options(**hudi_read_options).load(basePath) df_hudi.printSchema() # 输出Schema root |-- _hoodie_commit_time: string (nullable = true) |-- _hoodie_commit_seqno: string (nullable = true) |-- _hoodie_record_key: string (nullable = true) |-- _hoodie_partition_path: string (nullable = true) |-- _hoodie_file_name: string (nullable = true) |-- begin_lat: double (nullable = true) |-- begin_lon: double (nullable = true) |-- driver: string (nullable = true) |-- end_lat: double (nullable = true) |-- end_lon: double (nullable = true) |-- fare: double (nullable = true) |-- partitionpath: string (nullable = true) |-- rider: string (nullable = true) |-- ts: long (nullable = true) |-- uuid: string (nullable = true)
注:设置
'hoodie.file.index.enable': 'false'是因为不设置会报错:NoSuchMethodError: org.apache.spark.sql.execution.datasources.FileStatusCache.putLeafFiles(Lorg/apache/hadoop/fs/Path;[Lorg/apache/hadoop/fs/FileStatus;)V
正常读取的场景
- 直接读取特定分区的Parquet文件可获取数据:
df = spark.read.format("parquet").load(basePath + "/americas/brazil/sao_paulo/") - 使用Spark结构化流读取正常:
spark.readStream.format("hudi").load(basePath)
请问我遗漏了哪些配置?
内容的提问来源于stack exchange,提问作者jakeis
相关产品推荐
相关产品推荐

