【发布时间】: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