Skip to content
HN On Hacker News ↗

Writing Parquet Files Using Haskell

▲ 89 points • 33 comments • by cosmic_quanta • 3w ago • HN discussion ↗

Pangram verdict · v3.3

We believe that this entire text is human-written.

2 %

AI likelihood · overall

Human
100% human-written 0% AI-generated
SEGMENTS · HUMAN 1 of 1
SEGMENTS · AI 0 of 1
WORD COUNT 1,624
PEAK AI % 2% · §1
Analyzed
Sep 25
backend: pangram/v3.3
Segments scanned
1 windows
avg 1624 words each
Distribution
100 / 0%
human / AI fraction
Verdict
Human
Pangram v3.3

Article text · 1,624 words · 1 segments analyzed

Human AI-generated
§1 Human · 2%

We implemented a parquet writer in DataHaskell/Dataframe. Using it is simple; you need only pass your dataframe into the writeParquet function which writes a parquet file with sane defaults for row group and page sizes. An example: 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" If you need more fine-grained control over the parquet file you’ll want to use writeParquetWithOptions. Read on to see what those options are and how they affect the final file. For Haskell to interoperate with the data ecosystem, it must be able to understand the standard formats in use by that ecosystem. For a long time our options for serializing data1 in Haskell came down to CSV, JSON, or to simply dump it into a ByteString. While CSV and JSON have their place, they come with considerable disadvantages as the volume of data expands. They are slow and cumbersome to read, write, and query, and any custom homegrown formats tend to lack interoperability with standard data science tools. The parquet format trades simplicity for efficient storage and querying of the data stored therein. We want our data to have high compression ratios and we want to minimize reading data irrelevant to our specific query. The ability to read and write parquet files for long term storage, sending over the network, or for interop with other programs, especially given the universality of Parquet in the data science ecosystem, is a rather useful tool that we should like to have in a Dataframe library. Parqué? Feel free to skip this section if you already know the structure of a parquet file. Parquet files are a series of row groups followed by metadata at the end of the file with information that allows readers to locate relevant column chunks and pages. It also contains useful information like statistics and bloom filters so that, for example, a reader / query planner can decide whether or not to read a specific row group. Each row group is a collection of column chunks, each of which contain the same number of rows. Each column chunk is a series of data pages. Since each column chunk is a series of pages, and each row group is a series of column chunks, the final file simply looks like a series of pages from each column one after the other. We’re able to make sense of it all by using the metadata to identify the offset and size of each row group and column chunk. A data page is where we actually store all of our data. It consists of first the page metadata, describing its encoding, the number of values, the statistics, among other things. There are actually two versions of the data page with subtle differences. Next we have definition levels and repetition levels. They’re an inexpensive encoding of nullability and nested structure. A detailed description of definition levels and repetition levels is out of scope for this article; refer to the Dremel paper for that. For our purposes we currently only support writing definition levels up to one to denote nullable values. Finally we have our actual encoded and compressed data (depending on which data page we’re using we compress either just the data or both the definition/repetition levels and the data). The encoding is determined per page, while compression is determined per column chunk. The Implementation What follows is a technical discussion of the design of the Parquet writer and the tradeoffs we made. Write Options It follows from the description of the parquet format above that a writer needs to expose the appropriate set of fiddle factors to allow the user to tune the writer to produce a parquet file that is efficient for the specific data it is being applied to. That is, we want a parquet writer that compresses well with appropriately sized row groups and pages, so that a reader can effectively apply parallelism, projection pushdown (read only relevant column chunks), predicate pushdown (read only relevant row groups/ prune irrelevant pages), IO pushdown (read only relevant files, assuming the data is chunked into multiple files), etc2. data ParquetWriteOptions = ParquetWriteOptions { pageSize :: !Int , rowGroupSize :: !Int , batchRows :: !Int , subBatchRows :: !Int , compressionCodec :: !CompressionCodec , strategy :: !WriterStrategy , maxRowsPerFile :: !(Maybe Int) } Both pageSize and rowGroupSize are the target size in bytes of each page and each row group. But each column chunk in a row group must contain the same number of rows, and depending on the specific data being encoded, the encoding, and the compression algorithm we’re using, each column chunk will hold a different number of rows before reaching the target size. The same holds true for each page. So how do we ensure our page size target, our row group size target, and have the same number of rows in each column chunk? We must consider both the target pageSize and rowGroupSize to be best effort; they could be somewhat above or below the target. We run our columns through batches of size batchRows and check the state of the row group after each batch. So each row group contains an integer multiple of batchRows rows; pages also will contain an integer multiple of subBatchRows rows (with the exception of the final page and the final column chunk). Sub-batching at the page level allows us to reduce the amount of IORef book-keeping we have to do after each write, granting us significant speedups.3 Memory Due to the constraints described in the previous section, it’s difficult to know how big to make our buffers ahead of time. This is even more true of column chunk buffers – each column will fit a very different number of rows / pages in the same amount of data, and we may expect that some column chunk buffers will be significantly bigger than others and dominate the space per row group. So we must be able to grow our memory. One of the more convenient ways of handling raw memory without resorting to the FFI is MutableByteArray. data MemoryBuffer = MemoryBuffer { arrayRef :: !(IORef (MutableByteArray RealWorld)) , positionRef :: !(IORef Int) } For our implementation we use pinned ByteArrays, as we would like to convert it into a Ptr Word8 when it’s time to flush into either another buffer, for example when flushing a page buffer into a column chunk buffer, or into a file, as we would when flushing a row group to file. Using pinned ByteArrays results in a slight complication when trying to grow the memory buffer. We must not use the grow function provided by Data.Primitive, instead we must allocate a new pinned ByteArray and allow the old one to be GCed. One might be worried about heap fragmentation because a single pinned object in a 4KB GHC block can keep the whole block alive but we expect that our buffers will tend to be much larger than 4KB. Furthermore grows ought to be rare especially after the first few Pages and the first RowGroup. 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 #-} We also write helper functions for writing Word8s, Word32s, Word64s, Int32s, Int64s, Integers, Floats, Doubles, and ByteStrings to a buffer. Finally, we have a flushBufferToBuffer :: MemoryBuffer -> MemoryBuffer -> IO () and a flushBufferToFile :: WritableBinaryHandle -> MemoryBuffer -> IO ()4. The Core Loop Here we provide a high level sketch of the main loop of the Parquet Writer. There is considerable complexity we omit in favor of explaining the design as simply as possible. For the full implementation refer to Writer.hs. At a high level we can think of the Parquet writer as an effectful fold over the dataframe5 (Also known as reduce in some languages or, more rare, accumulate). A fold can be thought of as a loop, that is, we iterate over the rows in our dataframe. Folds are usually pure, but ‘effectful’ means that each iteration produces some kind of side-effect, in our case writing to and from memory and to disk. Finally the fold produces a final singular result from the iteration, and that is the file metadata that we must append to the end of the file. First we must define the state threaded through the writer to accomplish this: data ParquetWriterState = ParquetWriterState { outputFileHandle :: !WritableBinaryHandle -- newtype over a handle specifically for us to write to , columnChunks :: !(VB.Vector ColumnChunkState) -- effectively our row group buffer , 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 -- Each column has its own pageBuffer , definitionLevels :: !DefLevels , currentRowCount :: !(IORef Int) } -- In 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) -- in case cleanup or post processing is required }