【发布时间】:2020-05-30 07:12:15
【问题描述】:
我正在尝试将这样的 Dataframe 写入 Parquet:
| foo | bar |
|-----|-------------------|
| 1 | {"a": 1, "b": 10} |
| 2 | {"a": 2, "b": 20} |
| 3 | {"a": 3, "b": 30} |
我正在使用 Pandas 和 Fastparquet:
df = pd.DataFrame({
"foo": [1, 2, 3],
"bar": [{"a": 1, "b": 10}, {"a": 2, "b": 20}, {"a": 3, "b": 30}]
})
import fastparquet
fastparquet.write('/my/parquet/location/toy-fastparquet.parq', df)
我想在 (py)Spark 中加载 Parquet,并使用 Spark SQL 查询数据,例如:
df = spark.read.parquet("/my/parquet/location/")
df.registerTempTable('my_toy_table')
result = spark.sql("SELECT * FROM my_toy_table WHERE bar.b > 15")
我的问题是,即使fastparquet 可以正确读取其 Parquet 文件(bar 字段被正确反序列化为结构),在 Spark 中,bar 被读取为String 类型的列,仅包含原始结构的 JSON 表示:
In [2]: df.head()
Out[2]: Row(foo=1, bar='{"a": 1, "b": 10}')
我尝试从 PyArrow 编写 Parquet,但没有运气:ArrowNotImplementedError: Level generation for Struct not supported yet。我也尝试将file_scheme='hive' 传递给 Fastparquet,但我得到了相同的结果。将 Fastparquet 序列化更改为 BSON (object_encoding='bson') 会产生不可读的二进制字段。
[编辑]我看到了以下方法:
- [answered] 从 Spark 编写 Parquet
- [open] 查找实现Parquet's specification for nested types 的 Python 库,并且与 Spark 读取它们的方式兼容
- [answered] 使用特定的 JSON 反序列化读取 Spark 中的 Fastparquet 文件(我想这会对性能产生影响)
- 不要完全使用嵌套结构
【问题讨论】:
-
这确实是Arrow目前的局限,见issues.apache.org/jira/browse/ARROW-1644
-
谢谢@joris,我的 DF 不包含列表和结构的混合,只是一个结构字段(我使描述更清楚)。但是,目前似乎也不支持这种情况。
-
您是否尝试在加载数据时传递
schema? -
@cesar-a-mostacero 我试过了,但没有成功,因为我错过了 Alexandros 在下面的答案中解释的 JSON 解码
标签: apache-spark pyspark apache-spark-sql pyarrow fastparquet