dataframe-arrow-bridge (empty) → 1.0.0.0
raw patch · 4 files changed
+1137/−0 lines, 4 filesdep +aesondep +basedep +bytestring
Dependencies added: aeson, base, bytestring, containers, dataframe-core, dataframe-csv, dataframe-expr-serializer, dataframe-json, dataframe-lazy, dataframe-operations, dataframe-parquet, dataframe-parsing, text, vector
Files
- LICENSE +20/−0
- dataframe-arrow-bridge.cabal +61/−0
- src/DataFrame/IO/Arrow.hs +575/−0
- src/DataFrame/IR.hs +481/−0
+ LICENSE view
@@ -0,0 +1,20 @@+Copyright (c) 2026 Michael Chavinda++Permission is hereby granted, free of charge, to any person obtaining+a copy of this software and associated documentation files (the+"Software"), to deal in the Software without restriction, including+without limitation the rights to use, copy, modify, merge, publish,+distribute, sublicense, and/or sell copies of the Software, and to+permit persons to whom the Software is furnished to do so, subject to+the following conditions:++The above copyright notice and this permission notice shall be included+in all copies or substantial portions of the Software.++THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND,+EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF+MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT.+IN NO EVENT SHALL THE AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY+CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN AN ACTION OF CONTRACT,+TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION WITH THE+SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
+ dataframe-arrow-bridge.cabal view
@@ -0,0 +1,61 @@+cabal-version: 3.4+name: dataframe-arrow-bridge+version: 1.0.0.0+synopsis: Arrow C Data Interface bridge and plan IR for the dataframe ecosystem.++description:+ Zero-copy conversion between @DataFrame@ and the Arrow C Data+ Interface, plus the plan IR that the Python bindings and the+ @dataframe-arrow@ foreign library execute. Re-exports+ @DataFrame.IR.ExprJson@ from @dataframe-expr-serializer@ so+ consumers keep a single import.+ .+ Previously shipped as the @arrow-bridge@ public sublibrary of the+ @dataframe@ meta-package; it is a standalone package from 1.0.0.0 so+ that dependents resolve on Hackage.++bug-reports: https://github.com/mchav/dataframe/issues+license: MIT+license-file: LICENSE+author: Michael Chavinda+maintainer: mschavinda@gmail.com+copyright: (c) 2024-2026 Michael Chavinda+category: Data+tested-with: GHC ==9.4.8 || ==9.6.7 || ==9.8.4 || ==9.10.3 || ==9.12.2++source-repository head+ type: git+ location: https://github.com/mchav/dataframe++common warnings+ ghc-options:+ -Wincomplete-patterns+ -Wincomplete-uni-patterns+ -Wunused-imports+ -Wunused-local-binds+ -Wunused-packages++library+ import: warnings+ hs-source-dirs: src+ exposed-modules: DataFrame.IO.Arrow+ DataFrame.IR+ -- The expr/pipeline JSON codec lives in its own lightweight package;+ -- re-export it so the Python FFI and Haskell consumers keep importing+ -- @DataFrame.IR.ExprJson@ unchanged.+ reexported-modules: DataFrame.IR.ExprJson+ build-depends: base >= 4 && < 5,+ aeson >= 0.11 && < 3,+ bytestring >= 0.11 && < 0.14,+ containers >= 0.6.7 && < 0.10,+ dataframe-core >= 2.5 && < 2.6,+ dataframe-csv >= 2.3 && < 2.4,+ dataframe-expr-serializer >= 1.2.1 && < 1.3,+ dataframe-json >= 1.2.0.1 && < 1.3,+ dataframe-lazy >= 2.4.1 && < 2.5,+ dataframe-operations >= 2.5 && < 2.6,+ dataframe-parquet >= 1.5 && < 1.6,+ dataframe-parsing >= 2.2 && < 2.3,+ text >= 2.1 && < 3,+ vector >= 0.13 && < 0.15+ default-language: Haskell2010
+ src/DataFrame/IO/Arrow.hs view
@@ -0,0 +1,575 @@+{-# LANGUAGE ExplicitNamespaces #-}+{-# LANGUAGE ForeignFunctionInterface #-}+{-# LANGUAGE GADTs #-}+{-# LANGUAGE ScopedTypeVariables #-}+{-# LANGUAGE TypeApplications #-}++{- | Convert a 'DataFrame' to Arrow C Data Interface structs for zero-copy+ transfer to Python (or any other Arrow consumer).+-}+module DataFrame.IO.Arrow (+ dataframeToArrow,+ columnToArrow,+ arrowToDataframe,+ releaseSchemaImpl,+ releaseArrayImpl,+) where++import qualified Data.ByteString as BS+import qualified Data.Map as M+import qualified Data.Text as T+import qualified Data.Text.Encoding as TE+import qualified Data.Vector as V+import qualified Data.Vector.Unboxed as VU+import qualified DataFrame.Internal.Column as DI+import qualified DataFrame.Internal.Column.Bitmap as DI++import Control.Monad (foldM_, forM, join, when, zipWithM_)+import Data.Type.Equality (TestEquality (testEquality), type (:~:) (Refl))+import Foreign (+ Bits (popCount),+ FunPtr,+ Int32,+ Int64,+ Ptr,+ StablePtr,+ Storable (peek, peekElemOff, poke, pokeElemOff),+ Word8,+ castPtr,+ castPtrToStablePtr,+ castStablePtrToPtr,+ copyBytes,+ deRefStablePtr,+ free,+ freeStablePtr,+ mallocArray,+ mallocBytes,+ newStablePtr,+ nullFunPtr,+ nullPtr,+ plusPtr,+ )+import Foreign.C.String (CString, newCString, peekCString)+import Type.Reflection (typeRep)++import DataFrame.Internal.Column (Column (..))+import DataFrame.Internal.DataFrame (DataFrame (..), fromNamedColumns)++-- ---------------------------------------------------------------------------+-- Opaque phantom types for the Arrow structs+-- ---------------------------------------------------------------------------++data ArrowSchema+data ArrowArray++arrowSchemaSize :: Int+arrowSchemaSize = 72 -- 9 × 8 bytes++arrowArraySize :: Int+arrowArraySize = 80 -- 10 × 8 bytes++-- ArrowSchema field byte offsets+_schemaFormat+ , _schemaName+ , _schemaMetadata+ , _schemaFlags+ , _schemaNChildren+ , _schemaChildren+ , _schemaDictionary+ , _schemaRelease+ , _schemaPrivateData ::+ Int+_schemaFormat = 0+_schemaName = 8+_schemaMetadata = 16+_schemaFlags = 24+_schemaNChildren = 32+_schemaChildren = 40+_schemaDictionary = 48+_schemaRelease = 56+_schemaPrivateData = 64++-- ArrowArray field byte offsets+_arrayLength+ , _arrayNullCount+ , _arrayOffset+ , _arrayNBuffers+ , _arrayNChildren+ , _arrayBuffers+ , _arrayChildren+ , _arrayDictionary+ , _arrayRelease+ , _arrayPrivateData ::+ Int+_arrayLength = 0+_arrayNullCount = 8+_arrayOffset = 16+_arrayNBuffers = 24+_arrayNChildren = 32+_arrayBuffers = 40+_arrayChildren = 48+_arrayDictionary = 56+_arrayRelease = 64+_arrayPrivateData = 72++-- ---------------------------------------------------------------------------+-- Helpers+-- ---------------------------------------------------------------------------++-- Write a Storable value at a byte offset from a base pointer.+at :: (Storable a) => Ptr b -> Int -> a -> IO ()+at p off = poke (castPtr (p `plusPtr` off))++-- Read a Storable value at a byte offset from a base pointer.+readAt :: (Storable a) => Ptr b -> Int -> IO a+readAt p off = peek (castPtr (p `plusPtr` off))++-- ---------------------------------------------------------------------------+-- Release callbacks (self-import trick for compile-time-constant FunPtr)+-- ---------------------------------------------------------------------------++foreign export ccall "df_release_schema"+ releaseSchemaImpl :: Ptr ArrowSchema -> IO ()++foreign import ccall "&df_release_schema"+ pReleaseSchema :: FunPtr (Ptr ArrowSchema -> IO ())++foreign export ccall "df_release_array"+ releaseArrayImpl :: Ptr ArrowArray -> IO ()++foreign import ccall "&df_release_array"+ pReleaseArray :: FunPtr (Ptr ArrowArray -> IO ())++-- Dynamic wrappers to call producer's release callbacks after copying.+foreign import ccall "dynamic"+ callRelSchema :: FunPtr (Ptr ArrowSchema -> IO ()) -> Ptr ArrowSchema -> IO ()++foreign import ccall "dynamic"+ callRelArray :: FunPtr (Ptr ArrowArray -> IO ()) -> Ptr ArrowArray -> IO ()++releaseSchemaImpl :: Ptr ArrowSchema -> IO ()+releaseSchemaImpl p = do+ rawPriv <- peek (castPtr (p `plusPtr` _schemaPrivateData) :: Ptr (Ptr ()))+ let sp = castPtrToStablePtr rawPriv :: StablePtr (IO ())+ join (deRefStablePtr sp)+ freeStablePtr sp+ -- Arrow spec: release callback must set release to NULL to signal completion.+ -- p here is Arrow C++'s internal copy of the struct (not our mallocBytes+ -- allocation); our original allocation is freed inside the cleanup closure.+ p `at` _schemaRelease $ (nullFunPtr :: FunPtr (Ptr ArrowSchema -> IO ()))++releaseArrayImpl :: Ptr ArrowArray -> IO ()+releaseArrayImpl p = do+ rawPriv <- peek (castPtr (p `plusPtr` _arrayPrivateData) :: Ptr (Ptr ()))+ let sp = castPtrToStablePtr rawPriv :: StablePtr (IO ())+ join (deRefStablePtr sp)+ freeStablePtr sp+ -- Same reasoning as releaseSchemaImpl.+ p `at` _arrayRelease $ (nullFunPtr :: FunPtr (Ptr ArrowArray -> IO ()))++makeLeafSchema :: String -> T.Text -> IO (Ptr ArrowSchema)+makeLeafSchema fmt colName = do+ p <- mallocBytes arrowSchemaSize+ fmtStr <- newCString fmt+ nameStr <- newCString (T.unpack colName)+ p `at` _schemaFormat $ fmtStr+ p `at` _schemaName $ nameStr+ p `at` _schemaMetadata $ (nullPtr :: Ptr ())+ p `at` _schemaFlags $ (0 :: Int64)+ p `at` _schemaNChildren $ (0 :: Int64)+ p `at` _schemaChildren $ (nullPtr :: Ptr ())+ p `at` _schemaDictionary $ (nullPtr :: Ptr ())+ p `at` _schemaRelease $ pReleaseSchema+ -- Capture p so our original mallocBytes allocation is freed when release runs.+ cleanup <- newStablePtr (free fmtStr >> free nameStr >> free p)+ p `at` _schemaPrivateData $ castStablePtrToPtr cleanup+ return p++makeLeafArray :: Int -> Int64 -> [Ptr ()] -> IO () -> IO (Ptr ArrowArray)+makeLeafArray nRows nullCnt bufPtrs extraCleanup = do+ p <- mallocBytes arrowArraySize+ let nb = length bufPtrs+ bufArr <- mallocArray nb :: IO (Ptr (Ptr ()))+ zipWithM_ (pokeElemOff bufArr) [0 ..] bufPtrs+ p `at` _arrayLength $ (fromIntegral nRows :: Int64)+ p `at` _arrayNullCount $ nullCnt+ p `at` _arrayOffset $ (0 :: Int64)+ p `at` _arrayNBuffers $ (fromIntegral nb :: Int64)+ p `at` _arrayNChildren $ (0 :: Int64)+ p `at` _arrayBuffers $ bufArr+ p `at` _arrayChildren $ (nullPtr :: Ptr ())+ p `at` _arrayDictionary $ (nullPtr :: Ptr ())+ p `at` _arrayRelease $ pReleaseArray+ -- Capture p so our original mallocBytes allocation is freed when release runs.+ cleanup <- newStablePtr (free bufArr >> extraCleanup >> free p)+ p `at` _arrayPrivateData $ castStablePtrToPtr cleanup+ return p++{- | Allocate an Arrow-format validity bitmap from a 'DI.Bitmap'.+Returns (ptr, nullCount). Caller must 'free' the pointer.+-}+bitmapToPtr :: Int -> DI.Bitmap -> IO (Ptr Word8, Int)+bitmapToPtr n bm = do+ let numBytes = max 1 ((n + 7) `div` 8)+ validCount = VU.foldl' (\acc b -> acc + popCount b) 0 bm+ nullCount = n - validCount+ bitmapPtr <- mallocBytes numBytes :: IO (Ptr Word8)+ VU.imapM_ (pokeElemOff bitmapPtr) bm+ when (VU.length bm < numBytes) $+ mapM_+ (\i -> pokeElemOff bitmapPtr i (0 :: Word8))+ [VU.length bm .. numBytes - 1]+ return (bitmapPtr, nullCount)++-- | Read an Arrow validity bitmap into a 'DI.Bitmap'.+readArrowBitmap :: Ptr Word8 -> Int -> IO DI.Bitmap+readArrowBitmap bitmapPtr n = VU.generateM ((n + 7) `div` 8) (peekElemOff bitmapPtr)++columnToArrow :: T.Text -> Column -> IO (Ptr ArrowSchema, Ptr ArrowArray)+columnToArrow colName (UnboxedColumn _ (vec :: VU.Vector a))+ | Just Refl <- testEquality (typeRep @a) (typeRep @Int) = do+ let n = VU.length vec+ dataPtr <- mallocArray (max 1 n) :: IO (Ptr Int64)+ VU.imapM_ (\i v -> pokeElemOff dataPtr i (fromIntegral v)) vec+ sPtr <- makeLeafSchema "l" colName+ aPtr <- makeLeafArray n 0 [nullPtr, castPtr dataPtr] (free dataPtr)+ return (sPtr, aPtr)+columnToArrow colName (UnboxedColumn _ (vec :: VU.Vector a))+ | Just Refl <- testEquality (typeRep @a) (typeRep @Double) = do+ let n = VU.length vec+ dataPtr <- mallocArray (max 1 n) :: IO (Ptr Double)+ VU.imapM_ (pokeElemOff dataPtr) vec+ sPtr <- makeLeafSchema "g" colName+ aPtr <- makeLeafArray n 0 [nullPtr, castPtr dataPtr] (free dataPtr)+ return (sPtr, aPtr)+columnToArrow colName (BoxedColumn Nothing (vec :: V.Vector a))+ | Just Refl <- testEquality (typeRep @a) (typeRep @T.Text) = do+ let n = V.length vec+ bss = map TE.encodeUtf8 (V.toList vec)+ cumOff = scanl (+) 0 (map BS.length bss)+ total = last cumOff+ offPtr <- mallocArray (n + 1) :: IO (Ptr Int32)+ zipWithM_+ (\i o -> pokeElemOff offPtr i (fromIntegral o :: Int32))+ [0 ..]+ cumOff+ charsPtr <- mallocBytes (max 1 total) :: IO (Ptr Word8)+ foldM_+ ( \pos bs -> do+ BS.useAsCStringLen bs $ \(src, len) ->+ copyBytes (charsPtr `plusPtr` pos) (castPtr src) len+ return (pos + BS.length bs)+ )+ 0+ bss+ sPtr <- makeLeafSchema "u" colName+ aPtr <-+ makeLeafArray+ n+ 0+ [nullPtr, castPtr offPtr, castPtr charsPtr]+ (free offPtr >> free charsPtr)+ return (sPtr, aPtr)+columnToArrow colName (BoxedColumn Nothing (vec :: V.Vector a))+ | Just Refl <- testEquality (typeRep @a) (typeRep @Double) = do+ let n = V.length vec+ dataPtr <- mallocArray (max 1 n) :: IO (Ptr Double)+ V.imapM_ (pokeElemOff dataPtr) vec+ sPtr <- makeLeafSchema "g" colName+ aPtr <- makeLeafArray n 0 [nullPtr, castPtr dataPtr] (free dataPtr)+ return (sPtr, aPtr)+columnToArrow colName (BoxedColumn Nothing (vec :: V.Vector a))+ | Just Refl <- testEquality (typeRep @a) (typeRep @Int) = do+ let n = V.length vec+ dataPtr <- mallocArray (max 1 n) :: IO (Ptr Int64)+ V.imapM_ (\i v -> pokeElemOff dataPtr i (fromIntegral v)) vec+ sPtr <- makeLeafSchema "l" colName+ aPtr <- makeLeafArray n 0 [nullPtr, castPtr dataPtr] (free dataPtr)+ return (sPtr, aPtr)+-- Nullable Int (UnboxedColumn with bitmap)+columnToArrow colName (UnboxedColumn (Just bm) (vec :: VU.Vector a))+ | Just Refl <- testEquality (typeRep @a) (typeRep @Int) = do+ let n = VU.length vec+ (bitmapPtr, nullCount) <- bitmapToPtr n bm+ dataPtr <- mallocArray (max 1 n) :: IO (Ptr Int64)+ VU.imapM_ (\i v -> pokeElemOff dataPtr i (fromIntegral v :: Int64)) vec+ sPtr <- makeLeafSchema "l" colName+ aPtr <-+ makeLeafArray+ n+ (fromIntegral nullCount)+ [castPtr bitmapPtr, castPtr dataPtr]+ (free bitmapPtr >> free dataPtr)+ return (sPtr, aPtr)+-- Nullable Double (UnboxedColumn with bitmap)+columnToArrow colName (UnboxedColumn (Just bm) (vec :: VU.Vector a))+ | Just Refl <- testEquality (typeRep @a) (typeRep @Double) = do+ let n = VU.length vec+ (bitmapPtr, nullCount) <- bitmapToPtr n bm+ dataPtr <- mallocArray (max 1 n) :: IO (Ptr Double)+ VU.imapM_ (\i v -> pokeElemOff dataPtr i (realToFrac v :: Double)) vec+ sPtr <- makeLeafSchema "g" colName+ aPtr <-+ makeLeafArray+ n+ (fromIntegral nullCount)+ [castPtr bitmapPtr, castPtr dataPtr]+ (free bitmapPtr >> free dataPtr)+ return (sPtr, aPtr)+-- Nullable Text (BoxedColumn with bitmap)+columnToArrow colName (BoxedColumn (Just bm) (vec :: V.Vector a))+ | Just Refl <- testEquality (typeRep @a) (typeRep @T.Text) = do+ let n = V.length vec+ -- For null positions, use empty BS (null placeholder in vec is never evaluated)+ bss =+ map+ (\i -> if DI.bitmapTestBit bm i then TE.encodeUtf8 (vec V.! i) else BS.empty)+ [0 .. n - 1]+ cumOff = scanl (+) 0 (map BS.length bss)+ total = last cumOff+ (bitmapPtr, nullCount) <- bitmapToPtr n bm+ offPtr <- mallocArray (n + 1) :: IO (Ptr Int32)+ zipWithM_+ (\i o -> pokeElemOff offPtr i (fromIntegral o :: Int32))+ [0 ..]+ cumOff+ charsPtr <- mallocBytes (max 1 total) :: IO (Ptr Word8)+ foldM_+ ( \pos bs -> do+ BS.useAsCStringLen bs $ \(src, len) ->+ copyBytes (charsPtr `plusPtr` pos) (castPtr src) len+ return (pos + BS.length bs)+ )+ 0+ bss+ sPtr <- makeLeafSchema "u" colName+ aPtr <-+ makeLeafArray+ n+ (fromIntegral nullCount)+ [castPtr bitmapPtr, castPtr offPtr, castPtr charsPtr]+ (free bitmapPtr >> free offPtr >> free charsPtr)+ return (sPtr, aPtr)+columnToArrow colName _ =+ error $+ "DataFrame.IO.Arrow.columnToArrow: unsupported column type for '"+ ++ T.unpack colName+ ++ "'"++dataframeToArrow :: DataFrame -> IO (Ptr ArrowSchema, Ptr ArrowArray)+dataframeToArrow df = do+ let idxToName = M.fromList [(v, k) | (k, v) <- M.toList (columnIndices df)]+ ncols = M.size (columnIndices df)+ colsInOrder =+ [ (idxToName M.! i, columns df V.! i)+ | i <- [0 .. ncols - 1]+ ]++ childPairs <- forM colsInOrder (uncurry columnToArrow)+ let childSPtrs = map fst childPairs+ childAPtrs = map snd childPairs++ let nRows = case colsInOrder of+ [] -> 0+ (_, col) : _ -> DI.columnLength col+ topSchema <- mallocBytes arrowSchemaSize+ fmtStr <- newCString "+s"+ nameStr <- newCString ""+ childSArr <- mallocArray ncols :: IO (Ptr (Ptr ArrowSchema))+ zipWithM_ (pokeElemOff childSArr) [0 ..] childSPtrs+ topSchema `at` _schemaFormat $ fmtStr+ topSchema `at` _schemaName $ nameStr+ topSchema `at` _schemaMetadata $ (nullPtr :: Ptr ())+ topSchema `at` _schemaFlags $ (0 :: Int64)+ topSchema `at` _schemaNChildren $ (fromIntegral ncols :: Int64)+ topSchema `at` _schemaChildren $ childSArr+ topSchema `at` _schemaDictionary $ (nullPtr :: Ptr ())+ topSchema `at` _schemaRelease $ pReleaseSchema+ -- Do NOT loop over children here: Arrow C++ zeroes children[i]->release+ -- during import, so reading it would yield a null function pointer.+ -- Children are released independently by Arrow C++; their own cleanup+ -- closures free their buffers and struct memory.+ cleanupS <- newStablePtr $ do+ free childSArr+ free fmtStr+ free nameStr+ free topSchema -- free our original mallocBytes allocation+ topSchema `at` _schemaPrivateData $ castStablePtrToPtr cleanupS++ -- ── Top-level struct array ──────────────────────────────────────────────+ topArray <- mallocBytes arrowArraySize+ childAArr <- mallocArray ncols :: IO (Ptr (Ptr ArrowArray))+ zipWithM_ (pokeElemOff childAArr) [0 ..] childAPtrs+ topBufArr <- mallocArray 1 :: IO (Ptr (Ptr ()))+ pokeElemOff topBufArr 0 nullPtr+ topArray `at` _arrayLength $ (fromIntegral nRows :: Int64)+ topArray `at` _arrayNullCount $ (0 :: Int64)+ topArray `at` _arrayOffset $ (0 :: Int64)+ topArray `at` _arrayNBuffers $ (1 :: Int64)+ topArray `at` _arrayNChildren $ (fromIntegral ncols :: Int64)+ topArray `at` _arrayBuffers $ topBufArr+ topArray `at` _arrayChildren $ childAArr+ topArray `at` _arrayDictionary $ (nullPtr :: Ptr ())+ topArray `at` _arrayRelease $ pReleaseArray+ -- Same reasoning as cleanupS: Arrow C++ manages children independently.+ cleanupA <- newStablePtr $ do+ free childAArr+ free topBufArr+ free topArray -- free our original mallocBytes allocation+ topArray `at` _arrayPrivateData $ castStablePtrToPtr cleanupA++ return (topSchema, topArray)++{- | Import an Arrow RecordBatch from raw C Data Interface pointers.+ Copies all data into GC-managed Haskell vectors, then calls the+ producer's release callbacks.+-}+arrowToDataframe :: Ptr () -> Ptr () -> IO DataFrame+arrowToDataframe rawSchema rawArray = do+ let schemaPtr = castPtr rawSchema :: Ptr ArrowSchema+ arrayPtr = castPtr rawArray :: Ptr ArrowArray+ nCols <- readAt schemaPtr _schemaNChildren :: IO Int64+ childSArr <- readAt schemaPtr _schemaChildren :: IO (Ptr (Ptr ArrowSchema))+ childAArr <- readAt arrayPtr _arrayChildren :: IO (Ptr (Ptr ArrowArray))+ cols <- forM [0 .. fromIntegral nCols - 1] $ \i -> do+ cs <- peekElemOff childSArr i+ ca <- peekElemOff childAArr i+ readArrowColumn cs ca+ -- Call producer's release callbacks after all data has been copied.+ relA <- readAt arrayPtr _arrayRelease :: IO (FunPtr (Ptr ArrowArray -> IO ()))+ when (relA /= nullFunPtr) $ callRelArray relA arrayPtr+ relS <-+ readAt schemaPtr _schemaRelease :: IO (FunPtr (Ptr ArrowSchema -> IO ()))+ when (relS /= nullFunPtr) $ callRelSchema relS schemaPtr+ return $ fromNamedColumns cols++readArrowColumn :: Ptr ArrowSchema -> Ptr ArrowArray -> IO (T.Text, Column)+readArrowColumn schemaPtr arrayPtr = do+ fmtStr <- (readAt schemaPtr _schemaFormat :: IO CString) >>= peekCString+ nameStr <- (readAt schemaPtr _schemaName :: IO CString) >>= peekCString+ let name = T.pack nameStr+ len <- readAt arrayPtr _arrayLength :: IO Int64+ nullCnt <- readAt arrayPtr _arrayNullCount :: IO Int64+ bufArr <- readAt arrayPtr _arrayBuffers :: IO (Ptr (Ptr ()))+ let n = fromIntegral len+ col <- case fmtStr of+ "l" -> readInt64Col n nullCnt bufArr+ "i" -> readInt32Col n nullCnt bufArr+ "g" -> readFloat64Col n nullCnt bufArr+ "f" -> readFloat32Col n nullCnt bufArr+ "U" -> readLargeUtf8Col n nullCnt bufArr+ "u" -> readUtf8Col n nullCnt bufArr+ _ ->+ error $+ "DataFrame.IO.Arrow.readArrowColumn: unsupported format '"+ ++ fmtStr+ ++ "' for column '"+ ++ nameStr+ ++ "'"+ return (name, col)++readInt64Col :: Int -> Int64 -> Ptr (Ptr ()) -> IO Column+readInt64Col n nullCnt bufArr = do+ bitmapVoid <- peekElemOff bufArr 0+ dataVoid <- peekElemOff bufArr 1+ let dataPtr = castPtr dataVoid :: Ptr Int64+ if nullCnt > 0+ then do+ let bitmapPtr = castPtr bitmapVoid :: Ptr Word8+ bm <- readArrowBitmap bitmapPtr n+ vec <- VU.generateM n $ \i -> fmap fromIntegral (peekElemOff dataPtr i :: IO Int64)+ return $ UnboxedColumn (Just bm) (vec :: VU.Vector Int)+ else do+ vec <- VU.generateM n $ \i -> fmap fromIntegral (peekElemOff dataPtr i :: IO Int64)+ return $ UnboxedColumn Nothing (vec :: VU.Vector Int)++readInt32Col :: Int -> Int64 -> Ptr (Ptr ()) -> IO Column+readInt32Col n nullCnt bufArr = do+ bitmapVoid <- peekElemOff bufArr 0+ dataVoid <- peekElemOff bufArr 1+ let dataPtr = castPtr dataVoid :: Ptr Int32+ if nullCnt > 0+ then do+ let bitmapPtr = castPtr bitmapVoid :: Ptr Word8+ bm <- readArrowBitmap bitmapPtr n+ vec <- VU.generateM n $ \i -> fmap fromIntegral (peekElemOff dataPtr i :: IO Int32)+ return $ UnboxedColumn (Just bm) (vec :: VU.Vector Int)+ else do+ vec <- VU.generateM n $ \i -> fmap fromIntegral (peekElemOff dataPtr i :: IO Int32)+ return $ UnboxedColumn Nothing (vec :: VU.Vector Int)++readFloat64Col :: Int -> Int64 -> Ptr (Ptr ()) -> IO Column+readFloat64Col n nullCnt bufArr = do+ bitmapVoid <- peekElemOff bufArr 0+ dataVoid <- peekElemOff bufArr 1+ let dataPtr = castPtr dataVoid :: Ptr Double+ if nullCnt > 0+ then do+ let bitmapPtr = castPtr bitmapVoid :: Ptr Word8+ bm <- readArrowBitmap bitmapPtr n+ vec <- VU.generateM n (peekElemOff dataPtr)+ return $ UnboxedColumn (Just bm) (vec :: VU.Vector Double)+ else do+ vec <- VU.generateM n (peekElemOff dataPtr)+ return $ UnboxedColumn Nothing (vec :: VU.Vector Double)++readFloat32Col :: Int -> Int64 -> Ptr (Ptr ()) -> IO Column+readFloat32Col n nullCnt bufArr = do+ bitmapVoid <- peekElemOff bufArr 0+ dataVoid <- peekElemOff bufArr 1+ let dataPtr = castPtr dataVoid :: Ptr Float+ if nullCnt > 0+ then do+ let bitmapPtr = castPtr bitmapVoid :: Ptr Word8+ bm <- readArrowBitmap bitmapPtr n+ vec <- VU.generateM n $ \i -> fmap (realToFrac :: Float -> Double) (peekElemOff dataPtr i)+ return $ UnboxedColumn (Just bm) (vec :: VU.Vector Double)+ else do+ vec <- VU.generateM n $ \i -> fmap (realToFrac :: Float -> Double) (peekElemOff dataPtr i)+ return $ UnboxedColumn Nothing (vec :: VU.Vector Double)++-- | Read a large_string (format "U") column with int64 offsets.+readLargeUtf8Col :: Int -> Int64 -> Ptr (Ptr ()) -> IO Column+readLargeUtf8Col n nullCnt bufArr = do+ bitmapVoid <- peekElemOff bufArr 0+ offsetVoid <- peekElemOff bufArr 1+ charVoid <- peekElemOff bufArr 2+ let offsetPtr = castPtr offsetVoid :: Ptr Int64+ charPtr = castPtr charVoid :: Ptr Word8+ let readText i = do+ start <- fromIntegral <$> peekElemOff offsetPtr i+ end <- fromIntegral <$> peekElemOff offsetPtr (i + 1)+ TE.decodeUtf8+ <$> BS.packCStringLen (castPtr (charPtr `plusPtr` start), end - start)+ if nullCnt > 0+ then do+ let bitmapPtr = castPtr bitmapVoid :: Ptr Word8+ bm <- readArrowBitmap bitmapPtr n+ vec <- V.generateM n readText+ return $ BoxedColumn (Just bm) vec+ else do+ vec <- V.generateM n readText+ return $ BoxedColumn Nothing vec++-- | Read a utf8 (format "u") column with int32 offsets.+readUtf8Col :: Int -> Int64 -> Ptr (Ptr ()) -> IO Column+readUtf8Col n nullCnt bufArr = do+ bitmapVoid <- peekElemOff bufArr 0+ offsetVoid <- peekElemOff bufArr 1+ charVoid <- peekElemOff bufArr 2+ let offsetPtr = castPtr offsetVoid :: Ptr Int32+ charPtr = castPtr charVoid :: Ptr Word8+ readText i = do+ start <- fromIntegral <$> peekElemOff offsetPtr i+ end <- fromIntegral <$> peekElemOff offsetPtr (i + 1)+ TE.decodeUtf8+ <$> BS.packCStringLen (castPtr (charPtr `plusPtr` start), end - start)+ if nullCnt > 0+ then do+ let bitmapPtr = castPtr bitmapVoid :: Ptr Word8+ bm <- readArrowBitmap bitmapPtr n+ vec <- V.generateM n readText+ return $ BoxedColumn (Just bm) vec+ else do+ vec <- V.generateM n readText+ return $ BoxedColumn Nothing vec
+ src/DataFrame/IR.hs view
@@ -0,0 +1,481 @@+{-# LANGUAGE AllowAmbiguousTypes #-}+{-# LANGUAGE ExplicitNamespaces #-}+{-# LANGUAGE FlexibleContexts #-}+{-# LANGUAGE GADTs #-}+{-# LANGUAGE OverloadedStrings #-}+{-# LANGUAGE RankNTypes #-}+{-# LANGUAGE ScopedTypeVariables #-}+{-# LANGUAGE TypeApplications #-}++{- | Intermediate Representation for DataFrame query plans.+ JSON-decodable plan tree + interpreter.+-}+module DataFrame.IR (+ PlanNode (..),+ AggSpec (..),+ executePlan,+) where++import Data.Aeson (FromJSON (..), withObject, (.:))+import qualified Data.Aeson as Aeson+import Data.Aeson.Types (Parser)+import qualified Data.ByteString as BS+import Data.Int (Int16, Int32, Int64, Int8)+import qualified Data.Text as T+import Data.Type.Equality (+ TestEquality (testEquality),+ type (:~:) (Refl),+ type (:~~:) (HRefl),+ )+import qualified Data.Vector as V+import qualified Data.Vector.Unboxed as VU+import Data.Word (Word16, Word32, Word64, Word8)+import Foreign (wordPtrToPtr)+import Type.Reflection (SomeTypeRep (..), eqTypeRep, typeRep)++import DataFrame.Expression.Operators ((.=))+import DataFrame.Functions (count, mean, meanMaybe, sumMaybe)+import qualified DataFrame.Functions as Functions+import DataFrame.IO.Arrow (arrowToDataframe)+import DataFrame.IO.CSV (+ CsvReader,+ defaultReadOptions,+ readSeparated,+ readTsv,+ writeCsv,+ )+import DataFrame.IO.JSON (readJSON)+import qualified DataFrame.IO.Parquet as Parquet+import DataFrame.IR.ExprJson (SomeExpr (..), decodeExprAny, decodeExprAt)+import DataFrame.Internal.Column (+ Column (..),+ Columnable,+ mergedHead,+ )+import DataFrame.Internal.DataFrame (DataFrame, unsafeGetColumn)+import DataFrame.Internal.Expression (Expr (..), NamedExpr)+import qualified DataFrame.Lazy as Lazy+import DataFrame.Operations.Aggregation (aggregate, distinct, groupBy)+import DataFrame.Operations.Core (insertVector, renameMany)+import DataFrame.Operations.Join (JoinType (..), join)+import DataFrame.Operations.Permutation (SortOrder (..), sortBy)+import qualified DataFrame.Operations.Statistics as Stats+import DataFrame.Operations.Subset (exclude, filterWhere, range, select)+import qualified DataFrame.Operations.Subset as Subset+import DataFrame.Operations.Transformations (derive)+import DataFrame.Schema (Schema, makeSchema, schemaType)++-- ---------------------------------------------------------------------------+-- IR types+-- ---------------------------------------------------------------------------++data AggSpec = AggSpec+ { aggName :: T.Text+ , aggFn :: T.Text+ , aggCol :: T.Text+ }+ deriving (Show)++data PlanNode+ = ReadCsv FilePath+ | ReadTsv FilePath+ | -- | schema_addr array_addr+ FromArrow Word64 Word64+ | Select [T.Text] PlanNode+ | GroupBy [T.Text] [AggSpec] PlanNode+ | Sort [T.Text] Bool PlanNode+ | Limit Int PlanNode+ | -- | predicate JSON, child plan+ Filter Aeson.Value PlanNode+ | -- | column name, expr JSON, child plan+ Derive T.Text Aeson.Value PlanNode+ | Exclude [T.Text] PlanNode+ | Rename [(T.Text, T.Text)] PlanNode+ | Distinct PlanNode+ | TakeLast Int PlanNode+ | Drop Int PlanNode+ | DropLast Int PlanNode+ | Range Int Int PlanNode+ | -- | joinType ("inner"|"left"|"right"|"outer"), shared key columns, left, right+ Join T.Text [T.Text] PlanNode PlanNode+ | Describe PlanNode+ | -- | first column, second column, child plan+ Correlation T.Text T.Text PlanNode+ | Frequencies T.Text PlanNode+ | ReadParquet FilePath+ | ReadJson FilePath+ | -- | path, separator (single character), child plan; runs as a terminal op+ WriteCsv FilePath PlanNode+ | {- | path, schema (column name → type-tag map). Reads via the lazy+ engine with predicate / projection pushdown; subsequent ops+ currently still run eagerly on the materialized result.+ -}+ ScanCsv FilePath [(T.Text, T.Text)]+ | ScanParquet FilePath [(T.Text, T.Text)]+ deriving (Show)++-- ---------------------------------------------------------------------------+-- JSON decoding+-- ---------------------------------------------------------------------------++instance FromJSON AggSpec where+ parseJSON = withObject "AggSpec" $ \o ->+ AggSpec+ <$> o .: "name"+ <*> o .: "agg"+ <*> o .: "col"++instance FromJSON PlanNode where+ parseJSON = withObject "PlanNode" $ \o -> do+ op <- o .: "op" :: Parser T.Text+ case op of+ "ReadCsv" -> ReadCsv <$> o .: "path"+ "ReadTsv" -> ReadTsv <$> o .: "path"+ "FromArrow" -> FromArrow <$> o .: "schema" <*> o .: "array"+ "Select" -> Select <$> o .: "cols" <*> o .: "input"+ "GroupBy" -> GroupBy <$> o .: "keys" <*> o .: "aggregations" <*> o .: "input"+ "Sort" -> Sort <$> o .: "cols" <*> o .: "ascending" <*> o .: "input"+ "Limit" -> Limit <$> o .: "n" <*> o .: "input"+ "Filter" -> Filter <$> o .: "predicate" <*> o .: "input"+ "Derive" -> Derive <$> o .: "name" <*> o .: "expr" <*> o .: "input"+ "Exclude" -> Exclude <$> o .: "cols" <*> o .: "input"+ "Rename" -> Rename <$> o .: "pairs" <*> o .: "input"+ "Distinct" -> Distinct <$> o .: "input"+ "TakeLast" -> TakeLast <$> o .: "n" <*> o .: "input"+ "Drop" -> Drop <$> o .: "n" <*> o .: "input"+ "DropLast" -> DropLast <$> o .: "n" <*> o .: "input"+ "Range" -> Range <$> o .: "start" <*> o .: "end" <*> o .: "input"+ "Join" ->+ Join+ <$> o .: "how"+ <*> o .: "on"+ <*> o .: "left"+ <*> o .: "right"+ "Describe" -> Describe <$> o .: "input"+ "Correlation" ->+ Correlation+ <$> o .: "first"+ <*> o .: "second"+ <*> o .: "input"+ "Frequencies" -> Frequencies <$> o .: "col" <*> o .: "input"+ "ReadParquet" -> ReadParquet <$> o .: "path"+ "ReadJson" -> ReadJson <$> o .: "path"+ "WriteCsv" -> WriteCsv <$> o .: "path" <*> o .: "input"+ "ScanCsv" -> ScanCsv <$> o .: "path" <*> o .: "schema"+ "ScanParquet" -> ScanParquet <$> o .: "path" <*> o .: "schema"+ _ -> fail $ "DataFrame.IR: unknown op: " ++ T.unpack op++executePlan :: CsvReader -> PlanNode -> IO DataFrame+executePlan _reader (ReadCsv path) =+ readSeparated defaultReadOptions path+executePlan _reader (ReadTsv path) =+ readTsv path+executePlan _reader (FromArrow schemaAddr arrayAddr) =+ arrowToDataframe+ (wordPtrToPtr (fromIntegral schemaAddr))+ (wordPtrToPtr (fromIntegral arrayAddr))+executePlan reader (Select cols node) =+ select cols <$> executePlan reader node+executePlan reader (GroupBy keys aggs node) = do+ df <- executePlan reader node+ nes <- mapM (buildNamedExpr df) aggs+ return $ aggregate nes (groupBy keys df)+executePlan reader (Sort cols ascending node) = do+ df <- executePlan reader node+ let orders = map (\c -> mkSortOrder ascending c (unsafeGetColumn c df)) cols+ return $ sortBy orders df+executePlan reader (Limit k node) =+ Subset.take k <$> executePlan reader node+executePlan reader (Filter predJson node) = do+ df <- executePlan reader node+ case decodeExprAt @Bool predJson of+ Right pred_ -> return $ filterWhere pred_ df+ Left err -> ioError $ userError $ "DataFrame.IR.Filter: " <> err+executePlan reader (Derive name exprJson node) = do+ df <- executePlan reader node+ case decodeExprAny exprJson of+ Right (SomeExpr _trep expr) -> return $ derive name expr df+ Left err -> ioError $ userError $ "DataFrame.IR.Derive: " <> err+executePlan reader (Exclude cols node) =+ exclude cols <$> executePlan reader node+executePlan reader (Rename pairs node) =+ renameMany pairs <$> executePlan reader node+executePlan reader (Distinct node) =+ distinct <$> executePlan reader node+executePlan reader (TakeLast n node) =+ Subset.takeLast n <$> executePlan reader node+executePlan reader (Drop n node) =+ Subset.drop n <$> executePlan reader node+executePlan reader (DropLast n node) =+ Subset.dropLast n <$> executePlan reader node+executePlan reader (Range start end node) =+ range (start, end) <$> executePlan reader node+executePlan reader (Join how on leftPlan rightPlan) = do+ left <- executePlan reader leftPlan+ right <- executePlan reader rightPlan+ jt <- case how of+ "inner" -> return INNER+ "left" -> return LEFT+ "right" -> return RIGHT+ "outer" -> return FULL_OUTER+ "full_outer" -> return FULL_OUTER+ other ->+ ioError . userError $+ "DataFrame.IR.Join: unknown join type " <> T.unpack other+ return $ join jt on left right+executePlan reader (Describe node) = Stats.summarize <$> executePlan reader node+executePlan reader (Correlation a b node) = do+ df <- executePlan reader node+ let r = Stats.correlation a b df+ valueCol = case r of+ Just d -> V.singleton d+ Nothing -> V.singleton (0 / 0 :: Double)+ return $+ insertVector "first" (V.singleton a) $+ insertVector "second" (V.singleton b) $+ insertVector "correlation" valueCol mempty+executePlan reader (Frequencies colName node) = do+ df <- executePlan reader node+ runFrequencies colName df+executePlan _reader (ReadParquet path) = Parquet.readParquet path+executePlan _reader (ReadJson path) = readJSON path+executePlan reader (WriteCsv path node) = do+ df <- executePlan reader node+ writeCsv path df+ return df+executePlan reader (ScanCsv path schemaPairs) = do+ schema <- buildSchema schemaPairs+ Lazy.runDataFrame (Lazy.scanCsvWith reader schema (T.pack path))+executePlan _reader (ScanParquet path schemaPairs) = do+ schema <- buildSchema schemaPairs+ Lazy.runDataFrame (Lazy.scanParquet schema (T.pack path))++{- | Build a SortOrder from a column's runtime type.+Uses type dispatch to recover Ord for known column types.+-}+mkSortOrder :: Bool -> T.Text -> Column -> SortOrder+mkSortOrder isAsc name col = dispatchType (columnTypeRep col)+ where+ columnTypeRep :: Column -> SomeTypeRep+ columnTypeRep (UnboxedColumn _ (_ :: VU.Vector a)) = SomeTypeRep (typeRep @a)+ columnTypeRep (BoxedColumn _ (_ :: V.Vector a)) = SomeTypeRep (typeRep @a)+ columnTypeRep (PackedText _ _) = SomeTypeRep (typeRep @T.Text)+ columnTypeRep c@(MergedColumn _ _) = columnTypeRep (mergedHead c)+ mk :: (Columnable a, Ord a) => Expr a -> SortOrder+ mk = if isAsc then Asc else Desc+ dispatchType (SomeTypeRep tr)+ | Just HRefl <- eqTypeRep tr (typeRep @Int) = mk (Col @Int name)+ | Just HRefl <- eqTypeRep tr (typeRep @Int8) = mk (Col @Int8 name)+ | Just HRefl <- eqTypeRep tr (typeRep @Int16) = mk (Col @Int16 name)+ | Just HRefl <- eqTypeRep tr (typeRep @Int32) = mk (Col @Int32 name)+ | Just HRefl <- eqTypeRep tr (typeRep @Int64) = mk (Col @Int64 name)+ | Just HRefl <- eqTypeRep tr (typeRep @Word) = mk (Col @Word name)+ | Just HRefl <- eqTypeRep tr (typeRep @Word8) = mk (Col @Word8 name)+ | Just HRefl <- eqTypeRep tr (typeRep @Word16) = mk (Col @Word16 name)+ | Just HRefl <- eqTypeRep tr (typeRep @Word32) = mk (Col @Word32 name)+ | Just HRefl <- eqTypeRep tr (typeRep @Word64) = mk (Col @Word64 name)+ | Just HRefl <- eqTypeRep tr (typeRep @Integer) = mk (Col @Integer name)+ | Just HRefl <- eqTypeRep tr (typeRep @Double) = mk (Col @Double name)+ | Just HRefl <- eqTypeRep tr (typeRep @Float) = mk (Col @Float name)+ | Just HRefl <- eqTypeRep tr (typeRep @Bool) = mk (Col @Bool name)+ | Just HRefl <- eqTypeRep tr (typeRep @Char) = mk (Col @Char name)+ | Just HRefl <- eqTypeRep tr (typeRep @T.Text) = mk (Col @T.Text name)+ | Just HRefl <- eqTypeRep tr (typeRep @String) = mk (Col @String name)+ | Just HRefl <- eqTypeRep tr (typeRep @BS.ByteString) =+ mk (Col @BS.ByteString name)+ | otherwise = error $ "mkSortOrder: unsupported column type: " ++ show tr++-- | Dispatch aggregation by fn name and runtime column type.+buildNamedExpr :: DataFrame -> AggSpec -> IO NamedExpr+buildNamedExpr df (AggSpec name fn colName) =+ case fn of+ "count" -> countExpr name colName (unsafeGetColumn colName df)+ "sum" -> sumExpr name colName (unsafeGetColumn colName df)+ "mean" -> meanExpr name colName (unsafeGetColumn colName df)+ "min" -> minMaxExpr Functions.minimum name colName (unsafeGetColumn colName df)+ "max" -> minMaxExpr Functions.maximum name colName (unsafeGetColumn colName df)+ "median" -> doubleStatExpr Functions.median name colName (unsafeGetColumn colName df)+ "variance" -> doubleStatExpr Functions.variance name colName (unsafeGetColumn colName df)+ "std" -> doubleStatExpr stdDevExpr name colName (unsafeGetColumn colName df)+ other ->+ ioError $+ userError $+ "DataFrame.IR: unknown aggregation '" ++ T.unpack other ++ "'"++-- | Variance → standard deviation; sqrt of the underlying variance Expr.+stdDevExpr :: (Columnable a, Real a, VU.Unbox a) => Expr a -> Expr Double+stdDevExpr e = sqrt (Functions.variance e)++-- | Build a 'Schema' from a list of (col, type-tag) pairs sent over the wire.+buildSchema :: [(T.Text, T.Text)] -> IO Schema+buildSchema pairs = do+ sch <- mapM resolve pairs+ return (makeSchema sch)+ where+ resolve (name, tag) = case tag of+ "int" -> return (name, schemaType @Int)+ "int8" -> return (name, schemaType @Int8)+ "int16" -> return (name, schemaType @Int16)+ "int32" -> return (name, schemaType @Int32)+ "int64" -> return (name, schemaType @Int64)+ "double" -> return (name, schemaType @Double)+ "float" -> return (name, schemaType @Float)+ "bool" -> return (name, schemaType @Bool)+ "text" -> return (name, schemaType @T.Text)+ "string" -> return (name, schemaType @String)+ other ->+ ioError . userError $+ "DataFrame.IR.buildSchema: unsupported schema type tag '"+ ++ T.unpack other+ ++ "' for column '"+ ++ T.unpack name+ ++ "'"++-- | Dispatch 'frequencies' on the column's runtime element type.+runFrequencies :: T.Text -> DataFrame -> IO DataFrame+runFrequencies colName df = dispatchType (columnTypeRep (unsafeGetColumn colName df))+ where+ columnTypeRep :: Column -> SomeTypeRep+ columnTypeRep (UnboxedColumn _ (_ :: VU.Vector a)) = SomeTypeRep (typeRep @a)+ columnTypeRep (BoxedColumn _ (_ :: V.Vector a)) = SomeTypeRep (typeRep @a)+ columnTypeRep (PackedText _ _) = SomeTypeRep (typeRep @T.Text)+ columnTypeRep c@(MergedColumn _ _) = columnTypeRep (mergedHead c)++ fr :: forall a. (Columnable a, Ord a) => IO DataFrame+ fr = return $ Stats.frequencies (Col @a colName) df++ dispatchType :: SomeTypeRep -> IO DataFrame+ dispatchType (SomeTypeRep tr)+ | Just HRefl <- eqTypeRep tr (typeRep @Int) = fr @Int+ | Just HRefl <- eqTypeRep tr (typeRep @Int8) = fr @Int8+ | Just HRefl <- eqTypeRep tr (typeRep @Int16) = fr @Int16+ | Just HRefl <- eqTypeRep tr (typeRep @Int32) = fr @Int32+ | Just HRefl <- eqTypeRep tr (typeRep @Int64) = fr @Int64+ | Just HRefl <- eqTypeRep tr (typeRep @Word) = fr @Word+ | Just HRefl <- eqTypeRep tr (typeRep @Word8) = fr @Word8+ | Just HRefl <- eqTypeRep tr (typeRep @Word16) = fr @Word16+ | Just HRefl <- eqTypeRep tr (typeRep @Word32) = fr @Word32+ | Just HRefl <- eqTypeRep tr (typeRep @Word64) = fr @Word64+ | Just HRefl <- eqTypeRep tr (typeRep @Integer) = fr @Integer+ | Just HRefl <- eqTypeRep tr (typeRep @Double) = fr @Double+ | Just HRefl <- eqTypeRep tr (typeRep @Float) = fr @Float+ | Just HRefl <- eqTypeRep tr (typeRep @Bool) = fr @Bool+ | Just HRefl <- eqTypeRep tr (typeRep @Char) = fr @Char+ | Just HRefl <- eqTypeRep tr (typeRep @T.Text) = fr @T.Text+ | Just HRefl <- eqTypeRep tr (typeRep @String) = fr @String+ | otherwise =+ ioError . userError $+ "DataFrame.IR.Frequencies: unsupported column type for '"+ ++ T.unpack colName+ ++ "'"++countExpr :: T.Text -> T.Text -> Column -> IO NamedExpr+countExpr name colName (UnboxedColumn Nothing (_ :: VU.Vector a)) = return $ name .= count (Col @a colName)+countExpr name colName (UnboxedColumn (Just _) (_ :: VU.Vector a)) = return $ name .= count (Col @(Maybe a) colName)+countExpr name colName (BoxedColumn Nothing (_ :: V.Vector a)) = return $ name .= count (Col @a colName)+countExpr name colName (BoxedColumn (Just _) (_ :: V.Vector a)) = return $ name .= count (Col @(Maybe a) colName)+countExpr name colName (PackedText Nothing _) = return $ name .= count (Col @T.Text colName)+countExpr name colName (PackedText (Just _) _) = return $ name .= count (Col @(Maybe T.Text) colName)+countExpr name colName c@(MergedColumn _ _) = countExpr name colName (mergedHead c)++sumExpr :: T.Text -> T.Text -> Column -> IO NamedExpr+sumExpr name colName (UnboxedColumn Nothing (_ :: VU.Vector a))+ | Just Refl <- testEquality (typeRep @a) (typeRep @Int) =+ return $ name .= Functions.sum (Col @Int colName)+ | Just Refl <- testEquality (typeRep @a) (typeRep @Double) =+ return $ name .= Functions.sum (Col @Double colName)+sumExpr name colName (UnboxedColumn (Just _) (_ :: VU.Vector a))+ | Just Refl <- testEquality (typeRep @a) (typeRep @Int) =+ return $ name .= sumMaybe (Col @(Maybe Int) colName)+ | Just Refl <- testEquality (typeRep @a) (typeRep @Double) =+ return $ name .= sumMaybe (Col @(Maybe Double) colName)+sumExpr name colName (BoxedColumn Nothing (_ :: V.Vector a))+ | Just Refl <- testEquality (typeRep @a) (typeRep @Int) =+ return $ name .= Functions.sum (Col @Int colName)+ | Just Refl <- testEquality (typeRep @a) (typeRep @Double) =+ return $ name .= Functions.sum (Col @Double colName)+sumExpr name colName (BoxedColumn (Just _) (_ :: V.Vector a))+ | Just Refl <- testEquality (typeRep @a) (typeRep @Int) =+ return $ name .= sumMaybe (Col @(Maybe Int) colName)+ | Just Refl <- testEquality (typeRep @a) (typeRep @Double) =+ return $ name .= sumMaybe (Col @(Maybe Double) colName)+sumExpr _ colName _ =+ ioError $+ userError $+ "DataFrame.IR: sum: unsupported column type for '" ++ T.unpack colName ++ "'"++meanExpr :: T.Text -> T.Text -> Column -> IO NamedExpr+meanExpr name colName (UnboxedColumn Nothing (_ :: VU.Vector a))+ | Just Refl <- testEquality (typeRep @a) (typeRep @Int) =+ return $ name .= mean (Col @Int colName)+ | Just Refl <- testEquality (typeRep @a) (typeRep @Double) =+ return $ name .= mean (Col @Double colName)+meanExpr name colName (UnboxedColumn (Just _) (_ :: VU.Vector a))+ | Just Refl <- testEquality (typeRep @a) (typeRep @Double) =+ return $ name .= meanMaybe (Col @(Maybe Double) colName)+ | Just Refl <- testEquality (typeRep @a) (typeRep @Int) =+ return $ name .= meanMaybe (Col @(Maybe Int) colName)+meanExpr name colName (BoxedColumn Nothing (_ :: V.Vector a))+ | Just Refl <- testEquality (typeRep @a) (typeRep @Double) =+ return $ name .= mean (Col @Double colName)+meanExpr name colName (BoxedColumn (Just _) (_ :: V.Vector a))+ | Just Refl <- testEquality (typeRep @a) (typeRep @Double) =+ return $ name .= meanMaybe (Col @(Maybe Double) colName)+ | Just Refl <- testEquality (typeRep @a) (typeRep @Int) =+ return $ name .= meanMaybe (Col @(Maybe Int) colName)+meanExpr _ colName _ =+ ioError $+ userError $+ "DataFrame.IR: mean: unsupported column type for '" ++ T.unpack colName ++ "'"++-- | min / max — preserve column type, require Ord.+minMaxExpr ::+ (forall a. (Columnable a, Ord a) => Expr a -> Expr a) ->+ T.Text ->+ T.Text ->+ Column ->+ IO NamedExpr+minMaxExpr op name colName (UnboxedColumn Nothing (_ :: VU.Vector a))+ | Just Refl <- testEquality (typeRep @a) (typeRep @Int) =+ return $ name .= op (Col @Int colName)+ | Just Refl <- testEquality (typeRep @a) (typeRep @Double) =+ return $ name .= op (Col @Double colName)+ | Just Refl <- testEquality (typeRep @a) (typeRep @Float) =+ return $ name .= op (Col @Float colName)+minMaxExpr op name colName (BoxedColumn Nothing (_ :: V.Vector a))+ | Just Refl <- testEquality (typeRep @a) (typeRep @T.Text) =+ return $ name .= op (Col @T.Text colName)+ | Just Refl <- testEquality (typeRep @a) (typeRep @Int) =+ return $ name .= op (Col @Int colName)+ | Just Refl <- testEquality (typeRep @a) (typeRep @Double) =+ return $ name .= op (Col @Double colName)+minMaxExpr _ _ colName _ =+ ioError . userError $+ "DataFrame.IR: min/max: unsupported column type for '"+ ++ T.unpack colName+ ++ "'"++-- | median / variance / std — return Double, require Real + Unbox.+doubleStatExpr ::+ (forall a. (Columnable a, Real a, VU.Unbox a) => Expr a -> Expr Double) ->+ T.Text ->+ T.Text ->+ Column ->+ IO NamedExpr+doubleStatExpr op name colName (UnboxedColumn Nothing (_ :: VU.Vector a))+ | Just Refl <- testEquality (typeRep @a) (typeRep @Int) =+ return $ name .= op (Col @Int colName)+ | Just Refl <- testEquality (typeRep @a) (typeRep @Double) =+ return $ name .= op (Col @Double colName)+ | Just Refl <- testEquality (typeRep @a) (typeRep @Float) =+ return $ name .= op (Col @Float colName)+doubleStatExpr op name colName (BoxedColumn Nothing (_ :: V.Vector a))+ | Just Refl <- testEquality (typeRep @a) (typeRep @Int) =+ return $ name .= op (Col @Int colName)+ | Just Refl <- testEquality (typeRep @a) (typeRep @Double) =+ return $ name .= op (Col @Double colName)+doubleStatExpr _ _ colName _ =+ ioError . userError $+ "DataFrame.IR: median/variance/std: unsupported column type for '"+ ++ T.unpack colName+ ++ "'"