Pandas 无法读取在 PySpark 中创建的镶木地板文件

Tho*_*mas 7 python pandas apache-spark parquet pyspark

我正在通过以下方式从 Spark DataFrame 编写镶木地板文件:

df.write.parquet("path/myfile.parquet", mode = "overwrite", compression="gzip")
Run Code Online (Sandbox Code Playgroud)

这将创建一个包含多个文件的文件夹。

当我尝试将其读入 Pandas 时,出现以下错误,具体取决于我使用的解析器:

import pandas as pd
df = pd.read_parquet("path/myfile.parquet", engine="pyarrow")
Run Code Online (Sandbox Code Playgroud)

派箭:

文件“pyarrow\error.pxi”,第 83 行,在 pyarrow.lib.check_status

ArrowIOError: 镶木地板文件无效。损坏的页脚。

快速拼花:

文件“C:\Program Files\Anaconda3\lib\site-packages\fastparquet\util.py”,第 38 行,在 default_open 中 return open(f, mode)

PermissionError: [Errno 13] 权限被拒绝: 'path/myfile.parquet'

我正在使用以下版本:

  • 火花 2.4.0
  • 熊猫 0.23.4
  • pyarrow 0.10.0
  • 快速拼花 0.2.1

我尝试了 gzip 以及 snappy 压缩。两者都不起作用。我当然确保我将文件放在 Python 有权读/写的位置。

如果有人能够重现此错误,那将会有所帮助。

mar*_*oyo 5

问题是 Spark 由于其分布式特性而对文件进行分区(每个执行程序在接收文件名的目录中写入一个文件)。这不是 Pandas 支持的东西,它需要一个文件,而不是一个路径。

您可以通过不同方式规避此问题:

  • 谢谢您的回答。似乎读取单个文件(您的第二个要点)有效。但是,第一件事不起作用 - 看起来 pyarrow 无法处理 PySpark 的页脚(请参阅相关错误消息) (2认同)

Tho*_*mas 4

由于即使对于较新的 pandas 版本,这似乎仍然是一个问题,因此我编写了一些函数来规避此问题,作为更大的 pyspark 帮助程序库的一部分:

import pandas as pd
import datetime
import os

def read_parquet_folder_as_pandas(path, verbosity=1):
  files = [f for f in os.listdir(path) if f.endswith("parquet")]

  if verbosity > 0:
    print("{} parquet files found. Beginning reading...".format(len(files)), end="")
    start = datetime.datetime.now()

  df_list = [pd.read_parquet(os.path.join(path, f)) for f in files]
  df = pd.concat(df_list, ignore_index=True)

  if verbosity > 0:
    end = datetime.datetime.now()
    print(" Finished. Took {}".format(end-start))
  return df


def read_parquet_as_pandas(path, verbosity=1):
  """Workaround for pandas not being able to read folder-style parquet files.
  """
  if os.path.isdir(path):
    if verbosity>1: print("Parquet file is actually folder.")
    return read_parquet_folder_as_pandas(path, verbosity)
  else:
    return pd.read_parquet(path)
Run Code Online (Sandbox Code Playgroud)

这假设 parquet“文件”(实际上是一个文件夹)中的相关文件以“.parquet”结尾。这适用于由 databricks 导出的镶木地板文件,并且也可能与其他文件一起使用(未经测试,对评论中的反馈感到高兴)。

read_parquet_as_pandas()如果事先不知道它是否是文件夹,则可以使用该功能。