使用 Haskell 编写 Parquet 文件
摘要
本文详细介绍了如何在 Haskell 中为 DataHaskell/Dataframe 库实现 Parquet 文件写入器,以实现高效的数据序列化并与数据科学生态系统互操作。
暂无内容
查看缓存全文
缓存时间: 2026/09/25 01:10
# 使用Haskell编写Parquet文件
来源:https://www.datahaskell.org/blog/2026/09/18/writing-parquet-files-using-haskell.html
我们在 DataHaskell/Dataframe (https://github.com/DataHaskell/dataframe) 项目中实现了一个Parquet写入器。使用它很简单;你只需将数据框传递给 `writeParquet` 函数,它就会使用合理的行组和页大小默认值写入一个Parquet文件。示例如下:
```
import qualified DataFrame as D
import qualified DataFrame.Functions as F
import DataFrame (as, (|>))
main = do
sales <- D.readParquet "sales_data.parquet"
sales
|> D.groupBy ["product"]
|> D.aggregate [ F.sum (F.col @Int "amount") `as` "total"
, F.count (F.col @Int "amount") `as` "orders"
]
|> D.writeParquet "total_orders.parquet"
```
如果你需要对Parquet文件进行更细粒度的控制,你将需要使用 `writeParquetWithOptions`。继续阅读以了解这些选项及其如何影响最终文件。
为了让Haskell与数据生态系统进行互操作,它必须能够理解该生态系统使用的标准格式。长久以来,我们在Haskell中序列化数据的方式仅限于CSV、JSON,或直接将其转储到ByteString中。虽然CSV和JSON有其适用场景,但随着数据量的增长,它们会带来相当大的缺点:读取、写入和查询速度缓慢且笨重,而任何自定义的内部格式往往缺乏与标准数据科学工具的互操作性。
Parquet格式用一定的复杂性换来了高效的数据存储和查询。我们希望我们的数据具有高压缩比,并希望最大限度地减少读取与特定查询无关的数据。能够读写Parquet文件用于长期存储、网络传输或与其他程序(尤其是在数据科学生态系统中Parquet无处不在的情况下)互操作,是一个相当有用的工具,我们希望能在数据框库中拥有这个能力。
## Parquet文件结构?
如果你已经了解Parquet文件的结构,可以跳过本节。
Parquet文件是一系列行组,文件末尾是元数据,其中包含的信息允许读取器定位相关的列块和页。它还包含统计信息和布隆过滤器等有用信息,这样读取器/查询规划器就可以决定是否读取某个特定行组。
每个行组是一系列列块的集合,每个列块包含相同数量的行。每个列块是一系列数据页。由于每个列块是一系列页,而每个行组是一系列列块,最终文件看起来就像每个列的页依次排列。我们能够通过使用元数据识别每个行组和列块的偏移量和大小来理解全部内容。
数据页是我们实际存储所有数据的地方。它首先包含页元数据,描述其编码方式、值数量、统计信息等。实际上存在两个略有不同的数据页版本(https://github.com/apache/parquet-format/blob/master/src/main/thrift/parquet.thrift#L699, https://github.com/apache/parquet-format/blob/master/src/main/thrift/parquet.thrift#L752)。
接下来是定义级别和重复级别。它们是对可空性和嵌套结构的一种经济编码。关于定义级别和重复级别的详细描述超出了本文的范围;请参阅Dremel论文(https://static.googleusercontent.com/media/research.google.com/en//pubs/archive/36632.pdf)。目前,我们的实现仅支持写入定义级别最高到1,以表示可空值。
最后是我们实际编码(https://github.com/apache/parquet-format/blob/master/Encodings.md)并压缩(https://github.com/apache/parquet-format/blob/master/Compression.md)的数据(具体取决于我们使用的数据页类型,我们压缩数据本身,或者同时压缩定义/重复级别和数据)。编码是在每个页确定的,而压缩则是在每个列块确定的。
## 实现细节
接下来是关于Parquet写入器设计及其权衡的详细技术讨论。
### 写入选项
根据上述对Parquet格式的描述,写入器需要暴露一组适当的调节参数,允许用户调整写入器以生成对特定数据高效处理的Parquet文件。也就是说,我们希望有一个Parquet写入器,它能以适当大小的行组和页进行有效压缩,以便读取器可以有效应用并行性、投影下推(仅读取相关列块)、谓词下推(仅读取相关行组/剪枝不相关页)、IO下推(仅读取相关文件,假设数据被分片存储在多个文件中)等技术。
```
data ParquetWriteOptions = ParquetWriteOptions
{ pageSize :: !Int
, rowGroupSize :: !Int
, batchRows :: !Int
, subBatchRows :: !Int
, compressionCodec :: !CompressionCodec
, strategy :: !WriterStrategy
, maxRowsPerFile :: !(Maybe Int)
}
```
`pageSize` 和 `rowGroupSize` 分别是每个页和每个行组的目标大小(以字节为单位)。但行组中的每个列块必须包含相同数量的行,并且根据被编码的具体数据、所使用的编码方式和压缩算法,每个列块在达到目标大小之前会容纳不同数量的行。每个页也是如此。那么,我们如何确保实现页大小目标、行组大小目标,并且每个列块包含相同数量的行呢?
我们必须将目标 `pageSize` 和 `rowGroupSize` 视为尽力而为;它们可能略高于或低于目标。我们以 `batchRows` 大小的批次处理列,并在每批之后检查行组的状态。因此,每个行组包含 `batchRows` 行的整数倍;页也将包含 `subBatchRows` 行的整数倍(最终页和最终列块除外)。在页级别进行子批处理使我们能够减少每次写入后必须进行的IORef簿记工作量,从而显著提升速度。
### 内存管理
由于前述章节中描述的约束,很难提前知道缓冲区应设多大。列块缓冲区尤其如此——每个列在相同数据量中能容纳的行数/页数差异很大,我们可能预期一些列块缓冲区会明显大于其他缓冲区,并主导每个行组的空间占用。
因此,我们必须能够动态扩展内存。处理原始内存的一种便捷方式是使用 `MutableByteArray`。
```
data MemoryBuffer = MemoryBuffer
{ arrayRef :: !(IORef (MutableByteArray RealWorld))
, positionRef :: !(IORef Int)
}
```
在我们的实现中,我们使用固定(pinned)的 `ByteArray`,因为当需要刷新到另一个缓冲区(例如,当将页缓冲区刷新到列块缓冲区时)或文件(当将行组刷新到文件时)时,我们希望将其转换为 `Ptr Word8`。
使用固定的 `ByteArray` 在尝试扩展内存缓冲区时会稍微复杂一些。我们不能使用 `Data.Primitive` 提供的 `grow` 函数,而是必须分配一个新的固定 ByteArray 并允许旧的被垃圾回收。人们可能会担心堆碎片化,因为一个4KB GHC块中的单个固定对象可能使整个块保持活跃,但我们预期我们的缓冲区往往会比4KB大得多。此外,扩展操作应该很少发生,特别是在最初的几个页和第一个行组之后。
```
ensureCapacity :: MemoryBuffer -> Int -> IO (MutableByteArray RealWorld)
ensureCapacity buffer needed = do
array <- readIORef buffer.arrayRef
maxSize <- getSizeofMutableByteArray array
if needed <= maxSize
then pure array
else do
position <- readIORef buffer.positionRef
grown <- newPinnedByteArray (needed + (needed `div` 2))
copyMutableByteArray grown 0 array 0 position
writeIORef buffer.arrayRef grown
pure grown
{-# INLINE ensureCapacity #-}
```
我们还编写了辅助函数,用于将 `Word8`、`Word32`、`Word64`、`Int32`、`Int64`、`Integer`、`Float`、`Double` 和 `ByteString` 写入缓冲区。最后,我们有一个 `flushBufferToBuffer :: MemoryBuffer -> MemoryBuffer -> IO ()` 函数和一个 `flushBufferToFile :: WritableBinaryHandle -> MemoryBuffer -> IO ()` 函数。
### 核心循环
这里我们提供Parquet写入器主循环的高级概述。为了尽可能简单地解释设计,我们省略了大量复杂细节。完整实现请参阅 Writer.hs (https://github.com/DataHaskell/dataframe/blob/main/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs)。
在高级层面,我们可以将Parquet写入器视为对数据框的一个带副作用的折叠(fold)(也称为某些语言中的 `reduce`,或者更少见的 `accumulate`)。折叠可以看作是一个循环,也就是说,我们遍历数据框中的行。折叠通常是纯的,但“带副作用”意味着每次迭代都会产生某种副作用,在我们的情况下就是内存和磁盘的读写。最后,折叠从迭代中产生一个最终的单一结果,这就是我们必须附加到文件末尾的文件元数据。
首先,我们必须定义穿过写入器的状态以实现此目标:
```
data ParquetWriterState = ParquetWriterState
{ outputFileHandle :: !WritableBinaryHandle -- 一个专门为写入而设计的句柄的newtype包装
, columnChunks :: !(VB.Vector ColumnChunkState) -- 实际上就是我们的行组缓冲区
, currentFileOffsetRef :: !(IORef Int64)
, scratchBuffer :: !MemoryBuffer
, rowGroupMetadataRef :: !(IORef [RowGroup])
, rowNumberRef :: !(IORef Int)
}
data ColumnChunkState = ColumnChunkState
{ columnName :: !T.Text
, nullable :: !Bool
, schema :: !SchemaElement
, encoder :: !Encoder
, buffer :: !MemoryBuffer
, uncompressedBufferSize :: !(IORef Int64)
, pageState :: !PageState
}
data PageState = PageState
{ pageBuffer :: !MemoryBuffer -- 每列有自己的pageBuffer
, definitionLevels :: !DefLevels
, currentRowCount :: !(IORef Int)
}
-- 在 Encoder.hs 中
data Encoder = Encoder
{ encType :: !ThriftType
, convertedType :: !(Maybe ConvertedType)
, logicalType :: !(Maybe LogicalType)
, encodeValue :: !(MemoryBuffer -> Int -> Int -> IO (Int, Bool))
, finishValues :: !(MemoryBuffer -> Int -> IO Int) -- 在需要清理或后处理时使用
}
```
每个ColumnChunk的缓冲区在刷新到最终文件时被刷新,而单个页缓冲区则被刷新到ColumnChunk缓冲区(或者更准确地说,被用于组装最终刷新到ColumnChunk缓冲区的数据,因为我们还需要刷新页元数据和定义级别)。
如果我们省略设置初始 `writerState` 和其他变量的步骤,我们的核心循环非常简单:
```
loop :: Int -> IO ()
loop rowNum
| rowNum >= endRow = pure ()
| otherwise = do
let batchEnd = rowNum + min interval (endRow - rowNum)
writeBatch rowNum batchEnd
size <- bufferedSize writerState.columnChunks
when
(size >= options.rowGroupSize)
(flushRowGroup options writerState) -- 写入文件
loop batchEnd
```
然后是 `writeBatch`:
```
writeBatch :: Int -> Int -> IO ()
writeBatch rowNum batchEnd
| rowNum >= batchEnd = pure ()
| otherwise = do
let count = min options.subBatchRows (batchEnd - rowNum)
forM_ writerState.columnChunks (writeRows options scratchBuffer rowNum count)
modifyIORef writerState.rowNumberRef (+ count)
writeBatch (rowNum + count) batchEnd
```
`scratchBuffer` 是一个可复用的缓冲区,我们将用它来组装页,因为单个页缓冲区只包含值,而我们的定义级别位于不同的缓冲区中,并且页元数据此时还不存在。最终刷新到ColumnChunk缓冲区的就是这个 `scratchBuffer`。
`writeRows` 的内部工作原理与我们之前看到的 `loop` 有些相似。
```
rowWriterLoop !options !columnChunkState !end !size !position !row
| row >= end = writeIORef columnChunkState.pageState.pageBuffer.positionRef position
| position + options.pageSize > size = do
let page = columnChunkState.pageState
writeIORef page.pageBuffer.positionRef position
arr' <- ensureCapacity
page.pageBuffer
(position + max
options.pageSize
((end - row) * 64)
)
size' <- getSizeOfMutableByteArray arr'
rowWriterLoop options columnChunkState end size' position row
| otherwise = do
let page = columnChunkState.pageState
encode = columnChunkState.encoder.encodeValue
(position', notNull) <- encode page.pageBuffer position row
when columnChunkState.nullable $
pushDef page.definitionLevels (if notNull then 1 else 0)
rowWriterLoop options columnChunkState end size position' (row + 1)
```
注意,我们的写入器只处理定义级别最高到1的情况,并且不支持定义级别大于1或重复级别大于0的情况,因为目前我们特别只想支持具有扁平模式的数据框。
当一个页被填满时,我们可以将其刷新到父级 `columnChunkState.buffer`。
```
pageResidency <- readIORef page.pageBuffer.positionRef
defResidency <- readIORef page.definitionLevels.dlBuf.positionRef
when
(pageResidency + defResidency >= options.pageSize)
(flushPage options scratchBuffer columnChunkState)
```
## 后续工作
目前,Parquet写入器仅支持:
- Snappy和Uncompressed压缩
- Plain编码
- 定义级别最高到1
下一步首先是支持Parquet提供的全部压缩和编码方式。当前写入器完全是单线程的,因此我们可能会探索以有限的方式添加并发和/或并行性——只要存在性能上的合理理由。
目前,诸如列/页的统计信息和布隆过滤器等可选元数据字段未被记录。正如我们前面提到的,这些对于读取器优化查询Parquet文件中的数据非常有用。
最后,目前我们所有的缓冲区在等待刷新到磁盘时都存在于内存中。在内存受限的系统中,编写需要特别大的页/行组的Parquet文件时,我们最终可能会耗尽内存(因为一个完整的行组必须完全保存在内存中)。因此,我们计划实现一种两阶段策略,使用更小的内存缓冲区,并将数据写入磁盘临时文件,而不是保留在内存中。
相似文章
从单个 Parquet 文件构建快速下钻仪表板
本文展示了利用 Hyparquet JavaScript 阅读器和 HTTP 范围请求,从单个 Parquet 数据立方体构建快速下钻仪表板的方法,从而无需使用传统数据库。
Parquet 中定长列表的快速路径
这篇博客文章介绍了 Apache Parquet 的一项优化,用于高效存储和解码像向量嵌入这样的定长列表,通过绕过固定大小数据页的 Dremel 重构,实现了与扁平列相当的解码性能。
将Postgres数据以Parquet格式存储在S3上:LTAP架构解析
Databricks推出Lakebase LTAP架构,将Postgres数据以Parquet格式存储在S3上,无需CDC或镜像即可在单份数据上实现事务与分析。
Apache Parquet 中的自适应无损浮点编码
自适应无损浮点(ALP)编码是一种针对 Apache Parquet 中浮点数据的新轻量级编码,提供与 zstd 相似的压缩比,同时具有更快的解压速度、随机访问支持以及对 GPU/SIMD 友好的解码。
Haskell中的数据导向编程(SICP 2.4.3)
本文展示了在Haskell中实现数据导向编程来处理复数运算,遵循SICP 2.4.3的方法,避免在添加新表示时修改通用函数。