Skip to content

Use PrimBase - #235

Open
tomjaguarpaw wants to merge 46 commits into
DataHaskell:mainfrom
tomjaguarpaw:PrimBase
Open

tomjaguarpaw wants to merge 46 commits into
DataHaskell:mainfrom
tomjaguarpaw:PrimBase

Conversation

@tomjaguarpaw

Copy link
Copy Markdown

Hello! I noticed the recent blog post https://www.datahaskell.org/blog/2026/09/18/writing-parquet-files-using-haskell.html. Reading through the code I saw it is very IO heavy, because it manipulates arrays and IORefs in IO. The data processing code doesn't have externally visible side effects, however, so it doesn't need to run in IO. In Haskell we have two main ways of mutating state without IO. Firstly there's ST, but that's no good here, because we really do need to interleave data processing with external IO. Secondly there's PrimMonad/PrimBase, and that does work here, because IO is an instance of both of those classes.

This PR converts dataframe-parquet to use PrimMonad/PrimBase and thereby indicate statically which parts of the codebase do not perform externally-visible effects. In particular, these functions are moved out out IO:

  • DefLevels.hs
    • pushDef
  • RandomAccess.hs
    • writeFloatLE
    • writeDoubleLE
  • Encoder.hs
    • boolEncoder
    • timestampEncoder
  • Writer.hs
    • assemblePageBody
    • bufferedSize

For example, asesmblePageBody changes like this:

 assemblePageBody ::
-    MemoryBuffer ->
-    ColumnChunkState ->
-    IO MemoryBuffer
+    (PrimMonad m) =>
+    MemoryBuffer (PrimState m) ->
+    ColumnChunkState m ->
+    m (MemoryBuffer (PrimState m))

It no longer returns IO and it is clear that it does nothing externally visible.

N.B. This PR consists of many of whitespace commits (intended to make the payload commit diff smaller) followed by one payload commit (that performs the generalization).


Disclaimer: this PR was created with significant use of an AI agent

@mchav

mchav commented Sep 27, 2026

Copy link
Copy Markdown
Member

Thanks @tomjaguarpaw - tagging @sharmrj who wrote the writer implementation for review as well since this is informative.

@sharmrj sharmrj left a comment •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for your work!

I was looking for the correct abstraction to do all the memory/storage manipulation along with the IORef bookkeeping, but things got convoluted enough that I just did it in the IO monad. This is a great improvement.

I think the MonadUnliftIO constraint is likely unneeded. See the comments.

Could you also run the benchmark in the dataframe-parquet folder (the one that makes 10gb of parquet data and writes it; it should end up being ~2gb and change on disk) with and without your changes to make sure there isn't a performance regression?

Comment on lines 255 to +272
writeByteString ::
MemoryBuffer ->
(PrimBase m, MonadIO m, MonadUnliftIO m) =>
MemoryBuffer (PrimState m) ->
ByteString ->
IO ()
m ()
writeByteString buffer bs =
BU.unsafeUseAsCStringLen bs $ \(source, len) -> do
position <- readIORef buffer.positionRef
array <- ensureCapacity buffer (position + len)
withMutableByteArrayContents array $ \dst ->
copyBytes
(dst `plusPtr` position)
(castPtr source)
len
writeIORef buffer.positionRef (position + len)
withRunInIO $ \run ->
BU.unsafeUseAsCStringLen bs $ \(source, len) -> do
run $ do
position <- readMutVar buffer.positionRef
array <- ensureCapacity buffer (position + len)
withMutableByteArrayContents array $ \dst ->
liftIO $
copyBytes
(dst `plusPtr` position)
(castPtr source)
len
writeMutVar buffer.positionRef (position + len)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I'm not sure we need the extra MonadUnliftIO constraint here. We could instead make do with just liftIO and doing withMutableByteArrayContents before we do BS.unsafeUseAsCStringLen, and get the length by doing BS.length which ought to be cheap anyway. Roughly (I have not tried compiling this):

withMutableByteArrayContents array $ \destination ->
    position <- readMutVar buffer.positionRef
    let len = BS.length bs
    array <- ensureCapacity buffer (position + len)
    liftIO $ BS.unsafeUseAsCStringLen bs $ \(source, _) ->
        (copyBytes
            (destination`plusPtr` position)
            (castPtr source)
             len) >> writeMutVar buffer.positionRef (position + len)

Comment on lines 365 to +377
bufferToByteString ::
MemoryBuffer ->
IO ByteString
(PrimBase m, MonadIO m, MonadUnliftIO m) =>
MemoryBuffer (PrimState m) ->
m ByteString
bufferToByteString buffer = do
array <- readIORef buffer.arrayRef
position <- readIORef buffer.positionRef
create position $ \dst ->
withMutableByteArrayContents array $ \src ->
copyBytes dst (castPtr src) position

bufferResidency :: MemoryBuffer -> IO Int
bufferResidency buffer = readIORef buffer.positionRef
array <- readMutVar buffer.arrayRef
position <- readMutVar buffer.positionRef
bytes <- withRunInIO $ \run ->
create position $ \dst ->
run $
withMutableByteArrayContents array $ \src ->
liftIO $ copyBytes dst (castPtr src) position
pure bytes

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Same thing as the comment on writeByteString. We likely don't need the MonadUnliftIO constraint.

@tomjaguarpaw

Copy link
Copy Markdown
Author

Sure, I will look into removing the MonadUnliftIO as you suggest and get back to you, hopefully this weekend.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants