{-# LANGUAGE AllowAmbiguousTypes #-}
{-# LANGUAGE BangPatterns #-}
{-# LANGUAGE ExplicitNamespaces #-}
{-# LANGUAGE FlexibleContexts #-}
{-# LANGUAGE GADTs #-}
{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE ScopedTypeVariables #-}
{-# LANGUAGE TupleSections #-}
{-# LANGUAGE TypeApplications #-}

{- | Pull-based (iterator) execution engine: each operator returns a 'Stream'
yielding the next 'DataFrame' batch or 'Nothing' at end. Bounded local sources
materialise into eager whole-frame ops; unbounded sources stream in constant memory.
-}
module DataFrame.Lazy.Internal.Executor (
    CsvReader,
    execute,
    foldBatches,
) where

import Control.Concurrent (forkIO, getNumCapabilities)
import Control.Concurrent.Async (mapConcurrently)
import Control.Concurrent.STM (atomically)
import Control.Concurrent.STM.TBQueue (newTBQueueIO, readTBQueue, writeTBQueue)
import Control.Exception (evaluate)
import Control.Monad (filterM, forM, forM_, unless, when)
import qualified Data.ByteString as BS
import qualified Data.ByteString.Char8 as C8
import Data.IORef
import Data.Int (Int16, Int32, Int64, Int8)
import qualified Data.Map as M
import qualified Data.Maybe
import qualified Data.Set as S
import qualified Data.Text as T
import Data.Type.Equality (TestEquality (testEquality), type (:~:) (Refl))
import Data.Typeable (Typeable)
import qualified Data.Vector as VB
import qualified Data.Vector.Unboxed as VU
import Data.Word (Word16, Word32, Word64, Word8)
import DataFrame.IO.CSV (CsvReader, ReadOptions (..), schemaReadOptions)
import qualified DataFrame.IO.Parquet as Parquet
import qualified DataFrame.Internal.Column as C
import qualified DataFrame.Internal.DataFrame as D
import qualified DataFrame.Internal.Expression as E
import qualified DataFrame.Lazy.IO.Binary as Bin
import DataFrame.Lazy.Internal.LogicalPlan (DataSource (..), SortOrder (..))
import DataFrame.Lazy.Internal.PhysicalPlan
import qualified DataFrame.Operations.Aggregation as Agg
import qualified DataFrame.Operations.Core as Core
import qualified DataFrame.Operations.Join as Join
import DataFrame.Operations.Merge ()
import qualified DataFrame.Operations.Permutation as Perm
import qualified DataFrame.Operations.Subset as Sub
import qualified DataFrame.Operations.Transformations as Trans
import DataFrame.Schema (elements)
import System.Directory (doesDirectoryExist, removeFile)
import System.FilePath ((</>))
import System.FilePath.Glob (glob)
import System.IO (IOMode (ReadMode), hIsEOF, withFile)
import System.IO.Temp (emptySystemTempFile)
import Type.Reflection (typeRep)

-- ---------------------------------------------------------------------------
-- Stream abstraction
-- ---------------------------------------------------------------------------

{- | A pull-based stream: each call to the action yields the next batch or
'Nothing' when the stream is exhausted.  State is captured by the closure.
-}
newtype Stream = Stream {Stream -> IO (Maybe DataFrame)
pullBatch :: IO (Maybe D.DataFrame)}

{- | Wrap an already-computed 'DataFrame' as a single-batch stream: the first
pull yields it, every subsequent pull yields 'Nothing'.  Used by blocking
operators that materialise their whole result up front.
-}
materialized :: D.DataFrame -> IO Stream
materialized :: DataFrame -> IO Stream
materialized DataFrame
df = do
    IORef (Maybe DataFrame)
ref <- Maybe DataFrame -> IO (IORef (Maybe DataFrame))
forall a. a -> IO (IORef a)
newIORef (DataFrame -> Maybe DataFrame
forall a. a -> Maybe a
Just DataFrame
df)
    Stream -> IO Stream
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return (Stream -> IO Stream)
-> (IO (Maybe DataFrame) -> Stream)
-> IO (Maybe DataFrame)
-> IO Stream
forall b c a. (b -> c) -> (a -> b) -> a -> c
. IO (Maybe DataFrame) -> Stream
Stream (IO (Maybe DataFrame) -> IO Stream)
-> IO (Maybe DataFrame) -> IO Stream
forall a b. (a -> b) -> a -> b
$ do
        Maybe DataFrame
mb <- IORef (Maybe DataFrame) -> IO (Maybe DataFrame)
forall a. IORef a -> IO a
readIORef IORef (Maybe DataFrame)
ref
        IORef (Maybe DataFrame) -> Maybe DataFrame -> IO ()
forall a. IORef a -> a -> IO ()
writeIORef IORef (Maybe DataFrame)
ref Maybe DataFrame
forall a. Maybe a
Nothing
        Maybe DataFrame -> IO (Maybe DataFrame)
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return Maybe DataFrame
mb

{- | Drain all batches from a stream and concatenate them into one DataFrame.
Columns are concatenated in a single multi-way pass (O(rows)) rather than a
left-fold of @acc <> batch@ that would recopy the accumulator on every step.
-}
collectStream :: Stream -> IO D.DataFrame
collectStream :: Stream -> IO DataFrame
collectStream Stream
stream = [DataFrame] -> IO DataFrame
go []
  where
    go :: [DataFrame] -> IO DataFrame
go [DataFrame]
acc = do
        Maybe DataFrame
mb <- Stream -> IO (Maybe DataFrame)
pullBatch Stream
stream
        case Maybe DataFrame
mb of
            Maybe DataFrame
Nothing -> DataFrame -> IO DataFrame
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return ([DataFrame] -> DataFrame
concatBatches ([DataFrame] -> [DataFrame]
forall a. [a] -> [a]
reverse [DataFrame]
acc))
            Just DataFrame
df -> [DataFrame] -> IO DataFrame
go (DataFrame
df DataFrame -> [DataFrame] -> [DataFrame]
forall a. a -> [a] -> [a]
: [DataFrame]
acc)

-- | Concatenate a list of same-schema batches column-by-column in one pass.
concatBatches :: [D.DataFrame] -> D.DataFrame
concatBatches :: [DataFrame] -> DataFrame
concatBatches [] = DataFrame
D.empty
concatBatches [DataFrame
df] = DataFrame
df
concatBatches batches :: [DataFrame]
batches@(DataFrame
first : [DataFrame]
_) =
    [(Text, Column)] -> DataFrame
D.fromNamedColumns
        [ (Text
name, [Column] -> Column
C.concatManyColumns [Text -> DataFrame -> Column
D.unsafeGetColumn Text
name DataFrame
b | DataFrame
b <- [DataFrame]
batches])
        | Text
name <- DataFrame -> [Text]
D.columnNames DataFrame
first
        ]

-- ---------------------------------------------------------------------------
-- Boundedness analysis
-- ---------------------------------------------------------------------------

{- | True when every scan leaf is a finite local source, so the whole result can
be materialised and routed through the eager whole-frame ops. False for
online/streaming sources, which keep the constant-memory streaming paths.
-}
isBounded :: PhysicalPlan -> Bool
isBounded :: PhysicalPlan -> Bool
isBounded (PhysicalScan (CsvSource{}) ScanConfig
_) = Bool
True
isBounded (PhysicalScan (CsvSourceStreaming{}) ScanConfig
_) = Bool
False
isBounded (PhysicalScan (ParquetSource [Char]
_) ScanConfig
_) = Bool
True
isBounded (PhysicalProject [Text]
_ PhysicalPlan
c) = PhysicalPlan -> Bool
isBounded PhysicalPlan
c
isBounded (PhysicalFilter Expr Bool
_ PhysicalPlan
c) = PhysicalPlan -> Bool
isBounded PhysicalPlan
c
isBounded (PhysicalDerive Text
_ UExpr
_ PhysicalPlan
c) = PhysicalPlan -> Bool
isBounded PhysicalPlan
c
isBounded (PhysicalLimit Int
_ PhysicalPlan
c) = PhysicalPlan -> Bool
isBounded PhysicalPlan
c
isBounded (PhysicalSort [(Text, SortOrder)]
_ PhysicalPlan
c) = PhysicalPlan -> Bool
isBounded PhysicalPlan
c
isBounded (PhysicalHashAggregate [Text]
_ [(Text, UExpr)]
_ PhysicalPlan
c) = PhysicalPlan -> Bool
isBounded PhysicalPlan
c
isBounded (PhysicalSpill PhysicalPlan
c [Char]
_) = PhysicalPlan -> Bool
isBounded PhysicalPlan
c
isBounded (PhysicalSourceDF Int
_ DataFrame
_) = Bool
True
isBounded (PhysicalHashJoin JoinType
_ Text
_ Text
_ PhysicalPlan
l PhysicalPlan
r) = PhysicalPlan -> Bool
isBounded PhysicalPlan
l Bool -> Bool -> Bool
&& PhysicalPlan -> Bool
isBounded PhysicalPlan
r
isBounded (PhysicalSortMergeJoin JoinType
_ Text
_ Text
_ PhysicalPlan
l PhysicalPlan
r) = PhysicalPlan -> Bool
isBounded PhysicalPlan
l Bool -> Bool -> Bool
&& PhysicalPlan -> Bool
isBounded PhysicalPlan
r

-- ---------------------------------------------------------------------------
-- Top-level entry point
-- ---------------------------------------------------------------------------

{- | Execute a physical plan, returning the complete result as a single
'DataFrame'.
-}
execute :: PhysicalPlan -> IO D.DataFrame
execute :: PhysicalPlan -> IO DataFrame
execute PhysicalPlan
plan = PhysicalPlan -> IO Stream
buildStream PhysicalPlan
plan IO Stream -> (Stream -> IO DataFrame) -> IO DataFrame
forall a b. IO a -> (a -> IO b) -> IO b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= Stream -> IO DataFrame
collectStream

{- | Fold a function over every batch produced by a physical plan.
The fold is strict in the accumulator; each batch is discarded after folding.
-}
foldBatches ::
    (b -> D.DataFrame -> IO b) -> b -> PhysicalPlan -> IO b
foldBatches :: forall b. (b -> DataFrame -> IO b) -> b -> PhysicalPlan -> IO b
foldBatches b -> DataFrame -> IO b
f b
seed PhysicalPlan
plan = do
    Stream
stream <- PhysicalPlan -> IO Stream
buildStream PhysicalPlan
plan
    let loop :: b -> IO b
loop !b
acc = do
            Maybe DataFrame
mb <- Stream -> IO (Maybe DataFrame)
pullBatch Stream
stream
            case Maybe DataFrame
mb of
                Maybe DataFrame
Nothing -> b -> IO b
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return b
acc
                Just DataFrame
batch -> do
                    !b
acc' <- b -> DataFrame -> IO b
f b
acc DataFrame
batch
                    b -> IO b
loop b
acc'
    b -> IO b
loop b
seed

-- ---------------------------------------------------------------------------
-- Per-operator stream builders
-- ---------------------------------------------------------------------------

buildStream :: PhysicalPlan -> IO Stream
buildStream :: PhysicalPlan -> IO Stream
buildStream (PhysicalScan (CsvSource [Char]
path Char
sep CsvReader
reader) ScanConfig
cfg) =
    [Char] -> Char -> CsvReader -> ScanConfig -> IO Stream
executeCsvScan [Char]
path Char
sep CsvReader
reader ScanConfig
cfg
buildStream (PhysicalScan (CsvSourceStreaming [Char]
path Char
sep CsvReader
reader) ScanConfig
cfg) =
    [Char] -> Char -> CsvReader -> ScanConfig -> IO Stream
executeCsvScanStreaming [Char]
path Char
sep CsvReader
reader ScanConfig
cfg
buildStream (PhysicalScan (ParquetSource [Char]
path) ScanConfig
cfg) =
    [Char] -> ScanConfig -> IO Stream
executeParquetScan [Char]
path ScanConfig
cfg
buildStream (PhysicalSpill PhysicalPlan
child [Char]
path) = do
    DataFrame
df <- PhysicalPlan -> IO DataFrame
execute PhysicalPlan
child
    [Char] -> DataFrame -> IO ()
Bin.spillToDisk [Char]
path DataFrame
df
    DataFrame
df' <- [Char] -> IO DataFrame
Bin.readSpilled [Char]
path
    DataFrame -> IO Stream
materialized DataFrame
df'
buildStream (PhysicalFilter Expr Bool
p PhysicalPlan
child) = do
    Stream
childStream <- PhysicalPlan -> IO Stream
buildStream PhysicalPlan
child
    Stream -> IO Stream
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return (Stream -> IO Stream)
-> (IO (Maybe DataFrame) -> Stream)
-> IO (Maybe DataFrame)
-> IO Stream
forall b c a. (b -> c) -> (a -> b) -> a -> c
. IO (Maybe DataFrame) -> Stream
Stream (IO (Maybe DataFrame) -> IO Stream)
-> IO (Maybe DataFrame) -> IO Stream
forall a b. (a -> b) -> a -> b
$
        ( do
            Maybe DataFrame
mb <- Stream -> IO (Maybe DataFrame)
pullBatch Stream
childStream
            Maybe DataFrame -> IO (Maybe DataFrame)
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return (Maybe DataFrame -> IO (Maybe DataFrame))
-> Maybe DataFrame -> IO (Maybe DataFrame)
forall a b. (a -> b) -> a -> b
$ (DataFrame -> DataFrame) -> Maybe DataFrame -> Maybe DataFrame
forall a b. (a -> b) -> Maybe a -> Maybe b
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
fmap (Expr Bool -> DataFrame -> DataFrame
Sub.filterWhere Expr Bool
p) Maybe DataFrame
mb
        )
buildStream (PhysicalProject [Text]
cols PhysicalPlan
child) = do
    Stream
childStream <- PhysicalPlan -> IO Stream
buildStream PhysicalPlan
child
    Stream -> IO Stream
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return (Stream -> IO Stream)
-> (IO (Maybe DataFrame) -> Stream)
-> IO (Maybe DataFrame)
-> IO Stream
forall b c a. (b -> c) -> (a -> b) -> a -> c
. IO (Maybe DataFrame) -> Stream
Stream (IO (Maybe DataFrame) -> IO Stream)
-> IO (Maybe DataFrame) -> IO Stream
forall a b. (a -> b) -> a -> b
$
        ( do
            Maybe DataFrame
mb <- Stream -> IO (Maybe DataFrame)
pullBatch Stream
childStream
            Maybe DataFrame -> IO (Maybe DataFrame)
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return (Maybe DataFrame -> IO (Maybe DataFrame))
-> Maybe DataFrame -> IO (Maybe DataFrame)
forall a b. (a -> b) -> a -> b
$ (DataFrame -> DataFrame) -> Maybe DataFrame -> Maybe DataFrame
forall a b. (a -> b) -> Maybe a -> Maybe b
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
fmap ([Text] -> DataFrame -> DataFrame
Sub.select [Text]
cols) Maybe DataFrame
mb
        )
buildStream (PhysicalDerive Text
name UExpr
uexpr PhysicalPlan
child) = do
    Stream
childStream <- PhysicalPlan -> IO Stream
buildStream PhysicalPlan
child
    Stream -> IO Stream
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return (Stream -> IO Stream)
-> (IO (Maybe DataFrame) -> Stream)
-> IO (Maybe DataFrame)
-> IO Stream
forall b c a. (b -> c) -> (a -> b) -> a -> c
. IO (Maybe DataFrame) -> Stream
Stream (IO (Maybe DataFrame) -> IO Stream)
-> IO (Maybe DataFrame) -> IO Stream
forall a b. (a -> b) -> a -> b
$
        ( do
            Maybe DataFrame
mb <- Stream -> IO (Maybe DataFrame)
pullBatch Stream
childStream
            Maybe DataFrame -> IO (Maybe DataFrame)
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return (Maybe DataFrame -> IO (Maybe DataFrame))
-> Maybe DataFrame -> IO (Maybe DataFrame)
forall a b. (a -> b) -> a -> b
$ (DataFrame -> DataFrame) -> Maybe DataFrame -> Maybe DataFrame
forall a b. (a -> b) -> Maybe a -> Maybe b
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
fmap ([(Text, UExpr)] -> DataFrame -> DataFrame
Trans.deriveMany [(Text
name, UExpr
uexpr)]) Maybe DataFrame
mb
        )
buildStream (PhysicalLimit Int
n PhysicalPlan
child) = do
    Stream
childStream <- PhysicalPlan -> IO Stream
buildStream PhysicalPlan
child
    IORef Int
countRef <- Int -> IO (IORef Int)
forall a. a -> IO (IORef a)
newIORef (Int
0 :: Int)
    Stream -> IO Stream
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return (Stream -> IO Stream)
-> (IO (Maybe DataFrame) -> Stream)
-> IO (Maybe DataFrame)
-> IO Stream
forall b c a. (b -> c) -> (a -> b) -> a -> c
. IO (Maybe DataFrame) -> Stream
Stream (IO (Maybe DataFrame) -> IO Stream)
-> IO (Maybe DataFrame) -> IO Stream
forall a b. (a -> b) -> a -> b
$
        ( do
            Int
remaining <- IORef Int -> IO Int
forall a. IORef a -> IO a
readIORef IORef Int
countRef
            if Int
remaining Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
>= Int
n
                then Maybe DataFrame -> IO (Maybe DataFrame)
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return Maybe DataFrame
forall a. Maybe a
Nothing
                else do
                    Maybe DataFrame
mb <- Stream -> IO (Maybe DataFrame)
pullBatch Stream
childStream
                    case Maybe DataFrame
mb of
                        Maybe DataFrame
Nothing -> Maybe DataFrame -> IO (Maybe DataFrame)
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return Maybe DataFrame
forall a. Maybe a
Nothing
                        Just DataFrame
df -> do
                            let toTake :: Int
toTake = Int -> Int -> Int
forall a. Ord a => a -> a -> a
min (DataFrame -> Int
Core.nRows DataFrame
df) (Int
n Int -> Int -> Int
forall a. Num a => a -> a -> a
- Int
remaining)
                            IORef Int -> (Int -> Int) -> IO ()
forall a. IORef a -> (a -> a) -> IO ()
modifyIORef' IORef Int
countRef (Int -> Int -> Int
forall a. Num a => a -> a -> a
+ Int
toTake)
                            Maybe DataFrame -> IO (Maybe DataFrame)
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return (Maybe DataFrame -> IO (Maybe DataFrame))
-> Maybe DataFrame -> IO (Maybe DataFrame)
forall a b. (a -> b) -> a -> b
$ DataFrame -> Maybe DataFrame
forall a. a -> Maybe a
Just (Int -> DataFrame -> DataFrame
Sub.take Int
toTake DataFrame
df)
        )
buildStream (PhysicalSort [(Text, SortOrder)]
cols PhysicalPlan
child) = do
    DataFrame
df <- PhysicalPlan -> IO DataFrame
execute PhysicalPlan
child
    let sortOrds :: [SortOrder]
sortOrds = ((Text, SortOrder) -> SortOrder)
-> [(Text, SortOrder)] -> [SortOrder]
forall a b. (a -> b) -> [a] -> [b]
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
fmap (DataFrame -> (Text, SortOrder) -> SortOrder
toPermSortOrder DataFrame
df) [(Text, SortOrder)]
cols
    DataFrame -> IO Stream
materialized ([SortOrder] -> DataFrame -> DataFrame
Perm.sortBy [SortOrder]
sortOrds DataFrame
df)
buildStream (PhysicalHashAggregate [Text]
keys [(Text, UExpr)]
aggs PhysicalPlan
child)
    | PhysicalPlan -> Bool
isBounded PhysicalPlan
child = do
        DataFrame
df <- PhysicalPlan -> IO DataFrame
execute PhysicalPlan
child
        let result :: DataFrame
result = [(Text, UExpr)] -> GroupedDataFrame -> DataFrame
Agg.aggregate [(Text, UExpr)]
aggs ([Text] -> DataFrame -> GroupedDataFrame
Agg.groupBy [Text]
keys DataFrame
df)
        DataFrame -> IO Stream
materialized DataFrame
result
buildStream (PhysicalHashAggregate [Text]
keys [(Text, UExpr)]
aggs PhysicalPlan
child) = do
    Stream
childStream <- PhysicalPlan -> IO Stream
buildStream PhysicalPlan
child
    if ((Text, UExpr) -> Bool) -> [(Text, UExpr)] -> Bool
forall (t :: * -> *) a. Foldable t => (a -> Bool) -> t a -> Bool
all (UExpr -> Bool
isStreamableAgg (UExpr -> Bool)
-> ((Text, UExpr) -> UExpr) -> (Text, UExpr) -> Bool
forall b c a. (b -> c) -> (a -> b) -> a -> c
. (Text, UExpr) -> UExpr
forall a b. (a, b) -> b
snd) [(Text, UExpr)]
aggs
        then do
            let ([(Text, UExpr)]
partialAggs, [(Text, UExpr)]
mergeAggs, DataFrame -> DataFrame
finalizer) = [(Text, UExpr)]
-> ([(Text, UExpr)], [(Text, UExpr)], DataFrame -> DataFrame)
buildAggPlan [(Text, UExpr)]
aggs
            Int
nCaps <- IO Int
getNumCapabilities
            let workers :: Int
workers = Int -> Int -> Int
forall a. Ord a => a -> a -> a
max Int
1 Int
nCaps
            [Maybe DataFrame]
partials <-
                (Int -> IO (Maybe DataFrame)) -> [Int] -> IO [Maybe DataFrame]
forall (t :: * -> *) a b.
Traversable t =>
(a -> IO b) -> t a -> IO (t b)
mapConcurrently
                    (\Int
_ -> Stream
-> [Text]
-> [(Text, UExpr)]
-> [(Text, UExpr)]
-> IO (Maybe DataFrame)
workerLoop Stream
childStream [Text]
keys [(Text, UExpr)]
partialAggs [(Text, UExpr)]
mergeAggs)
                    [Int
1 .. Int
workers]
            Maybe DataFrame
mFinal <-
                let nonEmpty :: [DataFrame]
nonEmpty = [Maybe DataFrame] -> [DataFrame]
forall a. [Maybe a] -> [a]
Data.Maybe.catMaybes [Maybe DataFrame]
partials
                 in case [DataFrame]
nonEmpty of
                        [] -> Maybe DataFrame -> IO (Maybe DataFrame)
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return Maybe DataFrame
forall a. Maybe a
Nothing
                        [DataFrame
single] -> Maybe DataFrame -> IO (Maybe DataFrame)
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return (DataFrame -> Maybe DataFrame
forall a. a -> Maybe a
Just (DataFrame -> DataFrame
finalizer DataFrame
single))
                        (DataFrame
a : [DataFrame]
rest) -> do
                            !DataFrame
merged <- [Text]
-> [(Text, UExpr)] -> DataFrame -> [DataFrame] -> IO DataFrame
mergePartials [Text]
keys [(Text, UExpr)]
mergeAggs DataFrame
a [DataFrame]
rest
                            Maybe DataFrame -> IO (Maybe DataFrame)
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return (DataFrame -> Maybe DataFrame
forall a. a -> Maybe a
Just (DataFrame -> DataFrame
finalizer DataFrame
merged))
            IORef (Maybe DataFrame)
ref <- Maybe DataFrame -> IO (IORef (Maybe DataFrame))
forall a. a -> IO (IORef a)
newIORef Maybe DataFrame
mFinal
            Stream -> IO Stream
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return (Stream -> IO Stream)
-> (IO (Maybe DataFrame) -> Stream)
-> IO (Maybe DataFrame)
-> IO Stream
forall b c a. (b -> c) -> (a -> b) -> a -> c
. IO (Maybe DataFrame) -> Stream
Stream (IO (Maybe DataFrame) -> IO Stream)
-> IO (Maybe DataFrame) -> IO Stream
forall a b. (a -> b) -> a -> b
$ do
                Maybe DataFrame
mb <- IORef (Maybe DataFrame) -> IO (Maybe DataFrame)
forall a. IORef a -> IO a
readIORef IORef (Maybe DataFrame)
ref
                IORef (Maybe DataFrame) -> Maybe DataFrame -> IO ()
forall a. IORef a -> a -> IO ()
writeIORef IORef (Maybe DataFrame)
ref Maybe DataFrame
forall a. Maybe a
Nothing
                Maybe DataFrame -> IO (Maybe DataFrame)
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return Maybe DataFrame
mb
        else do
            DataFrame
df <- Stream -> IO DataFrame
collectStream Stream
childStream
            DataFrame -> IO Stream
materialized ([(Text, UExpr)] -> GroupedDataFrame -> DataFrame
Agg.aggregate [(Text, UExpr)]
aggs ([Text] -> DataFrame -> GroupedDataFrame
Agg.groupBy [Text]
keys DataFrame
df))
buildStream (PhysicalSourceDF Int
bs DataFrame
df) = do
    let total :: Int
total = DataFrame -> Int
Core.nRows DataFrame
df
    IORef Int
posRef <- Int -> IO (IORef Int)
forall a. a -> IO (IORef a)
newIORef (Int
0 :: Int)
    Stream -> IO Stream
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return (Stream -> IO Stream)
-> (IO (Maybe DataFrame) -> Stream)
-> IO (Maybe DataFrame)
-> IO Stream
forall b c a. (b -> c) -> (a -> b) -> a -> c
. IO (Maybe DataFrame) -> Stream
Stream (IO (Maybe DataFrame) -> IO Stream)
-> IO (Maybe DataFrame) -> IO Stream
forall a b. (a -> b) -> a -> b
$ do
        Int
i <- IORef Int -> IO Int
forall a. IORef a -> IO a
readIORef IORef Int
posRef
        if Int
i Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
>= Int
total
            then Maybe DataFrame -> IO (Maybe DataFrame)
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return Maybe DataFrame
forall a. Maybe a
Nothing
            else do
                let n :: Int
n = Int -> Int -> Int
forall a. Ord a => a -> a -> a
min Int
bs (Int
total Int -> Int -> Int
forall a. Num a => a -> a -> a
- Int
i)
                    batch :: DataFrame
batch = (Int, Int) -> DataFrame -> DataFrame
Sub.range (Int
i, Int
i Int -> Int -> Int
forall a. Num a => a -> a -> a
+ Int
n) DataFrame
df
                IORef Int -> Int -> IO ()
forall a. IORef a -> a -> IO ()
writeIORef IORef Int
posRef (Int
i Int -> Int -> Int
forall a. Num a => a -> a -> a
+ Int
n)
                Maybe DataFrame -> IO (Maybe DataFrame)
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return (DataFrame -> Maybe DataFrame
forall a. a -> Maybe a
Just DataFrame
batch)
buildStream (PhysicalHashJoin JoinType
jt Text
leftKey Text
rightKey PhysicalPlan
leftPlan PhysicalPlan
rightPlan)
    | PhysicalPlan -> Bool
isBounded PhysicalPlan
leftPlan = do
        DataFrame
leftDf <- PhysicalPlan -> IO DataFrame
execute PhysicalPlan
leftPlan
        DataFrame
rightDf <- PhysicalPlan -> IO DataFrame
execute PhysicalPlan
rightPlan
        DataFrame -> IO Stream
materialized (JoinType -> Text -> Text -> DataFrame -> DataFrame -> DataFrame
performJoin JoinType
jt Text
leftKey Text
rightKey DataFrame
leftDf DataFrame
rightDf)
buildStream (PhysicalHashJoin JoinType
jt Text
leftKey Text
rightKey PhysicalPlan
leftPlan PhysicalPlan
rightPlan) =
    case JoinType
jt of
        JoinType
Join.INNER -> (Set Text
 -> DataFrame -> DataFrame -> Vector Int -> Vector Int -> DataFrame)
-> IO Stream
streamingHashJoin Set Text
-> DataFrame -> DataFrame -> Vector Int -> Vector Int -> DataFrame
assembleInnerBatch
        JoinType
Join.LEFT -> (Set Text
 -> DataFrame -> DataFrame -> Vector Int -> Vector Int -> DataFrame)
-> IO Stream
streamingHashJoin Set Text
-> DataFrame -> DataFrame -> Vector Int -> Vector Int -> DataFrame
assembleLeftBatch
        JoinType
_ -> do
            DataFrame
leftDf <- PhysicalPlan -> IO DataFrame
execute PhysicalPlan
leftPlan
            DataFrame
rightDf <- PhysicalPlan -> IO DataFrame
execute PhysicalPlan
rightPlan
            DataFrame -> IO Stream
materialized (JoinType -> Text -> Text -> DataFrame -> DataFrame -> DataFrame
performJoin JoinType
jt Text
leftKey Text
rightKey DataFrame
leftDf DataFrame
rightDf)
  where
    streamingHashJoin :: (Set Text
 -> DataFrame -> DataFrame -> Vector Int -> Vector Int -> DataFrame)
-> IO Stream
streamingHashJoin Set Text
-> DataFrame -> DataFrame -> Vector Int -> Vector Int -> DataFrame
assembleFn = do
        DataFrame
rightDf <- PhysicalPlan -> IO DataFrame
execute PhysicalPlan
rightPlan
        let rightDf' :: DataFrame
rightDf' =
                if Text
leftKey Text -> Text -> Bool
forall a. Eq a => a -> a -> Bool
== Text
rightKey
                    then DataFrame
rightDf
                    else Text -> Text -> DataFrame -> DataFrame
Core.rename Text
rightKey Text
leftKey DataFrame
rightDf
            joinKey :: Text
joinKey = Text
leftKey
            csSet :: Set Text
csSet = [Text] -> Set Text
forall a. Ord a => [a] -> Set a
S.fromList [Text
joinKey]
            rightHashes :: Vector Int
rightHashes = [Text] -> DataFrame -> Vector Int
Join.buildHashColumn [Text
joinKey] DataFrame
rightDf'
            ci :: CompactIndex
ci = Vector Int -> CompactIndex
Join.buildCompactIndex Vector Int
rightHashes
        Stream
leftStream <- PhysicalPlan -> IO Stream
buildStream PhysicalPlan
leftPlan
        Stream -> IO Stream
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return (Stream -> IO Stream)
-> (IO (Maybe DataFrame) -> Stream)
-> IO (Maybe DataFrame)
-> IO Stream
forall b c a. (b -> c) -> (a -> b) -> a -> c
. IO (Maybe DataFrame) -> Stream
Stream (IO (Maybe DataFrame) -> IO Stream)
-> IO (Maybe DataFrame) -> IO Stream
forall a b. (a -> b) -> a -> b
$ do
            Maybe DataFrame
mBatch <- Stream -> IO (Maybe DataFrame)
pullBatch Stream
leftStream
            case Maybe DataFrame
mBatch of
                Maybe DataFrame
Nothing -> Maybe DataFrame -> IO (Maybe DataFrame)
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return Maybe DataFrame
forall a. Maybe a
Nothing
                Just DataFrame
probeBatch -> do
                    let probeHashes :: Vector Int
probeHashes = [Text] -> DataFrame -> Vector Int
Join.buildHashColumn [Text
joinKey] DataFrame
probeBatch
                        (Vector Int
probeIxs, Vector Int
buildIxs) = CompactIndex -> Vector Int -> (Vector Int, Vector Int)
Join.hashProbeKernel CompactIndex
ci Vector Int
probeHashes
                    Maybe DataFrame -> IO (Maybe DataFrame)
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return (Maybe DataFrame -> IO (Maybe DataFrame))
-> (DataFrame -> Maybe DataFrame)
-> DataFrame
-> IO (Maybe DataFrame)
forall b c a. (b -> c) -> (a -> b) -> a -> c
. DataFrame -> Maybe DataFrame
forall a. a -> Maybe a
Just (DataFrame -> IO (Maybe DataFrame))
-> DataFrame -> IO (Maybe DataFrame)
forall a b. (a -> b) -> a -> b
$ Set Text
-> DataFrame -> DataFrame -> Vector Int -> Vector Int -> DataFrame
assembleFn Set Text
csSet DataFrame
probeBatch DataFrame
rightDf' Vector Int
probeIxs Vector Int
buildIxs

    assembleLeftBatch :: Set Text
-> DataFrame -> DataFrame -> Vector Int -> Vector Int -> DataFrame
assembleLeftBatch Set Text
csSet DataFrame
probeBatch DataFrame
rightDf' Vector Int
probeIxs Vector Int
buildIxs =
        let batchN :: Int
batchN = DataFrame -> Int
Core.nRows DataFrame
probeBatch
            matched :: Vector Bool
matched =
                (Bool -> Bool -> Bool)
-> Vector Bool -> Vector (Int, Bool) -> Vector Bool
forall a b.
(Unbox a, Unbox b) =>
(a -> b -> a) -> Vector a -> Vector (Int, b) -> Vector a
VU.accumulate
                    (\Bool
_ Bool
b -> Bool
b)
                    (Int -> Bool -> Vector Bool
forall a. Unbox a => Int -> a -> Vector a
VU.replicate Int
batchN Bool
False)
                    ((Int -> (Int, Bool)) -> Vector Int -> Vector (Int, Bool)
forall a b. (Unbox a, Unbox b) => (a -> b) -> Vector a -> Vector b
VU.map (,Bool
True) Vector Int
probeIxs)
            unmatchedIxs :: Vector Int
unmatchedIxs = (Bool -> Bool) -> Vector Bool -> Vector Int
forall a. Unbox a => (a -> Bool) -> Vector a -> Vector Int
VU.findIndices Bool -> Bool
not Vector Bool
matched
            allProbeIxs :: Vector Int
allProbeIxs = Vector Int
probeIxs Vector Int -> Vector Int -> Vector Int
forall a. Unbox a => Vector a -> Vector a -> Vector a
VU.++ Vector Int
unmatchedIxs
            allBuildIxs :: Vector Int
allBuildIxs = Vector Int
buildIxs Vector Int -> Vector Int -> Vector Int
forall a. Unbox a => Vector a -> Vector a -> Vector a
VU.++ Int -> Int -> Vector Int
forall a. Unbox a => Int -> a -> Vector a
VU.replicate (Vector Int -> Int
forall a. Unbox a => Vector a -> Int
VU.length Vector Int
unmatchedIxs) (-Int
1)
         in Set Text
-> DataFrame -> DataFrame -> Vector Int -> Vector Int -> DataFrame
Join.assembleLeft Set Text
csSet DataFrame
probeBatch DataFrame
rightDf' Vector Int
allProbeIxs Vector Int
allBuildIxs

    assembleInnerBatch :: Set Text
-> DataFrame -> DataFrame -> Vector Int -> Vector Int -> DataFrame
assembleInnerBatch = Set Text
-> DataFrame -> DataFrame -> Vector Int -> Vector Int -> DataFrame
Join.assembleInner
buildStream (PhysicalSortMergeJoin JoinType
jt Text
leftKey Text
rightKey PhysicalPlan
leftPlan PhysicalPlan
rightPlan) = do
    DataFrame
leftDf <- PhysicalPlan -> IO DataFrame
execute PhysicalPlan
leftPlan
    DataFrame
rightDf <- PhysicalPlan -> IO DataFrame
execute PhysicalPlan
rightPlan
    DataFrame -> IO Stream
materialized (JoinType -> Text -> Text -> DataFrame -> DataFrame -> DataFrame
performJoin JoinType
jt Text
leftKey Text
rightKey DataFrame
leftDf DataFrame
rightDf)

-- ---------------------------------------------------------------------------
-- Streaming aggregation helpers
-- ---------------------------------------------------------------------------

{- | One worker's loop: pull batches off the shared child stream until
exhausted, building up a per-worker accumulator.
-}
workerLoop ::
    Stream ->
    [T.Text] ->
    [E.NamedExpr] ->
    [E.NamedExpr] ->
    IO (Maybe D.DataFrame)
workerLoop :: Stream
-> [Text]
-> [(Text, UExpr)]
-> [(Text, UExpr)]
-> IO (Maybe DataFrame)
workerLoop Stream
childStream [Text]
keys [(Text, UExpr)]
partialAggs [(Text, UExpr)]
mergeAggs = Maybe DataFrame -> IO (Maybe DataFrame)
loop Maybe DataFrame
forall a. Maybe a
Nothing
  where
    loop :: Maybe DataFrame -> IO (Maybe DataFrame)
loop !Maybe DataFrame
acc = do
        Maybe DataFrame
mb <- Stream -> IO (Maybe DataFrame)
pullBatch Stream
childStream
        case Maybe DataFrame
mb of
            Maybe DataFrame
Nothing -> Maybe DataFrame -> IO (Maybe DataFrame)
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return Maybe DataFrame
acc
            Just DataFrame
batch -> do
                !DataFrame
partial <-
                    DataFrame -> IO DataFrame
forall a. a -> IO a
evaluate (DataFrame -> IO DataFrame)
-> (DataFrame -> DataFrame) -> DataFrame -> IO DataFrame
forall b c a. (b -> c) -> (a -> b) -> a -> c
. DataFrame -> DataFrame
D.forceDataFrame (DataFrame -> IO DataFrame) -> DataFrame -> IO DataFrame
forall a b. (a -> b) -> a -> b
$
                        [(Text, UExpr)] -> GroupedDataFrame -> DataFrame
Agg.aggregate [(Text, UExpr)]
partialAggs ([Text] -> DataFrame -> GroupedDataFrame
Agg.groupBy [Text]
keys DataFrame
batch)
                !Maybe DataFrame
next <- case Maybe DataFrame
acc of
                    Maybe DataFrame
Nothing -> Maybe DataFrame -> IO (Maybe DataFrame)
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return (DataFrame -> Maybe DataFrame
forall a. a -> Maybe a
Just DataFrame
partial)
                    Just DataFrame
a -> do
                        !DataFrame
merged <-
                            DataFrame -> IO DataFrame
forall a. a -> IO a
evaluate (DataFrame -> IO DataFrame)
-> (DataFrame -> DataFrame) -> DataFrame -> IO DataFrame
forall b c a. (b -> c) -> (a -> b) -> a -> c
. DataFrame -> DataFrame
D.forceDataFrame (DataFrame -> IO DataFrame) -> DataFrame -> IO DataFrame
forall a b. (a -> b) -> a -> b
$
                                [(Text, UExpr)] -> GroupedDataFrame -> DataFrame
Agg.aggregate [(Text, UExpr)]
mergeAggs ([Text] -> DataFrame -> GroupedDataFrame
Agg.groupBy [Text]
keys (DataFrame
a DataFrame -> DataFrame -> DataFrame
forall a. Semigroup a => a -> a -> a
<> DataFrame
partial))
                        Maybe DataFrame -> IO (Maybe DataFrame)
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return (DataFrame -> Maybe DataFrame
forall a. a -> Maybe a
Just DataFrame
merged)
                Maybe DataFrame -> IO (Maybe DataFrame)
loop Maybe DataFrame
next

-- | Merge a head accumulator with the rest of the workers' partials.
mergePartials ::
    [T.Text] ->
    [E.NamedExpr] ->
    D.DataFrame ->
    [D.DataFrame] ->
    IO D.DataFrame
mergePartials :: [Text]
-> [(Text, UExpr)] -> DataFrame -> [DataFrame] -> IO DataFrame
mergePartials [Text]
keys [(Text, UExpr)]
mergeAggs = DataFrame -> [DataFrame] -> IO DataFrame
go
  where
    go :: DataFrame -> [DataFrame] -> IO DataFrame
go !DataFrame
acc [] = DataFrame -> IO DataFrame
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return DataFrame
acc
    go !DataFrame
acc (DataFrame
p : [DataFrame]
ps) = do
        !DataFrame
merged <-
            DataFrame -> IO DataFrame
forall a. a -> IO a
evaluate (DataFrame -> IO DataFrame)
-> (DataFrame -> DataFrame) -> DataFrame -> IO DataFrame
forall b c a. (b -> c) -> (a -> b) -> a -> c
. DataFrame -> DataFrame
D.forceDataFrame (DataFrame -> IO DataFrame) -> DataFrame -> IO DataFrame
forall a b. (a -> b) -> a -> b
$
                [(Text, UExpr)] -> GroupedDataFrame -> DataFrame
Agg.aggregate [(Text, UExpr)]
mergeAggs ([Text] -> DataFrame -> GroupedDataFrame
Agg.groupBy [Text]
keys (DataFrame
acc DataFrame -> DataFrame -> DataFrame
forall a. Semigroup a => a -> a -> a
<> DataFrame
p))
        DataFrame -> [DataFrame] -> IO DataFrame
go DataFrame
merged [DataFrame]
ps

isStreamableAgg :: E.UExpr -> Bool
isStreamableAgg :: UExpr -> Bool
isStreamableAgg (E.UExpr (E.Agg (E.CollectAgg Text
_ v b -> a
_) Expr b
_)) = Bool
False
isStreamableAgg (E.UExpr (E.Agg (E.FoldAgg Text
_ Maybe a
Nothing (a -> b -> a
_ :: a -> b -> a)) Expr b
_)) =
    case TypeRep a -> TypeRep b -> Maybe (a :~: b)
forall a b. TypeRep a -> TypeRep b -> Maybe (a :~: b)
forall {k} (f :: k -> *) (a :: k) (b :: k).
TestEquality f =>
f a -> f b -> Maybe (a :~: b)
testEquality (forall a. Typeable a => TypeRep a
forall {k} (a :: k). Typeable a => TypeRep a
typeRep @a) (forall a. Typeable a => TypeRep a
forall {k} (a :: k). Typeable a => TypeRep a
typeRep @b) of
        Just a :~: b
Refl -> Bool
True
        Maybe (a :~: b)
Nothing -> Bool
False
isStreamableAgg (E.UExpr (E.Agg (E.FoldAgg Text
_ (Just a
_) (a -> b -> a
_ :: a -> b -> a)) Expr b
_)) =
    case TypeRep a -> TypeRep Int -> Maybe (a :~: Int)
forall a b. TypeRep a -> TypeRep b -> Maybe (a :~: b)
forall {k} (f :: k -> *) (a :: k) (b :: k).
TestEquality f =>
f a -> f b -> Maybe (a :~: b)
testEquality (forall a. Typeable a => TypeRep a
forall {k} (a :: k). Typeable a => TypeRep a
typeRep @a) (forall a. Typeable a => TypeRep a
forall {k} (a :: k). Typeable a => TypeRep a
typeRep @Int) of
        Just a :~: Int
Refl -> Bool
True
        Maybe (a :~: Int)
Nothing ->
            case TypeRep a -> TypeRep b -> Maybe (a :~: b)
forall a b. TypeRep a -> TypeRep b -> Maybe (a :~: b)
forall {k} (f :: k -> *) (a :: k) (b :: k).
TestEquality f =>
f a -> f b -> Maybe (a :~: b)
testEquality (forall a. Typeable a => TypeRep a
forall {k} (a :: k). Typeable a => TypeRep a
typeRep @a) (forall a. Typeable a => TypeRep a
forall {k} (a :: k). Typeable a => TypeRep a
typeRep @b) of
                Just a :~: b
Refl -> Bool
True
                Maybe (a :~: b)
Nothing -> Bool
False
isStreamableAgg (E.UExpr (E.Agg (E.MergeAgg{}) Expr b
_)) = Bool
True
isStreamableAgg UExpr
_ = Bool
False

{- | Build the (partial, merge, finalizer) plan for a list of streamable
aggregates: @partialAggs@ run per batch, @mergeAggs@ combine two partial
results, and @finalizer@ post-processes (for 'MergeAgg' acc≠output types).
-}
buildAggPlan ::
    [(T.Text, E.UExpr)] ->
    ( [(T.Text, E.UExpr)]
    , [(T.Text, E.UExpr)]
    , D.DataFrame -> D.DataFrame
    )
buildAggPlan :: [(Text, UExpr)]
-> ([(Text, UExpr)], [(Text, UExpr)], DataFrame -> DataFrame)
buildAggPlan [(Text, UExpr)]
aggs = (([(Text, UExpr)], [(Text, UExpr)], DataFrame -> DataFrame)
 -> ([(Text, UExpr)], [(Text, UExpr)], DataFrame -> DataFrame)
 -> ([(Text, UExpr)], [(Text, UExpr)], DataFrame -> DataFrame))
-> ([(Text, UExpr)], [(Text, UExpr)], DataFrame -> DataFrame)
-> [([(Text, UExpr)], [(Text, UExpr)], DataFrame -> DataFrame)]
-> ([(Text, UExpr)], [(Text, UExpr)], DataFrame -> DataFrame)
forall b a. (b -> a -> b) -> b -> [a] -> b
forall (t :: * -> *) b a.
Foldable t =>
(b -> a -> b) -> b -> t a -> b
foldl ([(Text, UExpr)], [(Text, UExpr)], DataFrame -> DataFrame)
-> ([(Text, UExpr)], [(Text, UExpr)], DataFrame -> DataFrame)
-> ([(Text, UExpr)], [(Text, UExpr)], DataFrame -> DataFrame)
forall {a} {a} {b} {c} {a}.
([a], [a], b -> c) -> ([a], [a], a -> b) -> ([a], [a], a -> c)
combine ([], [], DataFrame -> DataFrame
forall a. a -> a
id) (((Text, UExpr)
 -> ([(Text, UExpr)], [(Text, UExpr)], DataFrame -> DataFrame))
-> [(Text, UExpr)]
-> [([(Text, UExpr)], [(Text, UExpr)], DataFrame -> DataFrame)]
forall a b. (a -> b) -> [a] -> [b]
map (Text, UExpr)
-> ([(Text, UExpr)], [(Text, UExpr)], DataFrame -> DataFrame)
processAgg [(Text, UExpr)]
aggs)
  where
    combine :: ([a], [a], b -> c) -> ([a], [a], a -> b) -> ([a], [a], a -> c)
combine ([a]
p1, [a]
m1, b -> c
f1) ([a]
p2, [a]
m2, a -> b
f2) = ([a]
p1 [a] -> [a] -> [a]
forall a. [a] -> [a] -> [a]
++ [a]
p2, [a]
m1 [a] -> [a] -> [a]
forall a. [a] -> [a] -> [a]
++ [a]
m2, b -> c
f1 (b -> c) -> (a -> b) -> a -> c
forall b c a. (b -> c) -> (a -> b) -> a -> c
. a -> b
f2)

    processAgg ::
        (T.Text, E.UExpr) ->
        ([(T.Text, E.UExpr)], [(T.Text, E.UExpr)], D.DataFrame -> D.DataFrame)
    processAgg :: (Text, UExpr)
-> ([(Text, UExpr)], [(Text, UExpr)], DataFrame -> DataFrame)
processAgg (Text
name, UExpr
ue) = case UExpr
ue of
        E.UExpr (E.Agg (E.FoldAgg Text
n Maybe a
Nothing (a -> b -> a
f :: a -> b -> a)) (Expr b
_ :: E.Expr b)) ->
            case TypeRep a -> TypeRep b -> Maybe (a :~: b)
forall a b. TypeRep a -> TypeRep b -> Maybe (a :~: b)
forall {k} (f :: k -> *) (a :: k) (b :: k).
TestEquality f =>
f a -> f b -> Maybe (a :~: b)
testEquality (forall a. Typeable a => TypeRep a
forall {k} (a :: k). Typeable a => TypeRep a
typeRep @a) (forall a. Typeable a => TypeRep a
forall {k} (a :: k). Typeable a => TypeRep a
typeRep @b) of
                Just a :~: b
Refl ->
                    ( [(Text
name, UExpr
ue)]
                    , [(Text
name, Expr a -> UExpr
forall a. Columnable a => Expr a -> UExpr
E.UExpr (AggStrategy a b -> Expr b -> Expr a
forall a b.
(Columnable a, Columnable b) =>
AggStrategy a b -> Expr b -> Expr a
E.Agg (Text -> Maybe a -> (a -> b -> a) -> AggStrategy a b
forall a b. Text -> Maybe a -> (a -> b -> a) -> AggStrategy a b
E.FoldAgg Text
n Maybe a
forall a. Maybe a
Nothing a -> b -> a
f) (forall a. Columnable a => Text -> Expr a
E.Col @a Text
name)))]
                    , DataFrame -> DataFrame
forall a. a -> a
id
                    )
                Maybe (a :~: b)
Nothing ->
                    case TypeRep a -> TypeRep Int -> Maybe (a :~: Int)
forall a b. TypeRep a -> TypeRep b -> Maybe (a :~: b)
forall {k} (f :: k -> *) (a :: k) (b :: k).
TestEquality f =>
f a -> f b -> Maybe (a :~: b)
testEquality (forall a. Typeable a => TypeRep a
forall {k} (a :: k). Typeable a => TypeRep a
typeRep @a) (forall a. Typeable a => TypeRep a
forall {k} (a :: k). Typeable a => TypeRep a
typeRep @Int) of
                        Just a :~: Int
Refl ->
                            ( [(Text
name, UExpr
ue)]
                            ,
                                [
                                    ( Text
name
                                    , Expr Int -> UExpr
forall a. Columnable a => Expr a -> UExpr
E.UExpr
                                        (AggStrategy Int Int -> Expr Int -> Expr Int
forall a b.
(Columnable a, Columnable b) =>
AggStrategy a b -> Expr b -> Expr a
E.Agg (Text -> Maybe Int -> (Int -> Int -> Int) -> AggStrategy Int Int
forall a b. Text -> Maybe a -> (a -> b -> a) -> AggStrategy a b
E.FoldAgg Text
"sum" Maybe Int
forall a. Maybe a
Nothing (Int -> Int -> Int
forall a. Num a => a -> a -> a
(+) :: Int -> Int -> Int)) (forall a. Columnable a => Text -> Expr a
E.Col @Int Text
name))
                                    )
                                ]
                            , DataFrame -> DataFrame
forall a. a -> a
id
                            )
                        Maybe (a :~: Int)
Nothing -> ([(Text
name, UExpr
ue)], [(Text
name, UExpr
ue)], DataFrame -> DataFrame
forall a. a -> a
id)
        E.UExpr (E.Agg (E.FoldAgg Text
n (Just a
_) (a -> b -> a
f :: a -> b -> a)) (Expr b
_ :: E.Expr b)) ->
            case TypeRep a -> TypeRep Int -> Maybe (a :~: Int)
forall a b. TypeRep a -> TypeRep b -> Maybe (a :~: b)
forall {k} (f :: k -> *) (a :: k) (b :: k).
TestEquality f =>
f a -> f b -> Maybe (a :~: b)
testEquality (forall a. Typeable a => TypeRep a
forall {k} (a :: k). Typeable a => TypeRep a
typeRep @a) (forall a. Typeable a => TypeRep a
forall {k} (a :: k). Typeable a => TypeRep a
typeRep @Int) of
                Just a :~: Int
Refl ->
                    ( [(Text
name, UExpr
ue)]
                    ,
                        [
                            ( Text
name
                            , Expr Int -> UExpr
forall a. Columnable a => Expr a -> UExpr
E.UExpr
                                (AggStrategy Int Int -> Expr Int -> Expr Int
forall a b.
(Columnable a, Columnable b) =>
AggStrategy a b -> Expr b -> Expr a
E.Agg (Text -> Maybe Int -> (Int -> Int -> Int) -> AggStrategy Int Int
forall a b. Text -> Maybe a -> (a -> b -> a) -> AggStrategy a b
E.FoldAgg Text
"sum" Maybe Int
forall a. Maybe a
Nothing (Int -> Int -> Int
forall a. Num a => a -> a -> a
(+) :: Int -> Int -> Int)) (forall a. Columnable a => Text -> Expr a
E.Col @Int Text
name))
                            )
                        ]
                    , DataFrame -> DataFrame
forall a. a -> a
id
                    )
                Maybe (a :~: Int)
Nothing ->
                    case TypeRep a -> TypeRep b -> Maybe (a :~: b)
forall a b. TypeRep a -> TypeRep b -> Maybe (a :~: b)
forall {k} (f :: k -> *) (a :: k) (b :: k).
TestEquality f =>
f a -> f b -> Maybe (a :~: b)
testEquality (forall a. Typeable a => TypeRep a
forall {k} (a :: k). Typeable a => TypeRep a
typeRep @a) (forall a. Typeable a => TypeRep a
forall {k} (a :: k). Typeable a => TypeRep a
typeRep @b) of
                        Just a :~: b
Refl ->
                            ( [(Text
name, UExpr
ue)]
                            , [(Text
name, Expr a -> UExpr
forall a. Columnable a => Expr a -> UExpr
E.UExpr (AggStrategy a b -> Expr b -> Expr a
forall a b.
(Columnable a, Columnable b) =>
AggStrategy a b -> Expr b -> Expr a
E.Agg (Text -> Maybe a -> (a -> b -> a) -> AggStrategy a b
forall a b. Text -> Maybe a -> (a -> b -> a) -> AggStrategy a b
E.FoldAgg Text
n Maybe a
forall a. Maybe a
Nothing a -> b -> a
f) (forall a. Columnable a => Text -> Expr a
E.Col @a Text
name)))]
                            , DataFrame -> DataFrame
forall a. a -> a
id
                            )
                        Maybe (a :~: b)
Nothing -> ([(Text
name, UExpr
ue)], [(Text
name, UExpr
ue)], DataFrame -> DataFrame
forall a. a -> a
id)
        E.UExpr
            ( E.Agg
                    ( E.MergeAgg
                            Text
n
                            acc
seed
                            (acc -> b -> acc
step :: acc -> b -> acc)
                            (acc -> acc -> acc
merge :: acc -> acc -> acc)
                            (acc -> a
fin :: acc -> a)
                        )
                    (Expr b
inner :: E.Expr b)
                ) ->
                let partialExpr :: UExpr
partialExpr =
                        Expr acc -> UExpr
forall a. Columnable a => Expr a -> UExpr
E.UExpr
                            ( AggStrategy acc b -> Expr b -> Expr acc
forall a b.
(Columnable a, Columnable b) =>
AggStrategy a b -> Expr b -> Expr a
E.Agg
                                (Text
-> acc
-> (acc -> b -> acc)
-> (acc -> acc -> acc)
-> (acc -> acc)
-> AggStrategy acc b
forall acc b a.
Columnable acc =>
Text
-> acc
-> (acc -> b -> acc)
-> (acc -> acc -> acc)
-> (acc -> a)
-> AggStrategy a b
E.MergeAgg Text
n acc
seed acc -> b -> acc
step acc -> acc -> acc
merge (acc -> acc
forall a. a -> a
id :: acc -> acc))
                                Expr b
inner
                            )
                    mergeExpr :: UExpr
mergeExpr =
                        Expr acc -> UExpr
forall a. Columnable a => Expr a -> UExpr
E.UExpr
                            ( AggStrategy acc acc -> Expr acc -> Expr acc
forall a b.
(Columnable a, Columnable b) =>
AggStrategy a b -> Expr b -> Expr a
E.Agg
                                (Text -> Maybe acc -> (acc -> acc -> acc) -> AggStrategy acc acc
forall a b. Text -> Maybe a -> (a -> b -> a) -> AggStrategy a b
E.FoldAgg (Text
"merge_" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
n) Maybe acc
forall a. Maybe a
Nothing acc -> acc -> acc
merge)
                                (forall a. Columnable a => Text -> Expr a
E.Col @acc Text
name)
                            )
                    finalize :: DataFrame -> DataFrame
finalize DataFrame
df =
                        let accCol :: Column
accCol = Text -> DataFrame -> Column
D.unsafeGetColumn Text
name DataFrame
df
                            finalCol :: Column
finalCol =
                                (DataFrameException -> Column)
-> (Column -> Column) -> Either DataFrameException Column -> Column
forall a c b. (a -> c) -> (b -> c) -> Either a b -> c
either
                                    ([Char] -> DataFrameException -> Column
forall a. HasCallStack => [Char] -> a
error [Char]
"buildAggPlan: MergeAgg finalize failed")
                                    Column -> Column
forall a. a -> a
id
                                    (forall b c.
(Columnable b, Columnable c) =>
(b -> c) -> Column -> Either DataFrameException Column
C.mapColumn @acc @a acc -> a
fin Column
accCol)
                         in Text -> Column -> DataFrame -> DataFrame
D.insertColumn Text
name Column
finalCol DataFrame
df
                 in ( [(Text
name, UExpr
partialExpr)]
                    , [(Text
name, UExpr
mergeExpr)]
                    , DataFrame -> DataFrame
finalize
                    )
        UExpr
_ -> ([(Text
name, UExpr
ue)], [(Text
name, UExpr
ue)], DataFrame -> DataFrame
forall a. a -> a
id)

-- ---------------------------------------------------------------------------
-- Parquet scan implementation
-- ---------------------------------------------------------------------------

{- | Scan a Parquet file, directory, or glob.  Each file becomes one batch.
Column projection and predicate pushdown are forwarded to 'readParquetWithOpts'
via 'ParquetReadOptions'.
-}
executeParquetScan :: FilePath -> ScanConfig -> IO Stream
executeParquetScan :: [Char] -> ScanConfig -> IO Stream
executeParquetScan [Char]
path ScanConfig
cfg = do
    Bool
isDir <- [Char] -> IO Bool
doesDirectoryExist [Char]
path
    let pat :: [Char]
pat = if Bool
isDir then [Char]
path [Char] -> [Char] -> [Char]
</> [Char]
"*" else [Char]
path
    [[Char]]
matches <- [Char] -> IO [[Char]]
glob [Char]
pat
    [[Char]]
files <- ([Char] -> IO Bool) -> [[Char]] -> IO [[Char]]
forall (m :: * -> *) a.
Applicative m =>
(a -> m Bool) -> [a] -> m [a]
filterM ((Bool -> Bool) -> IO Bool -> IO Bool
forall a b. (a -> b) -> IO a -> IO b
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
fmap Bool -> Bool
not (IO Bool -> IO Bool) -> ([Char] -> IO Bool) -> [Char] -> IO Bool
forall b c a. (b -> c) -> (a -> b) -> a -> c
. [Char] -> IO Bool
doesDirectoryExist) [[Char]]
matches
    Bool -> IO () -> IO ()
forall (f :: * -> *). Applicative f => Bool -> f () -> f ()
when ([[Char]] -> Bool
forall a. [a] -> Bool
forall (t :: * -> *) a. Foldable t => t a -> Bool
null [[Char]]
files) (IO () -> IO ()) -> IO () -> IO ()
forall a b. (a -> b) -> a -> b
$
        [Char] -> IO ()
forall a. HasCallStack => [Char] -> a
error ([Char]
"executeParquetScan: no parquet files found for " [Char] -> [Char] -> [Char]
forall a. [a] -> [a] -> [a]
++ [Char]
path)
    let opts :: ParquetReadOptions
opts =
            ParquetReadOptions
Parquet.defaultParquetReadOptions
                { Parquet.selectedColumns = Just (M.keys (elements (scanSchema cfg)))
                , Parquet.predicate = scanPushdownPredicate cfg
                }
    IORef [[Char]]
ref <- [[Char]] -> IO (IORef [[Char]])
forall a. a -> IO (IORef a)
newIORef [[Char]]
files
    Stream -> IO Stream
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return (Stream -> IO Stream)
-> (IO (Maybe DataFrame) -> Stream)
-> IO (Maybe DataFrame)
-> IO Stream
forall b c a. (b -> c) -> (a -> b) -> a -> c
. IO (Maybe DataFrame) -> Stream
Stream (IO (Maybe DataFrame) -> IO Stream)
-> IO (Maybe DataFrame) -> IO Stream
forall a b. (a -> b) -> a -> b
$ do
        [[Char]]
fs <- IORef [[Char]] -> IO [[Char]]
forall a. IORef a -> IO a
readIORef IORef [[Char]]
ref
        case [[Char]]
fs of
            [] -> Maybe DataFrame -> IO (Maybe DataFrame)
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return Maybe DataFrame
forall a. Maybe a
Nothing
            ([Char]
f : [[Char]]
rest) -> do
                IORef [[Char]] -> [[Char]] -> IO ()
forall a. IORef a -> a -> IO ()
writeIORef IORef [[Char]]
ref [[Char]]
rest
                DataFrame -> Maybe DataFrame
forall a. a -> Maybe a
Just (DataFrame -> Maybe DataFrame)
-> IO DataFrame -> IO (Maybe DataFrame)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> ParquetReadOptions -> [Char] -> IO DataFrame
Parquet.readParquetWithOpts ParquetReadOptions
opts [Char]
f

-- ---------------------------------------------------------------------------
-- CSV scan implementation
-- ---------------------------------------------------------------------------

{- | The options a CSV scan reads with: the plan's schema supplies the column
types and the projection, the source supplies the separator.
-}
scanReadOptions :: Char -> ScanConfig -> ReadOptions
scanReadOptions :: Char -> ScanConfig -> ReadOptions
scanReadOptions Char
sep ScanConfig
cfg =
    (Schema -> ReadOptions
schemaReadOptions (ScanConfig -> Schema
scanSchema ScanConfig
cfg)){columnSeparator = sep}

{- | SIMD-parallel CSV scan: the file is split at newline boundaries into one
slice per capability, parsed concurrently, sliced into batches, and fed through
a bounded queue. Pushdown predicates are applied per batch by the consumer.
-}
executeCsvScan :: FilePath -> Char -> CsvReader -> ScanConfig -> IO Stream
executeCsvScan :: [Char] -> Char -> CsvReader -> ScanConfig -> IO Stream
executeCsvScan [Char]
path Char
sep CsvReader
reader ScanConfig
cfg = do
    Int
nCaps <- IO Int
getNumCapabilities
    [[Char]]
chunkPaths <- Int -> [Char] -> IO [[Char]]
splitCsvAtNewlines (Int -> Int -> Int
forall a. Ord a => a -> a -> a
max Int
1 Int
nCaps) [Char]
path

    let opts :: ReadOptions
opts = Char -> ScanConfig -> ReadOptions
scanReadOptions Char
sep ScanConfig
cfg
        batchSz :: Int
batchSz = ScanConfig -> Int
scanBatchSize ScanConfig
cfg
    [DataFrame]
chunkDfs <- ([Char] -> IO DataFrame) -> [[Char]] -> IO [DataFrame]
forall (t :: * -> *) a b.
Traversable t =>
(a -> IO b) -> t a -> IO (t b)
mapConcurrently (CsvReader
reader ReadOptions
opts) [[Char]]
chunkPaths
    ([Char] -> IO ()) -> [[Char]] -> IO ()
forall (t :: * -> *) (m :: * -> *) a b.
(Foldable t, Monad m) =>
(a -> m b) -> t a -> m ()
mapM_ [Char] -> IO ()
removeFile [[Char]]
chunkPaths

    TBQueue (Maybe DataFrame)
queue <- Natural -> IO (TBQueue (Maybe DataFrame))
forall a. Natural -> IO (TBQueue a)
newTBQueueIO (Int -> Natural
forall a b. (Integral a, Num b) => a -> b
fromIntegral (Int -> Int -> Int
forall a. Ord a => a -> a -> a
max Int
4 (Int
2 Int -> Int -> Int
forall a. Num a => a -> a -> a
* Int
nCaps)))
    ThreadId
_ <- IO () -> IO ThreadId
forkIO (IO () -> IO ThreadId) -> IO () -> IO ThreadId
forall a b. (a -> b) -> a -> b
$ do
        [DataFrame] -> (DataFrame -> IO ()) -> IO ()
forall (t :: * -> *) (m :: * -> *) a b.
(Foldable t, Monad m) =>
t a -> (a -> m b) -> m ()
forM_ [DataFrame]
chunkDfs ((DataFrame -> IO ()) -> IO ()) -> (DataFrame -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \DataFrame
df ->
            [DataFrame] -> (DataFrame -> IO ()) -> IO ()
forall (t :: * -> *) (m :: * -> *) a b.
(Foldable t, Monad m) =>
t a -> (a -> m b) -> m ()
forM_ (Int -> DataFrame -> [DataFrame]
sliceIntoBatches Int
batchSz DataFrame
df) ((DataFrame -> IO ()) -> IO ()) -> (DataFrame -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \DataFrame
b ->
                STM () -> IO ()
forall a. STM a -> IO a
atomically (TBQueue (Maybe DataFrame) -> Maybe DataFrame -> STM ()
forall a. TBQueue a -> a -> STM ()
writeTBQueue TBQueue (Maybe DataFrame)
queue (DataFrame -> Maybe DataFrame
forall a. a -> Maybe a
Just DataFrame
b))
        STM () -> IO ()
forall a. STM a -> IO a
atomically (TBQueue (Maybe DataFrame) -> Maybe DataFrame -> STM ()
forall a. TBQueue a -> a -> STM ()
writeTBQueue TBQueue (Maybe DataFrame)
queue Maybe DataFrame
forall a. Maybe a
Nothing)
    Stream -> IO Stream
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return (Stream -> IO Stream)
-> (IO (Maybe DataFrame) -> Stream)
-> IO (Maybe DataFrame)
-> IO Stream
forall b c a. (b -> c) -> (a -> b) -> a -> c
. IO (Maybe DataFrame) -> Stream
Stream (IO (Maybe DataFrame) -> IO Stream)
-> IO (Maybe DataFrame) -> IO Stream
forall a b. (a -> b) -> a -> b
$
        ( do
            Maybe DataFrame
mb <- STM (Maybe DataFrame) -> IO (Maybe DataFrame)
forall a. STM a -> IO a
atomically (TBQueue (Maybe DataFrame) -> STM (Maybe DataFrame)
forall a. TBQueue a -> STM a
readTBQueue TBQueue (Maybe DataFrame)
queue)
            case Maybe DataFrame
mb of
                Maybe DataFrame
Nothing -> STM () -> IO ()
forall a. STM a -> IO a
atomically (TBQueue (Maybe DataFrame) -> Maybe DataFrame -> STM ()
forall a. TBQueue a -> a -> STM ()
writeTBQueue TBQueue (Maybe DataFrame)
queue Maybe DataFrame
forall a. Maybe a
Nothing) IO () -> IO (Maybe DataFrame) -> IO (Maybe DataFrame)
forall a b. IO a -> IO b -> IO b
forall (m :: * -> *) a b. Monad m => m a -> m b -> m b
>> Maybe DataFrame -> IO (Maybe DataFrame)
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return Maybe DataFrame
forall a. Maybe a
Nothing
                Just DataFrame
df ->
                    let df' :: DataFrame
df' = case ScanConfig -> Maybe (Expr Bool)
scanPushdownPredicate ScanConfig
cfg of
                            Maybe (Expr Bool)
Nothing -> DataFrame
df
                            Just Expr Bool
p -> Expr Bool -> DataFrame -> DataFrame
Sub.filterWhere Expr Bool
p DataFrame
df
                     in Maybe DataFrame -> IO (Maybe DataFrame)
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return (DataFrame -> Maybe DataFrame
forall a. a -> Maybe a
Just DataFrame
df')
        )

executeCsvScanStreaming ::
    FilePath -> Char -> CsvReader -> ScanConfig -> IO Stream
executeCsvScanStreaming :: [Char] -> Char -> CsvReader -> ScanConfig -> IO Stream
executeCsvScanStreaming [Char]
path Char
sep CsvReader
reader ScanConfig
cfg = do
    let opts :: ReadOptions
opts = Char -> ScanConfig -> ReadOptions
scanReadOptions Char
sep ScanConfig
cfg
        batchSz :: Int
batchSz = ScanConfig -> Int
scanBatchSize ScanConfig
cfg
        windowBytes :: Int
windowBytes = Int
64 Int -> Int -> Int
forall a. Num a => a -> a -> a
* Int
1024 Int -> Int -> Int
forall a. Num a => a -> a -> a
* Int
1024 :: Int
    TBQueue (Maybe DataFrame)
queue <- Natural -> IO (TBQueue (Maybe DataFrame))
forall a. Natural -> IO (TBQueue a)
newTBQueueIO Natural
8
    ThreadId
_ <- IO () -> IO ThreadId
forkIO (IO () -> IO ThreadId) -> IO () -> IO ThreadId
forall a b. (a -> b) -> a -> b
$
        [Char] -> IOMode -> (Handle -> IO ()) -> IO ()
forall r. [Char] -> IOMode -> (Handle -> IO r) -> IO r
withFile [Char]
path IOMode
ReadMode ((Handle -> IO ()) -> IO ()) -> (Handle -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \Handle
h -> do
            ByteString
header <- Handle -> IO ByteString
C8.hGetLine Handle
h
            let feed :: ByteString -> IO ()
feed ByteString
bytes =
                    Bool -> IO () -> IO ()
forall (f :: * -> *). Applicative f => Bool -> f () -> f ()
unless (ByteString -> Bool
BS.null ByteString
bytes) (IO () -> IO ()) -> IO () -> IO ()
forall a b. (a -> b) -> a -> b
$ do
                        [Char]
p <- [Char] -> IO [Char]
emptySystemTempFile [Char]
"lazy_csv_win_.csv"
                        [Char] -> ByteString -> IO ()
BS.writeFile [Char]
p (ByteString
header ByteString -> ByteString -> ByteString
forall a. Semigroup a => a -> a -> a
<> Word8 -> ByteString
BS.singleton Word8
nl ByteString -> ByteString -> ByteString
forall a. Semigroup a => a -> a -> a
<> ByteString
bytes)
                        DataFrame
df <- CsvReader
reader ReadOptions
opts [Char]
p
                        [Char] -> IO ()
removeFile [Char]
p
                        [DataFrame] -> (DataFrame -> IO ()) -> IO ()
forall (t :: * -> *) (m :: * -> *) a b.
(Foldable t, Monad m) =>
t a -> (a -> m b) -> m ()
forM_ (Int -> DataFrame -> [DataFrame]
sliceIntoBatches Int
batchSz DataFrame
df) ((DataFrame -> IO ()) -> IO ()) -> (DataFrame -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \DataFrame
b ->
                            STM () -> IO ()
forall a. STM a -> IO a
atomically (TBQueue (Maybe DataFrame) -> Maybe DataFrame -> STM ()
forall a. TBQueue a -> a -> STM ()
writeTBQueue TBQueue (Maybe DataFrame)
queue (DataFrame -> Maybe DataFrame
forall a. a -> Maybe a
Just DataFrame
b))
                loop :: ByteString -> IO ()
loop ByteString
leftover = do
                    Bool
eof <- Handle -> IO Bool
hIsEOF Handle
h
                    if Bool
eof
                        then ByteString -> IO ()
feed ByteString
leftover IO () -> IO () -> IO ()
forall a b. IO a -> IO b -> IO b
forall (m :: * -> *) a b. Monad m => m a -> m b -> m b
>> STM () -> IO ()
forall a. STM a -> IO a
atomically (TBQueue (Maybe DataFrame) -> Maybe DataFrame -> STM ()
forall a. TBQueue a -> a -> STM ()
writeTBQueue TBQueue (Maybe DataFrame)
queue Maybe DataFrame
forall a. Maybe a
Nothing)
                        else do
                            ByteString
chunk <- Handle -> Int -> IO ByteString
BS.hGetSome Handle
h Int
windowBytes
                            let buf :: ByteString
buf = ByteString
leftover ByteString -> ByteString -> ByteString
forall a. Semigroup a => a -> a -> a
<> ByteString
chunk
                            case Word8 -> ByteString -> Maybe Int
BS.elemIndexEnd Word8
nl ByteString
buf of
                                Maybe Int
Nothing -> ByteString -> IO ()
loop ByteString
buf
                                Just Int
i -> ByteString -> IO ()
feed (Int -> ByteString -> ByteString
BS.take Int
i ByteString
buf) IO () -> IO () -> IO ()
forall a b. IO a -> IO b -> IO b
forall (m :: * -> *) a b. Monad m => m a -> m b -> m b
>> ByteString -> IO ()
loop (Int -> ByteString -> ByteString
BS.drop (Int
i Int -> Int -> Int
forall a. Num a => a -> a -> a
+ Int
1) ByteString
buf)
            ByteString -> IO ()
loop ByteString
BS.empty
    Stream -> IO Stream
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return (Stream -> IO Stream)
-> (IO (Maybe DataFrame) -> Stream)
-> IO (Maybe DataFrame)
-> IO Stream
forall b c a. (b -> c) -> (a -> b) -> a -> c
. IO (Maybe DataFrame) -> Stream
Stream (IO (Maybe DataFrame) -> IO Stream)
-> IO (Maybe DataFrame) -> IO Stream
forall a b. (a -> b) -> a -> b
$ do
        Maybe DataFrame
mb <- STM (Maybe DataFrame) -> IO (Maybe DataFrame)
forall a. STM a -> IO a
atomically (TBQueue (Maybe DataFrame) -> STM (Maybe DataFrame)
forall a. TBQueue a -> STM a
readTBQueue TBQueue (Maybe DataFrame)
queue)
        case Maybe DataFrame
mb of
            Maybe DataFrame
Nothing -> STM () -> IO ()
forall a. STM a -> IO a
atomically (TBQueue (Maybe DataFrame) -> Maybe DataFrame -> STM ()
forall a. TBQueue a -> a -> STM ()
writeTBQueue TBQueue (Maybe DataFrame)
queue Maybe DataFrame
forall a. Maybe a
Nothing) IO () -> IO (Maybe DataFrame) -> IO (Maybe DataFrame)
forall a b. IO a -> IO b -> IO b
forall (m :: * -> *) a b. Monad m => m a -> m b -> m b
>> Maybe DataFrame -> IO (Maybe DataFrame)
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return Maybe DataFrame
forall a. Maybe a
Nothing
            Just DataFrame
df ->
                let df' :: DataFrame
df' = case ScanConfig -> Maybe (Expr Bool)
scanPushdownPredicate ScanConfig
cfg of
                        Maybe (Expr Bool)
Nothing -> DataFrame
df
                        Just Expr Bool
p -> Expr Bool -> DataFrame -> DataFrame
Sub.filterWhere Expr Bool
p DataFrame
df
                 in Maybe DataFrame -> IO (Maybe DataFrame)
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return (DataFrame -> Maybe DataFrame
forall a. a -> Maybe a
Just DataFrame
df')
  where
    nl :: Word8
    nl :: Word8
nl = Word8
0x0A

-- | Slice a 'DataFrame' into row-bounded batches of at most @n@ rows.
sliceIntoBatches :: Int -> D.DataFrame -> [D.DataFrame]
sliceIntoBatches :: Int -> DataFrame -> [DataFrame]
sliceIntoBatches Int
n DataFrame
df =
    let total :: Int
total = DataFrame -> Int
Core.nRows DataFrame
df
        starts :: [Int]
starts = [Int
0, Int
n .. Int
total Int -> Int -> Int
forall a. Num a => a -> a -> a
- Int
1]
     in [(Int, Int) -> DataFrame -> DataFrame
Sub.range (Int
s, Int -> Int -> Int
forall a. Ord a => a -> a -> a
min (Int
s Int -> Int -> Int
forall a. Num a => a -> a -> a
+ Int
n) Int
total) DataFrame
df | Int
s <- [Int]
starts]

{- | Split a CSV file at newline boundaries into @n@ temp files, each with the
original header plus a body slice. Returns the paths (caller removes them); the
per-file reader mmap's each, giving OS-paged reads instead of one monolithic read.
-}
splitCsvAtNewlines :: Int -> FilePath -> IO [FilePath]
splitCsvAtNewlines :: Int -> [Char] -> IO [[Char]]
splitCsvAtNewlines Int
n [Char]
path = do
    ByteString
bs <- [Char] -> IO ByteString
BS.readFile [Char]
path
    let (ByteString
header, ByteString
rest) = (Word8 -> Bool) -> ByteString -> (ByteString, ByteString)
BS.break (Word8 -> Word8 -> Bool
forall a. Eq a => a -> a -> Bool
== Word8
nl) ByteString
bs
        body :: ByteString
body = Int -> ByteString -> ByteString
BS.drop Int
1 ByteString
rest
        bodyLen :: Int
bodyLen = ByteString -> Int
BS.length ByteString
body
        rawOffsets :: [Int]
rawOffsets = [(Int
bodyLen Int -> Int -> Int
forall a. Num a => a -> a -> a
* Int
i) Int -> Int -> Int
forall a. Integral a => a -> a -> a
`div` Int
n | Int
i <- [Int
0 .. Int
n]]
        snapped :: [Int]
snapped = Int
0 Int -> [Int] -> [Int]
forall a. a -> [a] -> [a]
: (Int -> Int) -> [Int] -> [Int]
forall a b. (a -> b) -> [a] -> [b]
map (ByteString -> Int -> Int
snap ByteString
body) ([Int] -> [Int]
forall a. HasCallStack => [a] -> [a]
init (Int -> [Int] -> [Int]
forall a. Int -> [a] -> [a]
drop Int
1 [Int]
rawOffsets)) [Int] -> [Int] -> [Int]
forall a. [a] -> [a] -> [a]
++ [Int
bodyLen]
        ranges :: [(Int, Int)]
ranges = [Int] -> [Int] -> [(Int, Int)]
forall a b. [a] -> [b] -> [(a, b)]
zip [Int]
snapped (Int -> [Int] -> [Int]
forall a. Int -> [a] -> [a]
drop Int
1 [Int]
snapped)
        slices :: [ByteString]
slices =
            [ Int -> ByteString -> ByteString
BS.take (Int
hi Int -> Int -> Int
forall a. Num a => a -> a -> a
- Int
lo) (Int -> ByteString -> ByteString
BS.drop Int
lo ByteString
body)
            | (Int
lo, Int
hi) <- [(Int, Int)]
ranges
            , Int
hi Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
> Int
lo
            ]
    [ByteString] -> (ByteString -> IO [Char]) -> IO [[Char]]
forall (t :: * -> *) (m :: * -> *) a b.
(Traversable t, Monad m) =>
t a -> (a -> m b) -> m (t b)
forM [ByteString]
slices ((ByteString -> IO [Char]) -> IO [[Char]])
-> (ByteString -> IO [Char]) -> IO [[Char]]
forall a b. (a -> b) -> a -> b
$ \ByteString
chunk -> do
        [Char]
p <- [Char] -> IO [Char]
emptySystemTempFile [Char]
"lazy_csv_chunk_.csv"
        [Char] -> ByteString -> IO ()
BS.writeFile [Char]
p (ByteString
header ByteString -> ByteString -> ByteString
forall a. Semigroup a => a -> a -> a
<> Word8 -> ByteString
BS.singleton Word8
nl ByteString -> ByteString -> ByteString
forall a. Semigroup a => a -> a -> a
<> ByteString
chunk)
        [Char] -> IO [Char]
forall a. a -> IO a
forall (m :: * -> *) a. Monad m => a -> m a
return [Char]
p
  where
    nl :: Word8
    nl :: Word8
nl = Word8
0x0A
    snap :: ByteString -> Int -> Int
snap ByteString
body Int
off =
        case Word8 -> ByteString -> Maybe Int
BS.elemIndex Word8
nl (Int -> ByteString -> ByteString
BS.drop Int
off ByteString
body) of
            Just Int
i -> Int
off Int -> Int -> Int
forall a. Num a => a -> a -> a
+ Int
i Int -> Int -> Int
forall a. Num a => a -> a -> a
+ Int
1
            Maybe Int
Nothing -> ByteString -> Int
BS.length ByteString
body

-- ---------------------------------------------------------------------------
-- Join helper
-- ---------------------------------------------------------------------------

{- | Route a join to 'Operations.Join', renaming the right key when the names
differ. 'Join.join' keeps its first argument, so LEFT/RIGHT must pass 'leftDf'
first to retain the intended side; symmetric INNER/FULL_OUTER pass @rightDf@ first.
-}
performJoin ::
    Join.JoinType -> T.Text -> T.Text -> D.DataFrame -> D.DataFrame -> D.DataFrame
performJoin :: JoinType -> Text -> Text -> DataFrame -> DataFrame -> DataFrame
performJoin JoinType
jt Text
leftKey Text
rightKey DataFrame
leftDf DataFrame
rightDf =
    case JoinType
jt of
        JoinType
Join.LEFT -> JoinType -> [Text] -> DataFrame -> DataFrame -> DataFrame
Join.join JoinType
jt [Text
leftKey] DataFrame
leftDf DataFrame
rightRenamed
        JoinType
Join.RIGHT -> JoinType -> [Text] -> DataFrame -> DataFrame -> DataFrame
Join.join JoinType
jt [Text
leftKey] DataFrame
leftDf DataFrame
rightRenamed
        JoinType
_ -> JoinType -> [Text] -> DataFrame -> DataFrame -> DataFrame
Join.join JoinType
jt [Text
leftKey] DataFrame
rightRenamed DataFrame
leftDf
  where
    rightRenamed :: DataFrame
rightRenamed
        | Text
leftKey Text -> Text -> Bool
forall a. Eq a => a -> a -> Bool
== Text
rightKey = DataFrame
rightDf
        | Bool
otherwise = Text -> Text -> DataFrame -> DataFrame
Core.rename Text
rightKey Text
leftKey DataFrame
rightDf

-- ---------------------------------------------------------------------------
-- Sort order conversion
-- ---------------------------------------------------------------------------

{- | Convert a plan-level @(column, direction)@ into a Permutation 'SortOrder',
emitting @E.Col@ at the materialised column's element type so 'Perm.sortBy'
dispatches its comparator correctly; unknown/non-'Ord' types fall back to 'T.Text'.
-}
toPermSortOrder :: D.DataFrame -> (T.Text, SortOrder) -> Perm.SortOrder
toPermSortOrder :: DataFrame -> (Text, SortOrder) -> SortOrder
toPermSortOrder DataFrame
df (Text
col, SortOrder
dir) =
    case Text -> Map Text Int -> Maybe Int
forall k a. Ord k => k -> Map k a -> Maybe a
M.lookup Text
col (DataFrame -> Map Text Int
D.columnIndices DataFrame
df) of
        Maybe Int
Nothing -> forall a. (Columnable a, Ord a) => SortOrder
mk @T.Text
        Just Int
idx -> Column -> SortOrder
dispatch (DataFrame -> Vector Column
D.columns DataFrame
df Vector Column -> Int -> Column
forall a. Vector a -> Int -> a
VB.! Int
idx)
  where
    mk :: forall a. (C.Columnable a, Ord a) => Perm.SortOrder
    mk :: forall a. (Columnable a, Ord a) => SortOrder
mk = case SortOrder
dir of
        SortOrder
Ascending -> Expr a -> SortOrder
forall a. (Columnable a, Ord a) => Expr a -> SortOrder
Perm.Asc (forall a. Columnable a => Text -> Expr a
E.Col @a Text
col)
        SortOrder
Descending -> Expr a -> SortOrder
forall a. (Columnable a, Ord a) => Expr a -> SortOrder
Perm.Desc (forall a. Columnable a => Text -> Expr a
E.Col @a Text
col)

    dispatch :: C.Column -> Perm.SortOrder
    dispatch :: Column -> SortOrder
dispatch Column
column = case Column
column of
        c :: Column
c@(C.MergedColumn Column
_ Column
_) -> Column -> SortOrder
dispatch (Column -> Column
C.materializeMerged Column
c)
        C.PackedText{} -> forall a. (Columnable a, Ord a) => SortOrder
mk @T.Text
        C.BoxedColumn Maybe Bitmap
_ (Vector a
_ :: VB.Vector b) -> forall b. Typeable b => SortOrder
pick @b
        C.UnboxedColumn Maybe Bitmap
_ (Vector a
_ :: VU.Vector b) -> forall b. Typeable b => SortOrder
pick @b

    pick :: forall b. (Typeable b) => Perm.SortOrder
    pick :: forall b. Typeable b => SortOrder
pick =
        forall b a.
(Typeable b, Columnable a, Ord a) =>
SortOrder -> SortOrder
tryT @b @Int (SortOrder -> SortOrder) -> SortOrder -> SortOrder
forall a b. (a -> b) -> a -> b
$
            forall b a.
(Typeable b, Columnable a, Ord a) =>
SortOrder -> SortOrder
tryT @b @Double (SortOrder -> SortOrder) -> SortOrder -> SortOrder
forall a b. (a -> b) -> a -> b
$
                forall b a.
(Typeable b, Columnable a, Ord a) =>
SortOrder -> SortOrder
tryT @b @Float (SortOrder -> SortOrder) -> SortOrder -> SortOrder
forall a b. (a -> b) -> a -> b
$
                    forall b a.
(Typeable b, Columnable a, Ord a) =>
SortOrder -> SortOrder
tryT @b @Integer (SortOrder -> SortOrder) -> SortOrder -> SortOrder
forall a b. (a -> b) -> a -> b
$
                        forall b a.
(Typeable b, Columnable a, Ord a) =>
SortOrder -> SortOrder
tryT @b @Int8 (SortOrder -> SortOrder) -> SortOrder -> SortOrder
forall a b. (a -> b) -> a -> b
$
                            forall b a.
(Typeable b, Columnable a, Ord a) =>
SortOrder -> SortOrder
tryT @b @Int16 (SortOrder -> SortOrder) -> SortOrder -> SortOrder
forall a b. (a -> b) -> a -> b
$
                                forall b a.
(Typeable b, Columnable a, Ord a) =>
SortOrder -> SortOrder
tryT @b @Int32 (SortOrder -> SortOrder) -> SortOrder -> SortOrder
forall a b. (a -> b) -> a -> b
$
                                    forall b a.
(Typeable b, Columnable a, Ord a) =>
SortOrder -> SortOrder
tryT @b @Int64 (SortOrder -> SortOrder) -> SortOrder -> SortOrder
forall a b. (a -> b) -> a -> b
$
                                        forall b a.
(Typeable b, Columnable a, Ord a) =>
SortOrder -> SortOrder
tryT @b @Word8 (SortOrder -> SortOrder) -> SortOrder -> SortOrder
forall a b. (a -> b) -> a -> b
$
                                            forall b a.
(Typeable b, Columnable a, Ord a) =>
SortOrder -> SortOrder
tryT @b @Word16 (SortOrder -> SortOrder) -> SortOrder -> SortOrder
forall a b. (a -> b) -> a -> b
$
                                                forall b a.
(Typeable b, Columnable a, Ord a) =>
SortOrder -> SortOrder
tryT @b @Word32 (SortOrder -> SortOrder) -> SortOrder -> SortOrder
forall a b. (a -> b) -> a -> b
$
                                                    forall b a.
(Typeable b, Columnable a, Ord a) =>
SortOrder -> SortOrder
tryT @b @Word64 (SortOrder -> SortOrder) -> SortOrder -> SortOrder
forall a b. (a -> b) -> a -> b
$
                                                        forall b a.
(Typeable b, Columnable a, Ord a) =>
SortOrder -> SortOrder
tryT @b @Bool (SortOrder -> SortOrder) -> SortOrder -> SortOrder
forall a b. (a -> b) -> a -> b
$
                                                            forall b a.
(Typeable b, Columnable a, Ord a) =>
SortOrder -> SortOrder
tryT @b @Char (SortOrder -> SortOrder) -> SortOrder -> SortOrder
forall a b. (a -> b) -> a -> b
$
                                                                forall b a.
(Typeable b, Columnable a, Ord a) =>
SortOrder -> SortOrder
tryT @b @T.Text (forall a. (Columnable a, Ord a) => SortOrder
mk @T.Text)

    tryT ::
        forall b a.
        (Typeable b, C.Columnable a, Ord a) =>
        Perm.SortOrder ->
        Perm.SortOrder
    tryT :: forall b a.
(Typeable b, Columnable a, Ord a) =>
SortOrder -> SortOrder
tryT SortOrder
fallback = case TypeRep b -> TypeRep a -> Maybe (b :~: a)
forall a b. TypeRep a -> TypeRep b -> Maybe (a :~: b)
forall {k} (f :: k -> *) (a :: k) (b :: k).
TestEquality f =>
f a -> f b -> Maybe (a :~: b)
testEquality (forall a. Typeable a => TypeRep a
forall {k} (a :: k). Typeable a => TypeRep a
typeRep @b) (forall a. Typeable a => TypeRep a
forall {k} (a :: k). Typeable a => TypeRep a
typeRep @a) of
        Just b :~: a
Refl -> forall a. (Columnable a, Ord a) => SortOrder
mk @a
        Maybe (b :~: a)
Nothing -> SortOrder
fallback