{-# LANGUAGE AllowAmbiguousTypes #-}
{-# LANGUAGE BangPatterns #-}
{-# LANGUAGE ExplicitNamespaces #-}
{-# LANGUAGE FlexibleContexts #-}
{-# LANGUAGE GADTs #-}
{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE ScopedTypeVariables #-}
{-# LANGUAGE TupleSections #-}
{-# LANGUAGE TypeApplications #-}
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)
newtype Stream = Stream {Stream -> IO (Maybe DataFrame)
pullBatch :: IO (Maybe D.DataFrame)}
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
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)
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
]
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
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
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
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)
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
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
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)
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
scanReadOptions :: Char -> ScanConfig -> ReadOptions
scanReadOptions :: Char -> ScanConfig -> ReadOptions
scanReadOptions Char
sep ScanConfig
cfg =
(Schema -> ReadOptions
schemaReadOptions (ScanConfig -> Schema
scanSchema ScanConfig
cfg)){columnSeparator = sep}
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
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]
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
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
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