【问题标题】:dask read parquet and specify schemadask 读取镶木地板并指定架构
【发布时间】:2021-04-01 02:01:42
【问题描述】:

在读取 parquet 文件时,是否有与 spark 指定模式的能力相当的功能?可能使用传递给 pyarrow 的 kwargs?

我的存储桶中有一堆 parquet 文件,但其中一些字段的名称略有不一致。我可以在阅读这些情况后创建一个自定义延迟函数来处理这些情况,但我希望我可以在通过 globing 打开它们时指定架构。也许不是,因为我猜想通过 globing 打开然后尝试将它们连接起来。由于字段名称不一致,目前此操作失败。

创建拼花文件:

import dask.dataframe as dd

df = dd.demo.make_timeseries(
    start="2000-01-01",
    end="2000-01-03",
    dtypes={"id": int, "z": int},
    freq="1h",
    partition_freq="24h",
)

df.to_parquet("df.parquet", engine="pyarrow", overwrite=True)

通过dask读入,读入后指定schema:

df = dd.read_parquet("df.parquet", engine="pyarrow")
df["z"] = df["z"].astype("float")
df = df.rename(columns={"z": "a"})

通过 spark 读入并指定架构:

from pyspark.sql import SparkSession
import pyspark.sql.types as T
spark = SparkSession.builder.appName('App').getOrCreate()

schema = T.StructType(
    [
        T.StructField("id", T.IntegerType()),
        T.StructField("a", T.FloatType()),
        T.StructField("timestamp", T.TimestampType()),
    ]
)

df = spark.read.format("parquet").schema(schema).load("df.parquet")

【问题讨论】:

    标签: pandas apache-spark dask parquet pyarrow


    【解决方案1】:

    一些选项是:

    1. 加载后指定 dtypes(需要一致的列名):
    custom_dtypes = {"a": float, "id": int, "timestamp": pd.datetime}
    df = dd.read_parquet("df.parquet", engine="pyarrow").astype(custom_dtypes)
    

    目前由于字段名称不一致而失败。

    1. 如果文件中的列名不同,您可能希望在加载之前使用自定义delayed::
    @delayed
    def custom_load(path):
       df = pd.read_parquet(path)
       # some logic to ensure consistent columns
       # for example:
       if "z" in df.columns:
          df = df.rename(columns={"z": "a"}).astype(custom_dtypes)
       return df
    
    dask_df = dd.from_delayed([custom_load(path) for path in glob.glob("some_path/*parquet")])
    

    【讨论】:

      猜你喜欢
      • 2020-07-27
      • 1970-01-01
      • 1970-01-01
      • 2017-12-25
      • 2018-12-28
      • 2023-01-05
      • 2020-02-01
      • 2020-09-18
      • 1970-01-01
      相关资源
      最近更新 更多