プログラミング
HaskellでParquetファイルを書き出す
Writing Parquet files using Haskell (datahaskell.org)
要約
DataHaskellライブラリにParquetファイル書き込み機能が実装されました。これにより、Haskellから標準的なデータ形式であるParquetファイルを効率的に生成できるようになります。記事では、Parquetフォーマットの構造、書き込みオプション、およびメモリ管理の実装詳細について解説しています。
全文翻訳
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"
より詳細な制御が必要な場合は、writeParquetWithOptionsを使用します。これらのオプションが最終的なファイルにどのように影響するかは、読み進めてください。
Parquet For Haskell
Haskellがデータエコシステムと相互運用するためには、そのエコシステムで使用されている標準フォーマットを理解できる必要があります。長らくHaskellでデータをシリアライズする選択肢は、CSV、JSON、または単にByteStringにダンプするだけでした。CSVとJSONにはそれぞれの用途がありますが、データ量が増加するにつれて、かなりの欠点があります。それらは読み書きやクエリが遅く煩雑であり、カスタムの自作フォーマットは標準的なデータサイエンスツールとの相互運用性に欠ける傾向があります。
Parquetフォーマットは、効率的なストレージとクエリのためにシンプルさを犠牲にします。圧縮率が高く、特定のクエリに関連しないデータの読み込みを最小限に抑えたいと考えています。長期保存、ネットワーク経由での送信、または他のプログラムとの相互運用(特にデータサイエンスエコシステムにおけるParquetの普遍性を考えると)のためにParquetファイルを読み書きできることは、Dataframeライブラリに欲しい非常に便利なツールです。
Parqué?
Parquetファイルの構造をすでに知っている場合は、このセクションをスキップしても構いません。
Parquetファイルは、一連の行グループと、その後にファイルの末尾にあるメタデータで構成されており、メタデータには関連する列チャンクとページを特定するための情報が含まれています。また、統計情報やブルームフィルターなどの有用な情報も含まれているため、例えばリーダー/クエリプランナーは特定の行グループを読み取るべきかどうかを決定できます。
各行グループは列チャンクの集まりであり、各列チャンクは同じ数の行を含んでいます。各列チャンクはデータページのシリーズです。各列チャンクはページのシリーズであり、各行グループは列チャンクのシリーズであるため、最終的なファイルは単に各列のページが次々と並んだものになります。メタデータを使用して、各行グループと列チャンクのオフセットとサイズを特定することで、すべてを理解できます。
データページは、実際のデータを格納する場所です。まず、エンコーディング、値の数、統計情報などを記述するページメタデータで構成されます。データページには実際には2つのバージョンがあり、微妙な違いがあります。次に、定義レベルと繰り返しレベルがあります。これらは、null許容性(nullable)とネスト構造の非コストなエンコーディングです。定義レベルと繰り返しレベルの詳細な説明は、この記事の範囲外です。Dremel論文を参照してください。私たちの目的では、現在、null許容値を表すために定義レベルを1までしかサポートしていません。最後に、実際のエンコードおよび圧縮されたデータ(使用するデータページによっては、データのみを圧縮するか、定義/繰り返しレベルとデータの両方を圧縮します)があります。エンコーディングはページごとに決定され、圧縮は列チャンクごとに決定されます。
実装
以下は、Parquetライターの設計と、私たちが下したトレードオフに関する技術的な議論です。
書き込みオプション
上記のParquetフォーマットの説明から、ユーザーが特定のデータに効率的なParquetファイルを生成するためにライターを調整できるように、適切な調整可能な要素のセットを公開する必要があることがわかります。つまり、リーダーが並列処理、射影プッシュダウン(関連する列チャンクのみを読み取る)、述語プッシュダウン(関連する行グループのみを読み取る/無関係なページをプルーニングする)、IOプッシュダウン(データが複数のファイルにチャンク化されていると仮定して、関連するファイルのみを読み取る)などを効果的に適用できるように、適切にサイズ設定された行グループとページでうまく圧縮されるParquetライターを望んでいます。
data ParquetWriteOptions = ParquetWriteOptions {
pageSize :: !Int,
rowGroupSize :: !Int,
batchRows :: !Int,
subBatchRows :: !Int,
compressionCodec :: !CompressionCodec,
strategy :: !WriterStrategy,
maxRowsPerFile :: !(Maybe Int)
}
pageSizeとrowGroupSizeは、各ページと各行グループのターゲットバイトサイズです。しかし、行グループ内の各列チャンクは同じ数の行を含んでいる必要があり、エンコーディング、圧縮アルゴリズム、およびエンコードされるデータの種類によっては、ターゲットサイズに達する前に各列チャンクが異なる数の行を保持します。ページについても同様です。では、ページサイズのターゲットと行グループサイズのターゲットをどのように確保し、各列チャンクに同じ数の行を持たせるのでしょうか?ターゲットpageSizeとrowGroupSizeの両方を最善努力(best effort)と見なす必要があります。ターゲットよりやや大きいか小さい可能性があります。
バッチサイズbatchRowsで列を処理し、各バッチ後に the row group の状態を確認します。したがって、各行グループはbatchRowsの整数倍の行を含みます。ページもsubBatchRowsの整数倍の行を含みます(最後のページと最後の列チャンクを除く)。ページレベルでのサブバッチ処理により、書き込みごとにIORefのブックキーピングの量を減らすことができ、大幅な速度向上が得られます。
メモリ
前のセクションで説明した制約のため、事前にバッファのサイズをどれくらいにするかを知るのは困難です。これは、列チャンクバッファの場合にはさらに当てはまります。各列は、同じ量のデータで非常に異なる数の行/ページを保持でき、一部の列チャンクバッファが他のものより著しく大きくなり、行グループあたりのスペースを支配すると予想される場合があります。したがって、メモリを拡張できる必要があります。
FFIに頼ることなく生のメモリを扱う便利な方法の1つは、MutableByteArrayです。
data MemoryBuffer = MemoryBuffer {
arrayRef :: !(IORef (MutableByteArray RealWorld)),
positionRef :: !(IORef Int)
}
実装では、ページバッファを列チャンクバッファにフラッシュする場合や、行グループをファイルにフラッシュする場合のように、Ptr Word8に変換したいので、ピン留めされたByteArrayを使用します。ピン留めされたByteArrayを使用すると、メモリバッファを拡張する際にわずかな複雑さが生じます。Data.Primitiveが提供するgrow関数を使用することはできません。代わりに、新しいピン留めされたByteArrayを割り当て、古いByteArrayはGCに任せる必要があります。
単一のピン留めされたオブジェクトが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ライターのメインループのハイレベルなスケッチを提供します。かなりのコンピュ