Use PrimBase - #235
Use PrimBase#235
Conversation
|
Thanks @tomjaguarpaw - tagging @sharmrj who wrote the writer implementation for review as well since this is informative. |
There was a problem hiding this comment.
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?
| 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) |
There was a problem hiding this comment.
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)
| 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 |
There was a problem hiding this comment.
Same thing as the comment on writeByteString. We likely don't need the MonadUnliftIO constraint.
|
Sure, I will look into removing the |
OK,
Unfortunately, I don't think I can run that on my machine, given that "Memory usage for this benchmark will be north of 20 GB". |
I went ahead and ran the benchmark. Unfortunately it looks like there's a performance regression with this. But it might end up being solved by a few strategically placed specialize pragmas. I think you'll need to run a smaller benchmark with cost centers enabled to find those spots. BeforeAfter |
|
Wow, that is a big regression! Is it easy to make the benchmark run in less memory? Any particular suggestion for what I should do? |
|
The code that builds the dataframe lives in To be clear this isn't the roundtrip write a dataframe, read it back, and compare the old and new to make sure the writer is working correctly. This only writes a temporary dataframe file so it only consumes memory equal to the size of the dataframe + memory used by the writer. You run it with |
|
OK, I pushed to this PR some optimization settings that make things slightly better than before. BeforeAfterThese are the changes I made for benchmarking (only) - bench "write 10 GiB dataframe" $
+ bench "write 1 GiB dataframe" $-stressRows = 1_000_000
+stressRows = 100_000 |
|
I ran the 10gb benchmark again on the same device as before and I got So that solves the performance issue, I think. |
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
IOheavy, because it manipulates arrays andIORefsinIO. The data processing code doesn't have externally visible side effects, however, so it doesn't need to run inIO. In Haskell we have two main ways of mutating state withoutIO. Firstly there'sST, but that's no good here, because we really do need to interleave data processing with externalIO. Secondly there'sPrimMonad/PrimBase, and that does work here, becauseIOis an instance of both of those classes.This PR converts
dataframe-parquetto usePrimMonad/PrimBaseand thereby indicate statically which parts of the codebase do not perform externally-visible effects. In particular, these functions are moved out outIO:DefLevels.hspushDefRandomAccess.hswriteFloatLEwriteDoubleLEEncoder.hsboolEncodertimestampEncoderWriter.hsassemblePageBodybufferedSizeFor example,
asesmblePageBodychanges like this:It no longer returns
IOand 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