From 151b941bd3639ffaa67a249dc7c716f8978b623a Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sun, 27 Sep 2026 13:07:46 +0100 Subject: [PATCH 01/51] Remove unused ByteString import --- dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs | 1 - 1 file changed, 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs index 7acc7e8a..ab496c88 100644 --- a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs +++ b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs @@ -41,7 +41,6 @@ import Control.Monad.IO.Class (MonadIO (..)) import Control.Monad.Primitive (RealWorld) import Control.Monad.ST (stToIO) import Data.Bits (shiftR) -import qualified Data.ByteString as BS import Data.ByteString.Internal (ByteString (PS), create) import qualified Data.ByteString.Unsafe as BU import Data.IORef (IORef, newIORef, readIORef, writeIORef) From ba2247799bf7b2dabe0c79894e99eab4081cea2d Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 15:59:09 +0100 Subject: [PATCH 02/51] Wrap Encoder encodeValue type --- dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs index 0df7b76d..3eafba92 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs @@ -59,7 +59,8 @@ data Encoder = Encoder { encType :: !ThriftType , convertedType :: !(Maybe ConvertedType) , logicalType :: !(Maybe LogicalType) - , encodeValue :: !(MemoryBuffer -> Int -> Int -> IO (Int, Bool)) + , encodeValue :: + !(MemoryBuffer -> Int -> Int -> IO (Int, Bool)) , finishValues :: !(MemoryBuffer -> Int -> IO Int) } From ac7a6bf63aab72f05b83203811ddb484dd044130 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 15:59:50 +0100 Subject: [PATCH 03/51] Wrap ensureCapacity type --- dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs index ab496c88..a6aec4ef 100644 --- a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs +++ b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs @@ -220,7 +220,8 @@ mallocBuffer capacity -- is just a matter of adding a new buffer to the array (which we can -- pre-allocate to three elements to begin with and grow it only on the -- off chance that a buffer required more than three grows). -ensureCapacity :: MemoryBuffer -> Int -> IO (MutableByteArray RealWorld) +ensureCapacity :: + MemoryBuffer -> Int -> IO (MutableByteArray RealWorld) ensureCapacity buffer needed = do array <- readIORef buffer.arrayRef maxSize <- getSizeofMutableByteArray array From b91728f8a71b3314a1c57c40f9d51c79cfcd625d Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 16:01:35 +0100 Subject: [PATCH 04/51] Wrap Parquet writer reference import --- dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs index 2e8d4e9f..4065466b 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs @@ -14,7 +14,13 @@ module DataFrame.IO.Parquet.Writer ( import Control.Monad (forM_, unless, when) import qualified Data.ByteString as BS -import Data.IORef (IORef, modifyIORef', newIORef, readIORef, writeIORef) +import Data.IORef ( + IORef, + modifyIORef', + newIORef, + readIORef, + writeIORef, + ) import Data.Int (Int64) import Data.Maybe (fromJust) import Data.Primitive.ByteArray (getSizeofMutableByteArray) From 7a472eebbd698f2f489f1e764e7f290d92abf51b Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 16:01:42 +0100 Subject: [PATCH 05/51] Wrap writeRows type --- dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs index 4065466b..750b9b22 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs @@ -222,7 +222,12 @@ nativeTypeKeyValues names df = ] writeRows :: - ParquetWriteOptions -> MemoryBuffer -> Int -> Int -> ColumnChunkState -> IO () + ParquetWriteOptions -> + MemoryBuffer -> + Int -> + Int -> + ColumnChunkState -> + IO () writeRows options scratch firstRow count ccs = do let page = ccs.pageState buf = page.pageBuffer From 740c6576d818a2dbcd8faca6c46bca3162e188ab Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 16:01:49 +0100 Subject: [PATCH 06/51] Wrap flushPage type --- dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs index 750b9b22..9d4b4446 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs @@ -265,7 +265,8 @@ writeRows options scratch firstRow count ccs = do (pageRes + defRes >= options.pageSize) (flushPage options scratch ccs) -flushPage :: ParquetWriteOptions -> MemoryBuffer -> ColumnChunkState -> IO () +flushPage :: + ParquetWriteOptions -> MemoryBuffer -> ColumnChunkState -> IO () flushPage options scratch columnChunkState = do let page = columnChunkState.pageState numPageRows <- readIORef page.currentRowCount From 98b3dfc65687c421ce836980a7f8bfbf5bda088b Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 16:01:58 +0100 Subject: [PATCH 07/51] Wrap assemblePageBody type --- dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs index 9d4b4446..4d0e7bc9 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs @@ -281,7 +281,8 @@ flushPage options scratch columnChunkState = do resetPosition scratch writeIORef page.currentRowCount 0 -assemblePageBody :: MemoryBuffer -> ColumnChunkState -> IO MemoryBuffer +assemblePageBody :: + MemoryBuffer -> ColumnChunkState -> IO MemoryBuffer assemblePageBody scratch columnChunkState | not columnChunkState.nullable = pure columnChunkState.pageState.pageBuffer | otherwise = do From 45f9f84c51869dbbd4fd6bba30737220b055c269 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 16:02:07 +0100 Subject: [PATCH 08/51] Wrap flushRowGroup type --- dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs index 4d0e7bc9..07908357 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs @@ -316,7 +316,10 @@ writeDataPage codec numPageRows body columnChunkState = do columnChunkState.uncompressedBufferSize (+ fromIntegral (BS.length headerBytes + uncompressedPageSize)) -flushRowGroup :: ParquetWriteOptions -> ParquetWriterState -> IO () +flushRowGroup :: + ParquetWriteOptions -> + ParquetWriterState -> + IO () flushRowGroup options writerState = do rowNumber <- readIORef writerState.rowNumberRef when (rowNumber > 0) $ do From 2e2728f4bdd439ebc3a94ce6d8e2fb91bb5b0604 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 16:02:26 +0100 Subject: [PATCH 09/51] Wrap encoder text copy --- .../src/DataFrame/IO/Parquet/Writer/Encoder.hs | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs index 3eafba92..091d35c8 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs @@ -357,7 +357,13 @@ textEncoder col = writeWord32At buffer pos (fromIntegral count) arr <- readIORef buffer.arrayRef withMutableByteArrayContents arr $ \ptr -> - stToIO (TA.copyToPointer bytes offset (ptr `plusPtr` (pos + 4)) count) + stToIO + ( TA.copyToPointer + bytes + offset + (ptr `plusPtr` (pos + 4)) + count + ) pure (pos + 4 + count) mismatch = error From ac33c521650a1623b0afddb07a20690f8dfaf63d Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 16:02:35 +0100 Subject: [PATCH 10/51] Wrap mallocBuffer type --- dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs index a6aec4ef..4d657e12 100644 --- a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs +++ b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs @@ -198,7 +198,8 @@ data MemoryBuffer = MemoryBuffer , positionRef :: !(IORef Int) } -mallocBuffer :: Int -> IO MemoryBuffer +mallocBuffer :: + Int -> IO MemoryBuffer mallocBuffer capacity | capacity < 0 = ioError $ userError "mallocBuffer: negative capacity" | otherwise = do From 703dc53a1d404a5254b1d68551996449dddae20a Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 16:02:44 +0100 Subject: [PATCH 11/51] Wrap writeByteString type --- dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs index 4d657e12..4beff746 100644 --- a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs +++ b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs @@ -244,7 +244,8 @@ writeWord8 buffer b = do writeIORef buffer.positionRef (position + 1) {-# INLINE writeWord8 #-} -writeByteString :: MemoryBuffer -> ByteString -> IO () +writeByteString :: + MemoryBuffer -> ByteString -> IO () writeByteString buffer bs = BU.unsafeUseAsCStringLen bs $ \(source, len) -> do position <- readIORef buffer.positionRef From 3e3da7fe913369a77c523574b60acd703283443d Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 16:02:52 +0100 Subject: [PATCH 12/51] Wrap writeByteString copy --- dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs index 4beff746..7a847830 100644 --- a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs +++ b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs @@ -251,7 +251,10 @@ writeByteString buffer bs = position <- readIORef buffer.positionRef array <- ensureCapacity buffer (position + len) withMutableByteArrayContents array $ \dst -> - copyBytes (dst `plusPtr` position) (castPtr source) len + copyBytes + (dst `plusPtr` position) + (castPtr source) + len writeIORef buffer.positionRef (position + len) {-# INLINE writeByteString #-} From 35027d6360d23eadaec73071e93a62d9cf219c71 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 16:03:02 +0100 Subject: [PATCH 13/51] Wrap writeWord32At type --- dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs index 7a847830..10a0bb92 100644 --- a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs +++ b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs @@ -272,7 +272,8 @@ writeWord64LE buffer w = do writeIORef buffer.positionRef (position + 8) {-# INLINE writeWord64LE #-} -writeWord32At :: MemoryBuffer -> Int -> Word32 -> IO () +writeWord32At :: + MemoryBuffer -> Int -> Word32 -> IO () writeWord32At buffer position w = do array <- ensureCapacity buffer (position + 4) writeByteArray array position (fromIntegral w :: Word8) From b2e83eb97dac7227c6587bfcc6a7d6560d781981 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 16:03:11 +0100 Subject: [PATCH 14/51] Wrap writeWord64At type --- dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs index 10a0bb92..dc9221c5 100644 --- a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs +++ b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs @@ -282,7 +282,8 @@ writeWord32At buffer position w = do writeByteArray array (position + 3) (fromIntegral (w `shiftR` 24) :: Word8) {-# INLINE writeWord32At #-} -writeWord64At :: MemoryBuffer -> Int -> Word64 -> IO () +writeWord64At :: + MemoryBuffer -> Int -> Word64 -> IO () writeWord64At buffer position w = do array <- ensureCapacity buffer (position + 8) writeByteArray array position (fromIntegral w :: Word8) From 2d4e618d395e9c8a6b59a317dee89bbd08b50cb4 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 16:03:18 +0100 Subject: [PATCH 15/51] Wrap writeInteger64 type --- dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs index dc9221c5..7f360579 100644 --- a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs +++ b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs @@ -296,7 +296,8 @@ writeWord64At buffer position w = do writeByteArray array (position + 7) (fromIntegral (w `shiftR` 56) :: Word8) {-# INLINE writeWord64At #-} -writeInteger64 :: MemoryBuffer -> Integer -> IO () +writeInteger64 :: + MemoryBuffer -> Integer -> IO () writeInteger64 buffer value = do position <- readIORef buffer.positionRef newPosition <- writeInteger64At buffer position value From fcdb3cbb2855af57eb91b4f6bad6806339e5b72d Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 16:03:29 +0100 Subject: [PATCH 16/51] Wrap writeInteger64At type --- dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs index 7f360579..8eaafd04 100644 --- a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs +++ b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs @@ -304,7 +304,8 @@ writeInteger64 buffer value = do writeIORef buffer.positionRef newPosition {-# INLINE writeInteger64 #-} -writeInteger64At :: MemoryBuffer -> Int -> Integer -> IO Int +writeInteger64At :: + MemoryBuffer -> Int -> Integer -> IO Int writeInteger64At buffer position value | value < toInteger (minBound :: Int64) = outOfRange | value > toInteger (maxBound :: Int64) = outOfRange From b6616d81ba88b46030d572e354483152998f692f Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 16:03:41 +0100 Subject: [PATCH 17/51] Wrap flushBufferToBuffer type --- dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs index 8eaafd04..1a2d32d0 100644 --- a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs +++ b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs @@ -325,7 +325,8 @@ writeDoubleLE :: MemoryBuffer -> Double -> IO () writeDoubleLE buffer = writeWord64LE buffer . castDoubleToWord64 {-# INLINE writeDoubleLE #-} -flushBufferToBuffer :: MemoryBuffer -> MemoryBuffer -> IO () +flushBufferToBuffer :: + MemoryBuffer -> MemoryBuffer -> IO () flushBufferToBuffer source destination | source.arrayRef == destination.arrayRef = pure () | otherwise = do From 532b14745f24023ce4ef83df8077d08e0314732f Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 16:03:48 +0100 Subject: [PATCH 18/51] Wrap bufferToByteString type --- dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs index 1a2d32d0..8837fb19 100644 --- a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs +++ b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs @@ -345,7 +345,8 @@ flushBufferToBuffer source destination writeIORef source.positionRef 0 {-# INLINE flushBufferToBuffer #-} -bufferToByteString :: MemoryBuffer -> IO ByteString +bufferToByteString :: + MemoryBuffer -> IO ByteString bufferToByteString buffer = do array <- readIORef buffer.arrayRef position <- readIORef buffer.positionRef From 07f8067fe377b7dea08c84858e0ce2d7ac795b78 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 16:03:57 +0100 Subject: [PATCH 19/51] Wrap flushBufferToFile type --- dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs index 8837fb19..afd4cf95 100644 --- a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs +++ b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs @@ -383,7 +383,8 @@ resetPosition buffer = writeIORef buffer.positionRef 0 -- So when writing to a file to minimize syscall overhead while -- trying not to create dirty pages in the kernel page cache, we'll -- be flushing in 256 KiB chunks. -flushBufferToFile :: WritableBinaryHandle -> MemoryBuffer -> IO () +flushBufferToFile :: + WritableBinaryHandle -> MemoryBuffer -> IO () flushBufferToFile (WritableBinaryHandle h) buffer = do array <- readIORef buffer.arrayRef position <- readIORef buffer.positionRef From fc736371012ba7e775792039298f869a8615b714 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 16:04:05 +0100 Subject: [PATCH 20/51] Wrap appendTextArraySlice type --- dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs index afd4cf95..b30a5c40 100644 --- a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs +++ b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs @@ -411,7 +411,8 @@ writeByteStringToFile (WritableBinaryHandle h) bs = go (offset + n) go 0 -appendTextArraySlice :: MemoryBuffer -> TA.Array -> Int -> Int -> IO () +appendTextArraySlice :: + MemoryBuffer -> TA.Array -> Int -> Int -> IO () appendTextArraySlice buffer source offset count | count < 0 = ioError $ userError "appendTextArraySlice: negative length" | otherwise = do From c51d89584aa0c0c79645bc1b8394568dd37dbced Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 16:04:14 +0100 Subject: [PATCH 21/51] Wrap appendTextArraySlice error --- dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs index b30a5c40..ae958a6c 100644 --- a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs +++ b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs @@ -414,7 +414,8 @@ writeByteStringToFile (WritableBinaryHandle h) bs = appendTextArraySlice :: MemoryBuffer -> TA.Array -> Int -> Int -> IO () appendTextArraySlice buffer source offset count - | count < 0 = ioError $ userError "appendTextArraySlice: negative length" + | count < 0 = + ioError $ userError "appendTextArraySlice: negative length" | otherwise = do position <- readIORef buffer.positionRef array <- ensureCapacity buffer (position + count) From 7aa4e9b11ca4bc75d42975a0b1ebabc9a8c5a77e Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 16:04:23 +0100 Subject: [PATCH 22/51] Wrap appendTextArraySlice copy --- dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs index ae958a6c..92d1d10d 100644 --- a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs +++ b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs @@ -420,6 +420,12 @@ appendTextArraySlice buffer source offset count position <- readIORef buffer.positionRef array <- ensureCapacity buffer (position + count) withMutableByteArrayContents array $ \destination -> - stToIO (TA.copyToPointer source offset (destination `plusPtr` position) count) + stToIO + ( TA.copyToPointer + source + offset + (destination `plusPtr` position) + count + ) writeIORef buffer.positionRef (position + count) {-# INLINE appendTextArraySlice #-} From d835097e5309128aa497c1416f426c0fd98dc5eb Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 16:47:23 +0100 Subject: [PATCH 23/51] Wrap writeShard arguments --- dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs index 07908357..cc0d9f9c 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs @@ -150,7 +150,12 @@ shardPathFor pattern_ shardIndex = -- | Write rows @[startRow, endRow)@ of the frame to a single Parquet file. writeShard :: - ParquetWriteOptions -> FilePath -> DataFrame -> Int -> Int -> IO () + ParquetWriteOptions -> + FilePath -> + DataFrame -> + Int -> + Int -> + IO () writeShard options path_ df startRow endRow = do let names = columnNames df shardRows = max 0 (endRow - startRow) From ca7f0d3b25bee9778aa420c1a50bc57736902f31 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 16:48:29 +0100 Subject: [PATCH 24/51] Remove writeBatch local signature --- dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs | 1 - 1 file changed, 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs index cc0d9f9c..da836953 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs @@ -185,7 +185,6 @@ writeShard options path_ df startRow endRow = do rowNumberRef_ interval = max 1 options.batchRows subBatch = max 1 options.subBatchRows - writeBatch :: Int -> Int -> IO () writeBatch rowNum batchEnd | rowNum >= batchEnd = pure () | otherwise = do From 4f36c274ddc50a087542c1eb81c7e88fcdef9711 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 16:48:33 +0100 Subject: [PATCH 25/51] Remove loop local signature --- dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs | 1 - 1 file changed, 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs index da836953..f56f7484 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs @@ -192,7 +192,6 @@ writeShard options path_ df startRow endRow = do VB.forM_ columnChunks_ (writeRows options scratchBuffer_ rowNum count) modifyIORef' rowNumberRef_ (+ count) writeBatch (rowNum + count) batchEnd - loop :: Int -> IO () loop rowNum | rowNum >= endRow = pure () | otherwise = do From 8b80dc37af402046323c2a26ad43e6acf735b0c2 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 16:49:32 +0100 Subject: [PATCH 26/51] Wrap flushPage arguments --- dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs index f56f7484..0a162bc6 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs @@ -269,7 +269,10 @@ writeRows options scratch firstRow count ccs = do (flushPage options scratch ccs) flushPage :: - ParquetWriteOptions -> MemoryBuffer -> ColumnChunkState -> IO () + ParquetWriteOptions -> + MemoryBuffer -> + ColumnChunkState -> + IO () flushPage options scratch columnChunkState = do let page = columnChunkState.pageState numPageRows <- readIORef page.currentRowCount From 489ce928cc5f4032fa00c4403bb22d08c1696d87 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 16:49:37 +0100 Subject: [PATCH 27/51] Wrap assemblePageBody arguments --- dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs index 0a162bc6..3c88a0df 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs @@ -288,7 +288,9 @@ flushPage options scratch columnChunkState = do writeIORef page.currentRowCount 0 assemblePageBody :: - MemoryBuffer -> ColumnChunkState -> IO MemoryBuffer + MemoryBuffer -> + ColumnChunkState -> + IO MemoryBuffer assemblePageBody scratch columnChunkState | not columnChunkState.nullable = pure columnChunkState.pageState.pageBuffer | otherwise = do From 2ee36a7be3af05c004dac372522123519aa76f42 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 16:49:44 +0100 Subject: [PATCH 28/51] Wrap writeDataPage arguments --- dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs index 3c88a0df..4d670e08 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs @@ -304,7 +304,11 @@ assemblePageBody scratch columnChunkState pure scratch writeDataPage :: - CompressionCodec -> Int -> MemoryBuffer -> ColumnChunkState -> IO () + CompressionCodec -> + Int -> + MemoryBuffer -> + ColumnChunkState -> + IO () writeDataPage codec numPageRows body columnChunkState = do uncompressedPageSize <- bufferResidency body compressedBody <- case codec of From bab7f2a2bb0ce09e746e13876b8fc7bc82dc0200 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 17:59:51 +0100 Subject: [PATCH 29/51] Wrap writeByteString arguments --- dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs index 92d1d10d..6165ce58 100644 --- a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs +++ b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs @@ -245,7 +245,9 @@ writeWord8 buffer b = do {-# INLINE writeWord8 #-} writeByteString :: - MemoryBuffer -> ByteString -> IO () + MemoryBuffer -> + ByteString -> + IO () writeByteString buffer bs = BU.unsafeUseAsCStringLen bs $ \(source, len) -> do position <- readIORef buffer.positionRef From 223f00a201da20fc3f4fb6507d5849c7ec071340 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 17:59:58 +0100 Subject: [PATCH 30/51] Wrap bufferToByteString arguments --- dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs index 6165ce58..c33a4738 100644 --- a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs +++ b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs @@ -348,7 +348,8 @@ flushBufferToBuffer source destination {-# INLINE flushBufferToBuffer #-} bufferToByteString :: - MemoryBuffer -> IO ByteString + MemoryBuffer -> + IO ByteString bufferToByteString buffer = do array <- readIORef buffer.arrayRef position <- readIORef buffer.positionRef From cbb75b3429277a69c671d5ba36aeba3e52a7e905 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 18:14:01 +0100 Subject: [PATCH 31/51] Parenthesize integer range error --- dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs index c33a4738..92adb325 100644 --- a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs +++ b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs @@ -316,7 +316,7 @@ writeInteger64At buffer position value pure (position + 8) where outOfRange = - ioError (userError "writeParquet: Integer value is outside the INT64 range") + (ioError (userError "writeParquet: Integer value is outside the INT64 range")) {-# INLINE writeInteger64At #-} writeFloatLE :: MemoryBuffer -> Float -> IO () From 4dd8d1371fc59016f13e55dcd51af4a7b21cd026 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 19:58:17 +0100 Subject: [PATCH 32/51] Wrap row group flush call --- dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs index 4d670e08..70554e22 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs @@ -198,7 +198,8 @@ writeShard options path_ df startRow endRow = do let batchEnd = rowNum + min interval (endRow - rowNum) writeBatch rowNum batchEnd size <- bufferedSize columnChunks_ - when (size >= options.rowGroupSize) (flushRowGroup options writerState) + when (size >= options.rowGroupSize) $ + flushRowGroup options writerState loop batchEnd loop startRow flushRowGroup options writerState From 057bb929a4b35218448413cf0bea9aa28d3b0de6 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 19:58:26 +0100 Subject: [PATCH 33/51] Wrap row group metadata binding --- dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs index 70554e22..806f510a 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs @@ -203,7 +203,8 @@ writeShard options path_ df startRow endRow = do loop batchEnd loop startRow flushRowGroup options writerState - rowGroupMetadata <- reverse <$> readIORef rowGroupMetadataRef_ + rowGroupMetadata <- + reverse <$> readIORef rowGroupMetadataRef_ let schemaElements = rootSchemaElement (VB.length columnChunks_) : VB.toList (VB.map schema columnChunks_) From 5a4b445a81954fb64b82857866cbb01d1999474f Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 20:52:40 +0100 Subject: [PATCH 34/51] Wrap capacity growth arguments --- dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs index 806f510a..bc0f3260 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs @@ -250,7 +250,10 @@ writeRows options scratch firstRow count ccs = do | pos + margin > size = do -- Rare: buffer nearly full, grow it writeIORef buf.positionRef pos - arr' <- ensureCapacity buf (pos + max margin ((end - row) * 64)) + arr' <- + ensureCapacity + buf + (pos + max margin ((end - row) * 64)) size' <- getSizeofMutableByteArray arr' go size' pos row | otherwise = do From d81aec7fe5ce42f7ac11ac70e3ee472b30ee53f6 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 20:52:50 +0100 Subject: [PATCH 35/51] Wrap writeDataPage arguments --- dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs index bc0f3260..d54029f0 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs @@ -286,7 +286,11 @@ flushPage options scratch columnChunkState = do pos' <- columnChunkState.encoder.finishValues page.pageBuffer pos writeIORef page.pageBuffer.positionRef pos' body <- assemblePageBody scratch columnChunkState - writeDataPage options.compressionCodec numPageRows body columnChunkState + writeDataPage + options.compressionCodec + numPageRows + body + columnChunkState resetPosition page.pageBuffer resetPosition page.definitionLevels.dlBuf resetPosition scratch From b079c88e1777cfc72d0cc0f96b116bbe0af65cd9 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 20:53:04 +0100 Subject: [PATCH 36/51] Wrap row group flush operations --- .../src/DataFrame/IO/Parquet/Writer.hs | 18 +++++++++++++----- 1 file changed, 13 insertions(+), 5 deletions(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs index d54029f0..2f9cda54 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs @@ -350,14 +350,22 @@ flushRowGroup options writerState = do (reversedColumnChunks, totalCompressed, totalUncompressed) <- VB.foldM' ( \(acc, totalCompressedSize, totalUncompressedSize) columnChunkState -> do - offset <- readIORef writerState.currentFileOffsetRef - compressedSize <- bufferResidency columnChunkState.buffer - uncompressedSize <- readIORef columnChunkState.uncompressedBufferSize - flushBufferToFile writerState.outputFileHandle columnChunkState.buffer + offset <- + readIORef writerState.currentFileOffsetRef + compressedSize <- + bufferResidency columnChunkState.buffer + uncompressedSize <- + readIORef + columnChunkState.uncompressedBufferSize + flushBufferToFile + writerState.outputFileHandle + columnChunkState.buffer writeIORef writerState.currentFileOffsetRef (offset + fromIntegral compressedSize) - writeIORef columnChunkState.uncompressedBufferSize 0 + writeIORef + columnChunkState.uncompressedBufferSize + 0 let columnChunk = mkColumnChunk options.compressionCodec From a5676a1848b74cfc74dae2f36d4e4dbb4b3dac2b Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 20:53:14 +0100 Subject: [PATCH 37/51] Wrap page buffer residency call --- dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs index 2f9cda54..6a8a55f2 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs @@ -399,9 +399,11 @@ bufferedSize = VB.foldM' ( \total columnChunkState -> do chunkSize <- bufferResidency columnChunkState.buffer - valuesSize <- bufferResidency columnChunkState.pageState.pageBuffer + valuesSize <- + bufferResidency columnChunkState.pageState.pageBuffer defLevelsSize <- - bufferResidency columnChunkState.pageState.definitionLevels.dlBuf + bufferResidency + columnChunkState.pageState.definitionLevels.dlBuf pure (total + chunkSize + valuesSize + defLevelsSize) ) 0 From 7c901f8186e3a28aaa71d6d2a41a15bc469491b9 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 20:53:30 +0100 Subject: [PATCH 38/51] Extract fixed-width encoder writes --- .../DataFrame/IO/Parquet/Writer/Encoder.hs | 19 ++++++++++++------- 1 file changed, 12 insertions(+), 7 deletions(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs index 091d35c8..7fc323ab 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs @@ -72,7 +72,7 @@ buildEncoder col (INT32 enum) Nothing Nothing - (\buffer pos v -> writeWord32At buffer pos (fromIntegral v) >> pure (pos + 4)) + (\buffer pos v -> write32 buffer pos (fromIntegral v)) col | hasElemType @Int64 col = pure $ @@ -80,7 +80,7 @@ buildEncoder col (INT64 enum) Nothing Nothing - (\buffer pos v -> writeWord64At buffer pos (fromIntegral v) >> pure (pos + 8)) + (\buffer pos v -> write64 buffer pos (fromIntegral v)) col -- Ints in GHC can be 32 bit or 64 bit integers depending on the -- underlying computers architecture. So we'll do 64bit integers @@ -91,7 +91,7 @@ buildEncoder col (INT64 enum) Nothing Nothing - (\buffer pos v -> writeWord64At buffer pos (fromIntegral v) >> pure (pos + 8)) + (\buffer pos v -> write64 buffer pos (fromIntegral v)) col | hasElemType @Integer col = pure $ @@ -107,8 +107,7 @@ buildEncoder col (FLOAT enum) Nothing Nothing - ( \buffer pos v -> writeWord32At buffer pos (castFloatToWord32 v) >> pure (pos + 4) - ) + (\buffer pos v -> write32 buffer pos (castFloatToWord32 v)) col | hasElemType @Double col = pure $ @@ -116,14 +115,20 @@ buildEncoder col (DOUBLE enum) Nothing Nothing - ( \buffer pos v -> writeWord64At buffer pos (castDoubleToWord64 v) >> pure (pos + 8) - ) + (\buffer pos v -> write64 buffer pos (castDoubleToWord64 v)) col | hasElemType @Bool col = boolEncoder col | hasElemType @T.Text col = pure (textEncoder col) | hasElemType @UTCTime col = pure (timestampEncoder col) | otherwise = error ("writeParquet: unsupported column type " <> columnTypeString col) + where + write32 buffer pos value = do + writeWord32At buffer pos value + pure (pos + 4) + write64 buffer pos value = do + writeWord64At buffer pos value + pure (pos + 8) scalarEncoder :: forall a. From 49fbd4826a70cf01e35bf8c8a49af5f50773ebc0 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 20:53:42 +0100 Subject: [PATCH 39/51] Wrap textEncoder signature arguments --- dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs index 7fc323ab..e99eef52 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs @@ -326,7 +326,9 @@ boolEncoder col = do pure (Encoder (BOOLEAN enum) Nothing Nothing (columnWriter @Bool col addBit) finish) -textEncoder :: Column -> Encoder +textEncoder :: + Column -> + Encoder textEncoder col = Encoder (BYTE_ARRAY enum) From bb1ea7740ca8d1ecfc4a75f20e8748441a43b52e Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 20:53:49 +0100 Subject: [PATCH 40/51] Wrap timestampEncoder signature arguments --- dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs index e99eef52..141b8efa 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs @@ -376,7 +376,9 @@ textEncoder col = error ("writeParquet: incompatible text representation for " <> columnTypeString col) -timestampEncoder :: Column -> Encoder +timestampEncoder :: + Column -> + Encoder timestampEncoder col = Encoder (INT64 enum) From 6d421aec1198e34db865f43bc9d9e3e06eed5e28 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 20:54:03 +0100 Subject: [PATCH 41/51] Wrap SNAPPY compression branch --- dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs index 6a8a55f2..a146cebe 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs @@ -322,7 +322,8 @@ writeDataPage codec numPageRows body columnChunkState = do uncompressedPageSize <- bufferResidency body compressedBody <- case codec of UNCOMPRESSED _ -> pure Nothing - SNAPPY _ -> Just . Snappy.compress <$> bufferToByteString body + SNAPPY _ -> + Just . Snappy.compress <$> bufferToByteString body other -> error ("writeParquet: unsupported codec " <> show other) let compressedPageSize = maybe uncompressedPageSize BS.length compressedBody headerBytes = From 82026a71e6cf417a63277de4ddd8a1a58c383d96 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 21:40:49 +0100 Subject: [PATCH 42/51] Wrap writable binary file callback signature --- dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs index 92adb325..648fc4a0 100644 --- a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs +++ b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs @@ -187,7 +187,10 @@ atomicallyWriteFile path action = pure tmpFile ) -withWritableBinaryFile :: FilePath -> (WritableBinaryHandle -> IO a) -> IO a +withWritableBinaryFile :: + FilePath -> + (WritableBinaryHandle -> IO a) -> + IO a withWritableBinaryFile filepath = bracket (openWritableBinaryFile filepath) From 2caf74fea5008da2b61bc4fa13f770c1e87543c5 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 21:40:57 +0100 Subject: [PATCH 43/51] Wrap atomic file callback signature --- dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs index 648fc4a0..eb65cc39 100644 --- a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs +++ b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs @@ -162,7 +162,10 @@ openWritableBinaryFile filepath = do hSetBuffering h NoBuffering pure . WritableBinaryHandle $ h -atomicallyWriteFile :: FilePath -> (FilePath -> IO a) -> IO a +atomicallyWriteFile :: + FilePath -> + (FilePath -> IO a) -> + IO a atomicallyWriteFile path action = bracketOnError openAction From bd75b9f4292b741cd2ca08c2959a94764c0b1a25 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 21:43:03 +0100 Subject: [PATCH 44/51] Wrap bufferedSize signature --- dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs index a146cebe..5dda4a9a 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs @@ -395,7 +395,9 @@ flushRowGroup options writerState = do ) writeIORef writerState.rowNumberRef 0 -bufferedSize :: VB.Vector ColumnChunkState -> IO Int +bufferedSize :: + VB.Vector ColumnChunkState -> + IO Int bufferedSize = VB.foldM' ( \total columnChunkState -> do From 88b13aef79f761e6299c331a4fa942c3da1335e2 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sat, 26 Sep 2026 21:43:41 +0100 Subject: [PATCH 45/51] Wrap initColumnChunkState signature --- dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs index 5dda4a9a..244ac8ff 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs @@ -412,7 +412,10 @@ bufferedSize = 0 initColumnChunkState :: - ParquetWriteOptions -> T.Text -> Column -> IO ColumnChunkState + ParquetWriteOptions -> + T.Text -> + Column -> + IO ColumnChunkState initColumnChunkState options columnName_ column = do encoder_ <- buildEncoder column let nullable_ = hasMissing column From 7aa36314c649e135d421538de3576bc4f8ee9372 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Fri, 2 Oct 2026 12:38:28 +0100 Subject: [PATCH 46/51] Use ByteString length in writeByteString --- dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs index eb65cc39..b1ae8b24 100644 --- a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs +++ b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs @@ -41,6 +41,7 @@ import Control.Monad.IO.Class (MonadIO (..)) import Control.Monad.Primitive (RealWorld) import Control.Monad.ST (stToIO) import Data.Bits (shiftR) +import qualified Data.ByteString as BS import Data.ByteString.Internal (ByteString (PS), create) import qualified Data.ByteString.Unsafe as BU import Data.IORef (IORef, newIORef, readIORef, writeIORef) @@ -255,8 +256,9 @@ writeByteString :: ByteString -> IO () writeByteString buffer bs = - BU.unsafeUseAsCStringLen bs $ \(source, len) -> do + BU.unsafeUseAsCStringLen bs $ \(source, _) -> do position <- readIORef buffer.positionRef + let len = BS.length bs array <- ensureCapacity buffer (position + len) withMutableByteArrayContents array $ \dst -> copyBytes From 3493dac6cc7ba52b4f5c6b6eb313d3ec818209e5 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Fri, 2 Oct 2026 12:46:25 +0100 Subject: [PATCH 47/51] Add do block to writeByteString --- dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs index b1ae8b24..5d12d048 100644 --- a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs +++ b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs @@ -255,7 +255,7 @@ writeByteString :: MemoryBuffer -> ByteString -> IO () -writeByteString buffer bs = +writeByteString buffer bs = do BU.unsafeUseAsCStringLen bs $ \(source, _) -> do position <- readIORef buffer.positionRef let len = BS.length bs From 72358ccc63e30e6cd762e1fa5f308fa7c3514eae Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Fri, 2 Oct 2026 12:46:29 +0100 Subject: [PATCH 48/51] Reorder writeByteString pointer scopes --- .../src/DataFrame/IO/Utils/RandomAccess.hs | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs index 5d12d048..63251adf 100644 --- a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs +++ b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs @@ -256,16 +256,16 @@ writeByteString :: ByteString -> IO () writeByteString buffer bs = do - BU.unsafeUseAsCStringLen bs $ \(source, _) -> do - position <- readIORef buffer.positionRef - let len = BS.length bs - array <- ensureCapacity buffer (position + len) - withMutableByteArrayContents array $ \dst -> + position <- readIORef buffer.positionRef + let len = BS.length bs + array <- ensureCapacity buffer (position + len) + withMutableByteArrayContents array $ \dst -> + BU.unsafeUseAsCStringLen bs $ \(source, _) -> do copyBytes (dst `plusPtr` position) (castPtr source) len - writeIORef buffer.positionRef (position + len) + writeIORef buffer.positionRef (position + len) {-# INLINE writeByteString #-} writeWord32LE :: MemoryBuffer -> Word32 -> IO () From 8e2feb9ebf492d7c3bcb2f688abbe4db3d26d191 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Fri, 2 Oct 2026 12:32:53 +0100 Subject: [PATCH 49/51] Reorder bufferToByteString pointer scopes --- dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs index 63251adf..25db397d 100644 --- a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs +++ b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs @@ -361,8 +361,8 @@ bufferToByteString :: bufferToByteString buffer = do array <- readIORef buffer.arrayRef position <- readIORef buffer.positionRef - create position $ \dst -> - withMutableByteArrayContents array $ \src -> + withMutableByteArrayContents array $ \src -> + create position $ \dst -> copyBytes dst (castPtr src) position bufferResidency :: MemoryBuffer -> IO Int From 9d7a6ae3e569cd300913ca1307a2208ac294902a Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Fri, 2 Oct 2026 12:35:20 +0100 Subject: [PATCH 50/51] Generalize Parquet writer mutation to PrimBase --- .../src/DataFrame/IO/Parquet/Writer.hs | 136 ++++++++-------- .../DataFrame/IO/Parquet/Writer/DefLevels.hs | 42 ++--- .../DataFrame/IO/Parquet/Writer/Encoder.hs | 146 ++++++++--------- .../DataFrame/IO/Parquet/Writer/Metadata.hs | 6 +- .../src/DataFrame/IO/Utils/RandomAccess.hs | 153 ++++++++++-------- 5 files changed, 256 insertions(+), 227 deletions(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs index 244ac8ff..8dca6ca7 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer.hs @@ -1,4 +1,5 @@ {-# LANGUAGE BangPatterns #-} +{-# LANGUAGE FlexibleContexts #-} {-# LANGUAGE OverloadedRecordDot #-} {-# LANGUAGE OverloadedStrings #-} @@ -13,17 +14,19 @@ module DataFrame.IO.Parquet.Writer ( ) where import Control.Monad (forM_, unless, when) +import Control.Monad.IO.Class (MonadIO) +import Control.Monad.Primitive (PrimBase, PrimMonad, PrimState) import qualified Data.ByteString as BS -import Data.IORef ( - IORef, - modifyIORef', - newIORef, - readIORef, - writeIORef, - ) import Data.Int (Int64) import Data.Maybe (fromJust) import Data.Primitive.ByteArray (getSizeofMutableByteArray) +import Data.Primitive.MutVar ( + MutVar, + modifyMutVar', + newMutVar, + readMutVar, + writeMutVar, + ) import qualified Data.Text as T import qualified Data.Vector as VB import DataFrame.IO.Parquet.Thrift hiding (schema) @@ -77,29 +80,29 @@ import System.Directory (createDirectoryIfMissing) import System.FilePath (takeDirectory) import Text.Printf (printf) -data ParquetWriterState = ParquetWriterState +data ParquetWriterState m = ParquetWriterState { outputFileHandle :: !WritableBinaryHandle - , columnChunks :: !(VB.Vector ColumnChunkState) - , currentFileOffsetRef :: !(IORef Int64) - , scratchBuffer :: !MemoryBuffer - , rowGroupMetadataRef :: !(IORef [RowGroup]) - , rowNumberRef :: !(IORef Int) + , columnChunks :: !(VB.Vector (ColumnChunkState m)) + , currentFileOffsetRef :: !(MutVar (PrimState m) Int64) + , scratchBuffer :: !(MemoryBuffer (PrimState m)) + , rowGroupMetadataRef :: !(MutVar (PrimState m) [RowGroup]) + , rowNumberRef :: !(MutVar (PrimState m) Int) } -data ColumnChunkState = ColumnChunkState +data ColumnChunkState m = ColumnChunkState { columnName :: !T.Text , nullable :: !Bool , schema :: !SchemaElement - , encoder :: !Encoder - , buffer :: !MemoryBuffer - , uncompressedBufferSize :: !(IORef Int64) - , pageState :: !PageState + , encoder :: !(Encoder m) + , buffer :: !(MemoryBuffer (PrimState m)) + , uncompressedBufferSize :: !(MutVar (PrimState m) Int64) + , pageState :: !(PageState m) } -data PageState = PageState - { pageBuffer :: !MemoryBuffer - , definitionLevels :: !DefLevels - , currentRowCount :: !(IORef Int) +data PageState m = PageState + { pageBuffer :: !(MemoryBuffer (PrimState m)) + , definitionLevels :: !(DefLevels (PrimState m)) + , currentRowCount :: !(MutVar (PrimState m) Int) } writeParquet :: FilePath -> DataFrame -> IO () @@ -172,9 +175,9 @@ writeShard options path_ df startRow endRow = do scratchBuffer_ <- mallocBuffer (max 1 options.pageSize) atomicallyWriteFile path_ $ \path -> withWritableBinaryFile path $ \output -> do writeByteStringToFile output magic - currentFileOffsetRef_ <- newIORef 4 - rowGroupMetadataRef_ <- newIORef [] - rowNumberRef_ <- newIORef 0 + currentFileOffsetRef_ <- newMutVar 4 + rowGroupMetadataRef_ <- newMutVar [] + rowNumberRef_ <- newMutVar 0 let writerState = ParquetWriterState output @@ -190,7 +193,7 @@ writeShard options path_ df startRow endRow = do | otherwise = do let count = min subBatch (batchEnd - rowNum) VB.forM_ columnChunks_ (writeRows options scratchBuffer_ rowNum count) - modifyIORef' rowNumberRef_ (+ count) + modifyMutVar' rowNumberRef_ (+ count) writeBatch (rowNum + count) batchEnd loop rowNum | rowNum >= endRow = pure () @@ -204,7 +207,7 @@ writeShard options path_ df startRow endRow = do loop startRow flushRowGroup options writerState rowGroupMetadata <- - reverse <$> readIORef rowGroupMetadataRef_ + reverse <$> readMutVar rowGroupMetadataRef_ let schemaElements = rootSchemaElement (VB.length columnChunks_) : VB.toList (VB.map schema columnChunks_) @@ -227,12 +230,13 @@ nativeTypeKeyValues names df = ] writeRows :: + (PrimBase m, MonadIO m) => ParquetWriteOptions -> - MemoryBuffer -> + MemoryBuffer (PrimState m) -> Int -> Int -> - ColumnChunkState -> - IO () + ColumnChunkState m -> + m () writeRows options scratch firstRow count ccs = do let page = ccs.pageState buf = page.pageBuffer @@ -240,16 +244,16 @@ writeRows options scratch firstRow count ccs = do dl = page.definitionLevels end = firstRow + count - pos0 <- readIORef buf.positionRef + pos0 <- readMutVar buf.positionRef let margin = options.pageSize arr0 <- ensureCapacity buf (pos0 + max margin (count * 64)) size0 <- getSizeofMutableByteArray arr0 let go !size !pos !row - | row >= end = writeIORef buf.positionRef pos + | row >= end = writeMutVar buf.positionRef pos | pos + margin > size = do -- Rare: buffer nearly full, grow it - writeIORef buf.positionRef pos + writeMutVar buf.positionRef pos arr' <- ensureCapacity buf @@ -265,7 +269,7 @@ writeRows options scratch firstRow count ccs = do go size0 pos0 firstRow -- Batch bookkeeping: once per sub-batch instead of per value - modifyIORef' page.currentRowCount (+ count) + modifyMutVar' page.currentRowCount (+ count) flushDef dl pageRes <- bufferResidency buf defRes <- bufferResidency dl.dlBuf @@ -274,17 +278,18 @@ writeRows options scratch firstRow count ccs = do (flushPage options scratch ccs) flushPage :: + (PrimBase m, MonadIO m) => ParquetWriteOptions -> - MemoryBuffer -> - ColumnChunkState -> - IO () + MemoryBuffer (PrimState m) -> + ColumnChunkState m -> + m () flushPage options scratch columnChunkState = do let page = columnChunkState.pageState - numPageRows <- readIORef page.currentRowCount + numPageRows <- readMutVar page.currentRowCount when (numPageRows > 0) $ do - pos <- readIORef page.pageBuffer.positionRef + pos <- readMutVar page.pageBuffer.positionRef pos' <- columnChunkState.encoder.finishValues page.pageBuffer pos - writeIORef page.pageBuffer.positionRef pos' + writeMutVar page.pageBuffer.positionRef pos' body <- assemblePageBody scratch columnChunkState writeDataPage options.compressionCodec @@ -294,12 +299,13 @@ flushPage options scratch columnChunkState = do resetPosition page.pageBuffer resetPosition page.definitionLevels.dlBuf resetPosition scratch - writeIORef page.currentRowCount 0 + writeMutVar page.currentRowCount 0 assemblePageBody :: - MemoryBuffer -> - ColumnChunkState -> - IO MemoryBuffer + (PrimMonad m) => + MemoryBuffer (PrimState m) -> + ColumnChunkState m -> + m (MemoryBuffer (PrimState m)) assemblePageBody scratch columnChunkState | not columnChunkState.nullable = pure columnChunkState.pageState.pageBuffer | otherwise = do @@ -313,11 +319,12 @@ assemblePageBody scratch columnChunkState pure scratch writeDataPage :: + (PrimBase m, MonadIO m) => CompressionCodec -> Int -> - MemoryBuffer -> - ColumnChunkState -> - IO () + MemoryBuffer (PrimState m) -> + ColumnChunkState m -> + m () writeDataPage codec numPageRows body columnChunkState = do uncompressedPageSize <- bufferResidency body compressedBody <- case codec of @@ -334,16 +341,17 @@ writeDataPage codec numPageRows body columnChunkState = do case compressedBody of Nothing -> flushBufferToBuffer body columnChunkState.buffer Just bytes -> writeByteString columnChunkState.buffer bytes - modifyIORef' + modifyMutVar' columnChunkState.uncompressedBufferSize (+ fromIntegral (BS.length headerBytes + uncompressedPageSize)) flushRowGroup :: + (PrimBase m, MonadIO m) => ParquetWriteOptions -> - ParquetWriterState -> - IO () + ParquetWriterState m -> + m () flushRowGroup options writerState = do - rowNumber <- readIORef writerState.rowNumberRef + rowNumber <- readMutVar writerState.rowNumberRef when (rowNumber > 0) $ do VB.forM_ writerState.columnChunks @@ -352,19 +360,19 @@ flushRowGroup options writerState = do VB.foldM' ( \(acc, totalCompressedSize, totalUncompressedSize) columnChunkState -> do offset <- - readIORef writerState.currentFileOffsetRef + readMutVar writerState.currentFileOffsetRef compressedSize <- bufferResidency columnChunkState.buffer uncompressedSize <- - readIORef + readMutVar columnChunkState.uncompressedBufferSize flushBufferToFile writerState.outputFileHandle columnChunkState.buffer - writeIORef + writeMutVar writerState.currentFileOffsetRef (offset + fromIntegral compressedSize) - writeIORef + writeMutVar columnChunkState.uncompressedBufferSize 0 let columnChunk = @@ -384,7 +392,7 @@ flushRowGroup options writerState = do ) ([], 0 :: Int64, 0 :: Int64) writerState.columnChunks - modifyIORef' + modifyMutVar' writerState.rowGroupMetadataRef ( mkRowGroup (reverse reversedColumnChunks) @@ -393,11 +401,12 @@ flushRowGroup options writerState = do rowNumber : ) - writeIORef writerState.rowNumberRef 0 + writeMutVar writerState.rowNumberRef 0 bufferedSize :: - VB.Vector ColumnChunkState -> - IO Int + (PrimMonad m) => + VB.Vector (ColumnChunkState m) -> + m Int bufferedSize = VB.foldM' ( \total columnChunkState -> do @@ -412,10 +421,11 @@ bufferedSize = 0 initColumnChunkState :: + (PrimBase m, MonadIO m) => ParquetWriteOptions -> T.Text -> Column -> - IO ColumnChunkState + m (ColumnChunkState m) initColumnChunkState options columnName_ column = do encoder_ <- buildEncoder column let nullable_ = hasMissing column @@ -440,7 +450,7 @@ initColumnChunkState options columnName_ column = do -- is likely to hit the page limit, the others are liable to be -- much smaller than the limit. buffer_ <- mallocBuffer bufferSize - uncompressedBufferSize_ <- newIORef 0 + uncompressedBufferSize_ <- newMutVar 0 pageState_ <- initPageState bufferSize pure ColumnChunkState @@ -453,11 +463,11 @@ initColumnChunkState options columnName_ column = do , pageState = pageState_ } -initPageState :: Int -> IO PageState +initPageState :: (PrimMonad m, MonadIO m) => Int -> m (PageState m) initPageState bufferSize = do pageBuffer_ <- mallocBuffer bufferSize definitionLevels_ <- newDefLevels - currentRowCount_ <- newIORef 0 + currentRowCount_ <- newMutVar 0 pure PageState { pageBuffer = pageBuffer_ diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/DefLevels.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/DefLevels.hs index 2204c1e2..65efe19b 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/DefLevels.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/DefLevels.hs @@ -8,51 +8,53 @@ module DataFrame.IO.Parquet.Writer.DefLevels ( ) where import Control.Monad (when) +import Control.Monad.IO.Class (MonadIO) +import Control.Monad.Primitive (PrimMonad, PrimState) import Data.Bits (shiftL, shiftR, (.&.), (.|.)) -import Data.IORef (IORef, newIORef, readIORef, writeIORef) +import Data.Primitive.MutVar (MutVar, newMutVar, readMutVar, writeMutVar) import Data.Word (Word64) import DataFrame.IO.Utils.RandomAccess (MemoryBuffer, mallocBuffer, writeWord8) -data DefLevels = DefLevels - { dlBuf :: !MemoryBuffer - , dlValue :: !(IORef Int) - , dlCount :: !(IORef Int) +data DefLevels s = DefLevels + { dlBuf :: !(MemoryBuffer s) + , dlValue :: !(MutVar s Int) + , dlCount :: !(MutVar s Int) } -newDefLevels :: IO DefLevels -newDefLevels = DefLevels <$> mallocBuffer 64 <*> newIORef 0 <*> newIORef 0 +newDefLevels :: (PrimMonad m, MonadIO m) => m (DefLevels (PrimState m)) +newDefLevels = DefLevels <$> mallocBuffer 64 <*> newMutVar 0 <*> newMutVar 0 -pushDef :: DefLevels -> Int -> IO () +pushDef :: (PrimMonad m) => DefLevels (PrimState m) -> Int -> m () pushDef dl value = do - count <- readIORef dl.dlCount + count <- readMutVar dl.dlCount if count == 0 - then writeIORef dl.dlValue value >> writeIORef dl.dlCount 1 + then writeMutVar dl.dlValue value >> writeMutVar dl.dlCount 1 else do - current <- readIORef dl.dlValue + current <- readMutVar dl.dlValue if current == value - then writeIORef dl.dlCount (count + 1) + then writeMutVar dl.dlCount (count + 1) else do writeDefRun dl current count - writeIORef dl.dlValue value - writeIORef dl.dlCount 1 + writeMutVar dl.dlValue value + writeMutVar dl.dlCount 1 {-# INLINE pushDef #-} -flushDef :: DefLevels -> IO () +flushDef :: (PrimMonad m) => DefLevels (PrimState m) -> m () flushDef dl = do - count <- readIORef dl.dlCount + count <- readMutVar dl.dlCount when (count > 0) $ do - value <- readIORef dl.dlValue + value <- readMutVar dl.dlValue writeDefRun dl value count - writeIORef dl.dlCount 0 + writeMutVar dl.dlCount 0 {-# INLINE flushDef #-} -writeDefRun :: DefLevels -> Int -> Int -> IO () +writeDefRun :: (PrimMonad m) => DefLevels (PrimState m) -> Int -> Int -> m () writeDefRun dl value count = do writeLeb128 dl.dlBuf (fromIntegral (count `shiftL` 1)) writeWord8 dl.dlBuf (fromIntegral value) {-# INLINE writeDefRun #-} -writeLeb128 :: MemoryBuffer -> Word64 -> IO () +writeLeb128 :: (PrimMonad m) => MemoryBuffer (PrimState m) -> Word64 -> m () writeLeb128 buffer value | value < 0x80 = writeWord8 buffer (fromIntegral value) | otherwise = do diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs index 141b8efa..de9c1fef 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs @@ -10,14 +10,13 @@ module DataFrame.IO.Parquet.Writer.Encoder ( buildEncoder, ) where +import Control.Monad.IO.Class (MonadIO, liftIO) +import Control.Monad.Primitive (PrimBase, PrimMonad, PrimState, RealWorld) import Control.Monad.ST (stToIO) import Data.Bits (shiftL, (.|.)) -import Data.IORef (newIORef, readIORef, writeIORef) import Data.Int (Int32, Int64) -import Data.Primitive.ByteArray ( - withMutableByteArrayContents, - writeByteArray, - ) +import Data.Primitive.ByteArray (withMutableByteArrayContents, writeByteArray) +import Data.Primitive.MutVar (newMutVar, readMutVar, writeMutVar) import qualified Data.Text as T import qualified Data.Text.Array as TA import Data.Text.Internal (Text (Text)) @@ -55,16 +54,16 @@ import GHC.Float (castDoubleToWord64, castFloatToWord32) import Pinch (enum, putField) import Type.Reflection (typeRep) -data Encoder = Encoder +data Encoder m = Encoder { encType :: !ThriftType , convertedType :: !(Maybe ConvertedType) , logicalType :: !(Maybe LogicalType) , encodeValue :: - !(MemoryBuffer -> Int -> Int -> IO (Int, Bool)) - , finishValues :: !(MemoryBuffer -> Int -> IO Int) + !(MemoryBuffer (PrimState m) -> Int -> Int -> m (Int, Bool)) + , finishValues :: !(MemoryBuffer (PrimState m) -> Int -> m Int) } -buildEncoder :: Column -> IO Encoder +buildEncoder :: (PrimBase m, MonadIO m) => Column -> m (Encoder m) buildEncoder col | hasElemType @Int32 col = pure $ @@ -131,14 +130,14 @@ buildEncoder col pure (pos + 8) scalarEncoder :: - forall a. - (Columnable a) => + forall a m. + (Columnable a, Monad m) => ThriftType -> Maybe ConvertedType -> Maybe LogicalType -> - (MemoryBuffer -> Int -> a -> IO Int) -> + (MemoryBuffer (PrimState m) -> Int -> a -> m Int) -> Column -> - Encoder + Encoder m scalarEncoder tt conv logical writePrim col = Encoder tt conv logical (columnWriter @a col writePrim) (\_ pos -> pure pos) {-# INLINEABLE scalarEncoder #-} @@ -146,60 +145,60 @@ scalarEncoder tt conv logical writePrim col = ThriftType -> Maybe ConvertedType -> Maybe LogicalType -> - (MemoryBuffer -> Int -> Int32 -> IO Int) -> + (MemoryBuffer RealWorld -> Int -> Int32 -> IO Int) -> Column -> - Encoder + Encoder IO #-} {-# SPECIALIZE scalarEncoder :: ThriftType -> Maybe ConvertedType -> Maybe LogicalType -> - (MemoryBuffer -> Int -> Int64 -> IO Int) -> + (MemoryBuffer RealWorld -> Int -> Int64 -> IO Int) -> Column -> - Encoder + Encoder IO #-} {-# SPECIALIZE scalarEncoder :: ThriftType -> Maybe ConvertedType -> Maybe LogicalType -> - (MemoryBuffer -> Int -> Float -> IO Int) -> + (MemoryBuffer RealWorld -> Int -> Float -> IO Int) -> Column -> - Encoder + Encoder IO #-} {-# SPECIALIZE scalarEncoder :: ThriftType -> Maybe ConvertedType -> Maybe LogicalType -> - (MemoryBuffer -> Int -> Double -> IO Int) -> + (MemoryBuffer RealWorld -> Int -> Double -> IO Int) -> Column -> - Encoder + Encoder IO #-} {-# SPECIALIZE scalarEncoder :: ThriftType -> Maybe ConvertedType -> Maybe LogicalType -> - (MemoryBuffer -> Int -> Int -> IO Int) -> + (MemoryBuffer RealWorld -> Int -> Int -> IO Int) -> Column -> - Encoder + Encoder IO #-} {-# SPECIALIZE scalarEncoder :: ThriftType -> Maybe ConvertedType -> Maybe LogicalType -> - (MemoryBuffer -> Int -> Integer -> IO Int) -> + (MemoryBuffer RealWorld -> Int -> Integer -> IO Int) -> Column -> - Encoder + Encoder IO #-} columnWriter :: - forall a. - (Columnable a) => + forall a m. + (Columnable a, Monad m) => Column -> - (MemoryBuffer -> Int -> a -> IO Int) -> - MemoryBuffer -> + (MemoryBuffer (PrimState m) -> Int -> a -> m Int) -> + MemoryBuffer (PrimState m) -> Int -> Int -> - IO (Int, Bool) + m (Int, Bool) columnWriter col writePrim = case col of BoxedColumn bitmap (values :: VB.Vector b) -> case testEquality (typeRep @a) (typeRep @b) of @@ -222,64 +221,64 @@ columnWriter col writePrim = case col of {-# INLINEABLE columnWriter #-} {-# SPECIALIZE columnWriter :: Column -> - (MemoryBuffer -> Int -> Int32 -> IO Int) -> - MemoryBuffer -> + (MemoryBuffer RealWorld -> Int -> Int32 -> IO Int) -> + MemoryBuffer RealWorld -> Int -> Int -> IO (Int, Bool) #-} {-# SPECIALIZE columnWriter :: Column -> - (MemoryBuffer -> Int -> Int64 -> IO Int) -> - MemoryBuffer -> + (MemoryBuffer RealWorld -> Int -> Int64 -> IO Int) -> + MemoryBuffer RealWorld -> Int -> Int -> IO (Int, Bool) #-} {-# SPECIALIZE columnWriter :: Column -> - (MemoryBuffer -> Int -> Float -> IO Int) -> - MemoryBuffer -> + (MemoryBuffer RealWorld -> Int -> Float -> IO Int) -> + MemoryBuffer RealWorld -> Int -> Int -> IO (Int, Bool) #-} {-# SPECIALIZE columnWriter :: Column -> - (MemoryBuffer -> Int -> Double -> IO Int) -> - MemoryBuffer -> + (MemoryBuffer RealWorld -> Int -> Double -> IO Int) -> + MemoryBuffer RealWorld -> Int -> Int -> IO (Int, Bool) #-} {-# SPECIALIZE columnWriter :: Column -> - (MemoryBuffer -> Int -> Bool -> IO Int) -> - MemoryBuffer -> + (MemoryBuffer RealWorld -> Int -> Bool -> IO Int) -> + MemoryBuffer RealWorld -> Int -> Int -> IO (Int, Bool) #-} {-# SPECIALIZE columnWriter :: Column -> - (MemoryBuffer -> Int -> UTCTime -> IO Int) -> - MemoryBuffer -> + (MemoryBuffer RealWorld -> Int -> UTCTime -> IO Int) -> + MemoryBuffer RealWorld -> Int -> Int -> IO (Int, Bool) #-} {-# SPECIALIZE columnWriter :: Column -> - (MemoryBuffer -> Int -> Int -> IO Int) -> - MemoryBuffer -> + (MemoryBuffer RealWorld -> Int -> Int -> IO Int) -> + MemoryBuffer RealWorld -> Int -> Int -> IO (Int, Bool) #-} {-# SPECIALIZE columnWriter :: Column -> - (MemoryBuffer -> Int -> Integer -> IO Int) -> - MemoryBuffer -> + (MemoryBuffer RealWorld -> Int -> Integer -> IO Int) -> + MemoryBuffer RealWorld -> Int -> Int -> IO (Int, Bool) @@ -290,45 +289,46 @@ isPresent Nothing _ = True isPresent (Just bitmap) row = bitmapTestBit bitmap row {-# INLINE isPresent #-} -boolEncoder :: Column -> IO Encoder +boolEncoder :: (PrimMonad m) => Column -> m (Encoder m) boolEncoder col = do - bitsRef <- newIORef (0 :: Word8) - countRef <- newIORef (0 :: Int) + bitsRef <- newMutVar (0 :: Word8) + countRef <- newMutVar (0 :: Int) let addBit buffer pos value = do - bits <- readIORef bitsRef - count <- readIORef countRef + bits <- readMutVar bitsRef + count <- readMutVar countRef let bits' = if value then bits .|. ((1 :: Word8) `shiftL` count) else bits count' = count + 1 if count' == 8 then do - arr <- readIORef buffer.arrayRef + arr <- readMutVar buffer.arrayRef writeByteArray arr pos bits' - writeIORef bitsRef 0 - writeIORef countRef 0 + writeMutVar bitsRef 0 + writeMutVar countRef 0 pure (pos + 1) else do - writeIORef bitsRef bits' - writeIORef countRef count' + writeMutVar bitsRef bits' + writeMutVar countRef count' pure pos finish buffer pos = do - count <- readIORef countRef + count <- readMutVar countRef pos' <- if count > 0 then do - bits <- readIORef bitsRef - arr <- readIORef buffer.arrayRef + bits <- readMutVar bitsRef + arr <- readMutVar buffer.arrayRef writeByteArray arr pos bits pure (pos + 1) else pure pos - writeIORef bitsRef 0 - writeIORef countRef 0 + writeMutVar bitsRef 0 + writeMutVar countRef 0 pure pos' pure (Encoder (BOOLEAN enum) Nothing Nothing (columnWriter @Bool col addBit) finish) textEncoder :: + (PrimBase m, MonadIO m) => Column -> - Encoder + Encoder m textEncoder col = Encoder (BYTE_ARRAY enum) @@ -359,26 +359,28 @@ textEncoder col = pure (pos', True) | otherwise = pure (pos, False) writeTextSlice buffer pos bytes offset count = do - writeIORef buffer.positionRef pos + writeMutVar buffer.positionRef pos _ <- ensureCapacity buffer (pos + 4 + count) writeWord32At buffer pos (fromIntegral count) - arr <- readIORef buffer.arrayRef + arr <- readMutVar buffer.arrayRef withMutableByteArrayContents arr $ \ptr -> - stToIO - ( TA.copyToPointer - bytes - offset - (ptr `plusPtr` (pos + 4)) - count - ) + liftIO $ + stToIO + ( TA.copyToPointer + bytes + offset + (ptr `plusPtr` (pos + 4)) + count + ) pure (pos + 4 + count) mismatch = error ("writeParquet: incompatible text representation for " <> columnTypeString col) timestampEncoder :: + (PrimMonad m) => Column -> - Encoder + Encoder m timestampEncoder col = Encoder (INT64 enum) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Metadata.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Metadata.hs index 016b925c..6092b9bc 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Metadata.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Metadata.hs @@ -1,3 +1,4 @@ +{-# LANGUAGE FlexibleContexts #-} {-# LANGUAGE OverloadedStrings #-} module DataFrame.IO.Parquet.Writer.Metadata ( @@ -10,6 +11,8 @@ module DataFrame.IO.Parquet.Writer.Metadata ( magic, ) where +import Control.Monad.IO.Class (MonadIO) +import Control.Monad.Primitive (PrimBase) import qualified Data.ByteString as BS import Data.Int (Int64) import qualified Data.Text as T @@ -137,12 +140,13 @@ mkRowGroup chunks totalCompressed totalUncompressed rgRows = } writeFooter :: + (PrimBase m, MonadIO m) => WritableBinaryHandle -> [SchemaElement] -> Int -> [RowGroup] -> [(T.Text, T.Text)] -> - IO () + m () writeFooter output schemaElements numRows rowGroupMetadata keyValues = do let metadata = FileMetadata diff --git a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs index 25db397d..30d86bae 100644 --- a/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs +++ b/dataframe-parquet/src/DataFrame/IO/Utils/RandomAccess.hs @@ -38,13 +38,12 @@ module DataFrame.IO.Utils.RandomAccess ( import Control.Exception (bracket, bracketOnError, finally) import Control.Monad (when) import Control.Monad.IO.Class (MonadIO (..)) -import Control.Monad.Primitive (RealWorld) +import Control.Monad.Primitive (PrimBase, PrimMonad, PrimState) import Control.Monad.ST (stToIO) import Data.Bits (shiftR) import qualified Data.ByteString as BS import Data.ByteString.Internal (ByteString (PS), create) import qualified Data.ByteString.Unsafe as BU -import Data.IORef (IORef, newIORef, readIORef, writeIORef) import Data.Int (Int64) import Data.Primitive.ByteArray ( MutableByteArray, @@ -54,6 +53,7 @@ import Data.Primitive.ByteArray ( withMutableByteArrayContents, writeByteArray, ) +import Data.Primitive.MutVar (MutVar, newMutVar, readMutVar, writeMutVar) import qualified Data.Text.Array as TA import qualified Data.Vector.Storable as VS import Data.Word (Word32, Word64, Word8) @@ -200,18 +200,18 @@ withWritableBinaryFile filepath = (openWritableBinaryFile filepath) (hClose . unHandle) -data MemoryBuffer = MemoryBuffer - { arrayRef :: !(IORef (MutableByteArray RealWorld)) - , positionRef :: !(IORef Int) +data MemoryBuffer s = MemoryBuffer + { arrayRef :: !(MutVar s (MutableByteArray s)) + , positionRef :: !(MutVar s Int) } mallocBuffer :: - Int -> IO MemoryBuffer + (PrimMonad m, MonadIO m) => Int -> m (MemoryBuffer (PrimState m)) mallocBuffer capacity - | capacity < 0 = ioError $ userError "mallocBuffer: negative capacity" + | capacity < 0 = liftIO $ ioError $ userError "mallocBuffer: negative capacity" | otherwise = do array <- newPinnedByteArray capacity - MemoryBuffer <$> newIORef array <*> newIORef 0 + MemoryBuffer <$> newMutVar array <*> newMutVar 0 -- We're using pinned ByteArrays so we must -- not use the grow function brovided by Data.Primitive @@ -229,61 +229,64 @@ mallocBuffer capacity -- pre-allocate to three elements to begin with and grow it only on the -- off chance that a buffer required more than three grows). ensureCapacity :: - MemoryBuffer -> Int -> IO (MutableByteArray RealWorld) + (PrimMonad m) => + MemoryBuffer (PrimState m) -> Int -> m (MutableByteArray (PrimState m)) ensureCapacity buffer needed = do - array <- readIORef buffer.arrayRef + array <- readMutVar buffer.arrayRef maxSize <- getSizeofMutableByteArray array if needed <= maxSize then pure array else do - position <- readIORef buffer.positionRef + position <- readMutVar buffer.positionRef grown <- newPinnedByteArray (needed + (needed `div` 2)) copyMutableByteArray grown 0 array 0 position - writeIORef buffer.arrayRef grown + writeMutVar buffer.arrayRef grown pure grown {-# INLINE ensureCapacity #-} -writeWord8 :: MemoryBuffer -> Word8 -> IO () +writeWord8 :: (PrimMonad m) => MemoryBuffer (PrimState m) -> Word8 -> m () writeWord8 buffer b = do - position <- readIORef buffer.positionRef + position <- readMutVar buffer.positionRef array <- ensureCapacity buffer (position + 1) writeByteArray array position b - writeIORef buffer.positionRef (position + 1) + writeMutVar buffer.positionRef (position + 1) {-# INLINE writeWord8 #-} writeByteString :: - MemoryBuffer -> + (PrimBase m, MonadIO m) => + MemoryBuffer (PrimState m) -> ByteString -> - IO () + m () writeByteString buffer bs = do - position <- readIORef buffer.positionRef + position <- readMutVar buffer.positionRef let len = BS.length bs array <- ensureCapacity buffer (position + len) withMutableByteArrayContents array $ \dst -> - BU.unsafeUseAsCStringLen bs $ \(source, _) -> do - copyBytes - (dst `plusPtr` position) - (castPtr source) - len - writeIORef buffer.positionRef (position + len) + liftIO $ + BU.unsafeUseAsCStringLen bs $ \(source, _) -> do + copyBytes + (dst `plusPtr` position) + (castPtr source) + len + writeMutVar buffer.positionRef (position + len) {-# INLINE writeByteString #-} -writeWord32LE :: MemoryBuffer -> Word32 -> IO () +writeWord32LE :: (PrimMonad m) => MemoryBuffer (PrimState m) -> Word32 -> m () writeWord32LE buffer w = do - position <- readIORef buffer.positionRef + position <- readMutVar buffer.positionRef writeWord32At buffer position w - writeIORef buffer.positionRef (position + 4) + writeMutVar buffer.positionRef (position + 4) {-# INLINE writeWord32LE #-} -writeWord64LE :: MemoryBuffer -> Word64 -> IO () +writeWord64LE :: (PrimMonad m) => MemoryBuffer (PrimState m) -> Word64 -> m () writeWord64LE buffer w = do - position <- readIORef buffer.positionRef + position <- readMutVar buffer.positionRef writeWord64At buffer position w - writeIORef buffer.positionRef (position + 8) + writeMutVar buffer.positionRef (position + 8) {-# INLINE writeWord64LE #-} writeWord32At :: - MemoryBuffer -> Int -> Word32 -> IO () + (PrimMonad m) => MemoryBuffer (PrimState m) -> Int -> Word32 -> m () writeWord32At buffer position w = do array <- ensureCapacity buffer (position + 4) writeByteArray array position (fromIntegral w :: Word8) @@ -293,7 +296,7 @@ writeWord32At buffer position w = do {-# INLINE writeWord32At #-} writeWord64At :: - MemoryBuffer -> Int -> Word64 -> IO () + (PrimMonad m) => MemoryBuffer (PrimState m) -> Int -> Word64 -> m () writeWord64At buffer position w = do array <- ensureCapacity buffer (position + 8) writeByteArray array position (fromIntegral w :: Word8) @@ -307,15 +310,16 @@ writeWord64At buffer position w = do {-# INLINE writeWord64At #-} writeInteger64 :: - MemoryBuffer -> Integer -> IO () + (PrimMonad m, MonadIO m) => MemoryBuffer (PrimState m) -> Integer -> m () writeInteger64 buffer value = do - position <- readIORef buffer.positionRef + position <- readMutVar buffer.positionRef newPosition <- writeInteger64At buffer position value - writeIORef buffer.positionRef newPosition + writeMutVar buffer.positionRef newPosition {-# INLINE writeInteger64 #-} writeInteger64At :: - MemoryBuffer -> Int -> Integer -> IO Int + (PrimMonad m, MonadIO m) => + MemoryBuffer (PrimState m) -> Int -> Integer -> m Int writeInteger64At buffer position value | value < toInteger (minBound :: Int64) = outOfRange | value > toInteger (maxBound :: Int64) = outOfRange @@ -324,25 +328,27 @@ writeInteger64At buffer position value pure (position + 8) where outOfRange = - (ioError (userError "writeParquet: Integer value is outside the INT64 range")) + liftIO + (ioError (userError "writeParquet: Integer value is outside the INT64 range")) {-# INLINE writeInteger64At #-} -writeFloatLE :: MemoryBuffer -> Float -> IO () +writeFloatLE :: (PrimMonad m) => MemoryBuffer (PrimState m) -> Float -> m () writeFloatLE buffer = writeWord32LE buffer . castFloatToWord32 {-# INLINE writeFloatLE #-} -writeDoubleLE :: MemoryBuffer -> Double -> IO () +writeDoubleLE :: (PrimMonad m) => MemoryBuffer (PrimState m) -> Double -> m () writeDoubleLE buffer = writeWord64LE buffer . castDoubleToWord64 {-# INLINE writeDoubleLE #-} flushBufferToBuffer :: - MemoryBuffer -> MemoryBuffer -> IO () + (PrimMonad m) => + MemoryBuffer (PrimState m) -> MemoryBuffer (PrimState m) -> m () flushBufferToBuffer source destination | source.arrayRef == destination.arrayRef = pure () | otherwise = do - sourceArray <- readIORef source.arrayRef - sourcePosition <- readIORef source.positionRef - destinationPosition <- readIORef destination.positionRef + sourceArray <- readMutVar source.arrayRef + sourcePosition <- readMutVar source.positionRef + destinationPosition <- readMutVar destination.positionRef destinationArray <- ensureCapacity destination (destinationPosition + sourcePosition) copyMutableByteArray @@ -351,26 +357,28 @@ flushBufferToBuffer source destination sourceArray 0 sourcePosition - writeIORef destination.positionRef (destinationPosition + sourcePosition) - writeIORef source.positionRef 0 + writeMutVar destination.positionRef (destinationPosition + sourcePosition) + writeMutVar source.positionRef 0 {-# INLINE flushBufferToBuffer #-} bufferToByteString :: - MemoryBuffer -> - IO ByteString + (PrimBase m, MonadIO m) => + MemoryBuffer (PrimState m) -> + m ByteString bufferToByteString buffer = do - array <- readIORef buffer.arrayRef - position <- readIORef buffer.positionRef + array <- readMutVar buffer.arrayRef + position <- readMutVar buffer.positionRef withMutableByteArrayContents array $ \src -> - create position $ \dst -> - copyBytes dst (castPtr src) position + liftIO $ + create position $ \dst -> + copyBytes dst (castPtr src) position -bufferResidency :: MemoryBuffer -> IO Int -bufferResidency buffer = readIORef buffer.positionRef +bufferResidency :: (PrimMonad m) => MemoryBuffer (PrimState m) -> m Int +bufferResidency buffer = readMutVar buffer.positionRef {-# INLINE bufferResidency #-} -resetPosition :: MemoryBuffer -> IO () -resetPosition buffer = writeIORef buffer.positionRef 0 +resetPosition :: (PrimMonad m) => MemoryBuffer (PrimState m) -> m () +resetPosition buffer = writeMutVar buffer.positionRef 0 {-# INLINE resetPosition #-} -- I tested write speeds by doing (on Apple Silicon) @@ -395,11 +403,12 @@ resetPosition buffer = writeIORef buffer.positionRef 0 -- trying not to create dirty pages in the kernel page cache, we'll -- be flushing in 256 KiB chunks. flushBufferToFile :: - WritableBinaryHandle -> MemoryBuffer -> IO () + (PrimBase m, MonadIO m) => + WritableBinaryHandle -> MemoryBuffer (PrimState m) -> m () flushBufferToFile (WritableBinaryHandle h) buffer = do - array <- readIORef buffer.arrayRef - position <- readIORef buffer.positionRef - withMutableByteArrayContents array $ \ptr -> do + array <- readMutVar buffer.arrayRef + position <- readMutVar buffer.positionRef + withMutableByteArrayContents array $ \ptr -> liftIO $ do let chunkSize = 262144 go offset | offset >= position = pure () @@ -408,7 +417,7 @@ flushBufferToFile (WritableBinaryHandle h) buffer = do hPutBuf h (ptr `plusPtr` offset) n go (offset + n) go 0 - writeIORef buffer.positionRef 0 + writeMutVar buffer.positionRef 0 writeByteStringToFile :: WritableBinaryHandle -> ByteString -> IO () writeByteStringToFile (WritableBinaryHandle h) bs = @@ -423,20 +432,22 @@ writeByteStringToFile (WritableBinaryHandle h) bs = go 0 appendTextArraySlice :: - MemoryBuffer -> TA.Array -> Int -> Int -> IO () + (PrimBase m, MonadIO m) => + MemoryBuffer (PrimState m) -> TA.Array -> Int -> Int -> m () appendTextArraySlice buffer source offset count | count < 0 = - ioError $ userError "appendTextArraySlice: negative length" + liftIO $ ioError $ userError "appendTextArraySlice: negative length" | otherwise = do - position <- readIORef buffer.positionRef + position <- readMutVar buffer.positionRef array <- ensureCapacity buffer (position + count) withMutableByteArrayContents array $ \destination -> - stToIO - ( TA.copyToPointer - source - offset - (destination `plusPtr` position) - count - ) - writeIORef buffer.positionRef (position + count) + liftIO $ + stToIO + ( TA.copyToPointer + source + offset + (destination `plusPtr` position) + count + ) + writeMutVar buffer.positionRef (position + count) {-# INLINE appendTextArraySlice #-} From 13fd800cd4833e756daad118e846e8d1cf631fa4 Mon Sep 17 00:00:00 2001 From: Tom Ellis Date: Sun, 4 Oct 2026 18:49:29 +0100 Subject: [PATCH 51/51] Pragmas for optimization --- .../DataFrame/IO/Parquet/Writer/Encoder.hs | 121 +----------------- 1 file changed, 5 insertions(+), 116 deletions(-) diff --git a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs index de9c1fef..3b5b6afe 100644 --- a/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs +++ b/dataframe-parquet/src/DataFrame/IO/Parquet/Writer/Encoder.hs @@ -4,6 +4,7 @@ {-# LANGUAGE OverloadedStrings #-} {-# LANGUAGE ScopedTypeVariables #-} {-# LANGUAGE TypeApplications #-} +{-# OPTIONS_GHC -funfolding-use-threshold=1000 #-} module DataFrame.IO.Parquet.Writer.Encoder ( Encoder (..), @@ -11,7 +12,7 @@ module DataFrame.IO.Parquet.Writer.Encoder ( ) where import Control.Monad.IO.Class (MonadIO, liftIO) -import Control.Monad.Primitive (PrimBase, PrimMonad, PrimState, RealWorld) +import Control.Monad.Primitive (PrimBase, PrimMonad, PrimState) import Control.Monad.ST (stToIO) import Data.Bits (shiftL, (.|.)) import Data.Int (Int32, Int64) @@ -129,6 +130,8 @@ buildEncoder col writeWord64At buffer pos value pure (pos + 8) +{-# SPECIALIZE buildEncoder :: Column -> IO (Encoder IO) #-} + scalarEncoder :: forall a m. (Columnable a, Monad m) => @@ -140,56 +143,6 @@ scalarEncoder :: Encoder m scalarEncoder tt conv logical writePrim col = Encoder tt conv logical (columnWriter @a col writePrim) (\_ pos -> pure pos) -{-# INLINEABLE scalarEncoder #-} -{-# SPECIALIZE scalarEncoder :: - ThriftType -> - Maybe ConvertedType -> - Maybe LogicalType -> - (MemoryBuffer RealWorld -> Int -> Int32 -> IO Int) -> - Column -> - Encoder IO - #-} -{-# SPECIALIZE scalarEncoder :: - ThriftType -> - Maybe ConvertedType -> - Maybe LogicalType -> - (MemoryBuffer RealWorld -> Int -> Int64 -> IO Int) -> - Column -> - Encoder IO - #-} -{-# SPECIALIZE scalarEncoder :: - ThriftType -> - Maybe ConvertedType -> - Maybe LogicalType -> - (MemoryBuffer RealWorld -> Int -> Float -> IO Int) -> - Column -> - Encoder IO - #-} -{-# SPECIALIZE scalarEncoder :: - ThriftType -> - Maybe ConvertedType -> - Maybe LogicalType -> - (MemoryBuffer RealWorld -> Int -> Double -> IO Int) -> - Column -> - Encoder IO - #-} -{-# SPECIALIZE scalarEncoder :: - ThriftType -> - Maybe ConvertedType -> - Maybe LogicalType -> - (MemoryBuffer RealWorld -> Int -> Int -> IO Int) -> - Column -> - Encoder IO - #-} -{-# SPECIALIZE scalarEncoder :: - ThriftType -> - Maybe ConvertedType -> - Maybe LogicalType -> - (MemoryBuffer RealWorld -> Int -> Integer -> IO Int) -> - Column -> - Encoder IO - #-} - columnWriter :: forall a m. (Columnable a, Monad m) => @@ -218,71 +171,7 @@ columnWriter col writePrim = case col of mismatch = error ("writeParquet: incompatible column representation for " <> columnTypeString col) -{-# INLINEABLE columnWriter #-} -{-# SPECIALIZE columnWriter :: - Column -> - (MemoryBuffer RealWorld -> Int -> Int32 -> IO Int) -> - MemoryBuffer RealWorld -> - Int -> - Int -> - IO (Int, Bool) - #-} -{-# SPECIALIZE columnWriter :: - Column -> - (MemoryBuffer RealWorld -> Int -> Int64 -> IO Int) -> - MemoryBuffer RealWorld -> - Int -> - Int -> - IO (Int, Bool) - #-} -{-# SPECIALIZE columnWriter :: - Column -> - (MemoryBuffer RealWorld -> Int -> Float -> IO Int) -> - MemoryBuffer RealWorld -> - Int -> - Int -> - IO (Int, Bool) - #-} -{-# SPECIALIZE columnWriter :: - Column -> - (MemoryBuffer RealWorld -> Int -> Double -> IO Int) -> - MemoryBuffer RealWorld -> - Int -> - Int -> - IO (Int, Bool) - #-} -{-# SPECIALIZE columnWriter :: - Column -> - (MemoryBuffer RealWorld -> Int -> Bool -> IO Int) -> - MemoryBuffer RealWorld -> - Int -> - Int -> - IO (Int, Bool) - #-} -{-# SPECIALIZE columnWriter :: - Column -> - (MemoryBuffer RealWorld -> Int -> UTCTime -> IO Int) -> - MemoryBuffer RealWorld -> - Int -> - Int -> - IO (Int, Bool) - #-} -{-# SPECIALIZE columnWriter :: - Column -> - (MemoryBuffer RealWorld -> Int -> Int -> IO Int) -> - MemoryBuffer RealWorld -> - Int -> - Int -> - IO (Int, Bool) - #-} -{-# SPECIALIZE columnWriter :: - Column -> - (MemoryBuffer RealWorld -> Int -> Integer -> IO Int) -> - MemoryBuffer RealWorld -> - Int -> - Int -> - IO (Int, Bool) - #-} +{-# INLINE columnWriter #-} isPresent :: Maybe Bitmap -> Int -> Bool isPresent Nothing _ = True