{-# LANGUAGE CPP #-}
{-# LANGUAGE FlexibleContexts #-}
{-# LANGUAGE MonoLocalBinds #-}
{-# LANGUAGE NumericUnderscores #-}
{-# LANGUAGE OverloadedRecordDot #-}
{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE ScopedTypeVariables #-}
{-# LANGUAGE TypeApplications #-}
module DataFrame.IO.Parquet where
import Control.Exception (throw)
import Control.Monad
import Control.Monad.IO.Class (MonadIO (..))
import Control.Monad.ST (stToIO)
import Data.Bits (Bits (shiftL), (.|.))
import qualified Data.ByteString as BS
import Data.Either (fromRight)
import Data.Functor ((<&>))
import Data.Int (Int32, Int64)
#if !MIN_VERSION_base(4,20,0)
import Data.List (foldl', transpose)
#else
import Data.List (transpose)
#endif
import qualified Data.List as L
import qualified Data.Map as Map
import qualified Data.Text as T
import Data.Time.Calendar (Day (ModifiedJulianDay))
import Data.Time.Clock (UTCTime (UTCTime), picosecondsToDiffTime)
import qualified Data.Vector as Vector
import qualified Data.Vector.Unboxed as VU
import DataFrame.Errors (DataFrameException (ColumnsNotFoundException))
import DataFrame.IO.Parquet.Page (
PageDecoder,
UnboxedPageDecoder,
appendNullableStringPageIO,
appendStringPageIO,
boolDecoder,
byteArrayDecoder,
doubleDecoder,
fixedLenByteArrayDecoder,
floatDecoder,
foldColumnDataPagesM,
foldColumnPagesM,
int32Decoder,
int64Decoder,
int96Decoder,
)
import DataFrame.IO.Parquet.Seeking (
FileBufferedOrSeekable,
ForceNonSeekable,
withFileBufferedOrSeekable,
)
import DataFrame.IO.Parquet.Thrift (
ColumnChunk (..),
DecimalType (..),
FileMetadata (..),
LogicalType (..),
RowGroup (..),
ThriftType (..),
TimeUnit (..),
TimestampType (..),
unField,
)
import DataFrame.IO.Parquet.Utils (
ColumnDescription (..),
foldNonNullable,
foldNonNullableUnboxed,
foldNullable,
foldNullableUnboxed,
foldRepeated,
foldRepeatedUnboxed,
generateColumnDescriptions,
getColumnNames,
)
import DataFrame.IO.Utils.RandomAccess (
RandomAccess (..),
ReaderIO (runReaderIO),
)
import DataFrame.Internal.Column (Column, Columnable)
import qualified DataFrame.Internal.Column as DI
import DataFrame.Internal.ColumnBuilder (
freezeTextChunk,
mergeTextChunks,
newTextBuilder,
)
import DataFrame.Internal.DataFrame (DataFrame (..))
import DataFrame.Internal.Expression (Expr, getColumns)
import DataFrame.Operations.Merge ()
import qualified DataFrame.Operations.Subset as DS
import qualified Pinch
import System.Directory (doesDirectoryExist)
import System.FilePath ((</>))
import System.FilePath.Glob (glob)
import System.IO (IOMode (ReadMode))
data ParquetReadOptions = ParquetReadOptions
{ ParquetReadOptions -> Maybe [Text]
selectedColumns :: Maybe [T.Text]
, ParquetReadOptions -> Maybe (Expr Bool)
predicate :: Maybe (Expr Bool)
, ParquetReadOptions -> Maybe (Int, Int)
rowRange :: Maybe (Int, Int)
, ParquetReadOptions -> Bool
safeColumns :: Bool
}
deriving (Int -> ParquetReadOptions -> ShowS
[ParquetReadOptions] -> ShowS
ParquetReadOptions -> [Char]
(Int -> ParquetReadOptions -> ShowS)
-> (ParquetReadOptions -> [Char])
-> ([ParquetReadOptions] -> ShowS)
-> Show ParquetReadOptions
forall a.
(Int -> a -> ShowS) -> (a -> [Char]) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> ParquetReadOptions -> ShowS
showsPrec :: Int -> ParquetReadOptions -> ShowS
$cshow :: ParquetReadOptions -> [Char]
show :: ParquetReadOptions -> [Char]
$cshowList :: [ParquetReadOptions] -> ShowS
showList :: [ParquetReadOptions] -> ShowS
Show)
defaultParquetReadOptions :: ParquetReadOptions
defaultParquetReadOptions :: ParquetReadOptions
defaultParquetReadOptions =
ParquetReadOptions
{ selectedColumns :: Maybe [Text]
selectedColumns = Maybe [Text]
forall a. Maybe a
Nothing
, predicate :: Maybe (Expr Bool)
predicate = Maybe (Expr Bool)
forall a. Maybe a
Nothing
, rowRange :: Maybe (Int, Int)
rowRange = Maybe (Int, Int)
forall a. Maybe a
Nothing
, safeColumns :: Bool
safeColumns = Bool
False
}
readParquet :: FilePath -> IO DataFrame
readParquet :: [Char] -> IO DataFrame
readParquet = ParquetReadOptions -> [Char] -> IO DataFrame
readParquetWithOpts ParquetReadOptions
defaultParquetReadOptions
readParquetWithOpts :: ParquetReadOptions -> FilePath -> IO DataFrame
readParquetWithOpts :: ParquetReadOptions -> [Char] -> IO DataFrame
readParquetWithOpts = ForceNonSeekable -> ParquetReadOptions -> [Char] -> IO DataFrame
_readParquetWithOpts ForceNonSeekable
forall a. Maybe a
Nothing
_readParquetWithOpts ::
ForceNonSeekable -> ParquetReadOptions -> FilePath -> IO DataFrame
_readParquetWithOpts :: ForceNonSeekable -> ParquetReadOptions -> [Char] -> IO DataFrame
_readParquetWithOpts ForceNonSeekable
extraConfig ParquetReadOptions
opts [Char]
path =
ForceNonSeekable
-> [Char]
-> IOMode
-> (FileBufferedOrSeekable -> IO DataFrame)
-> IO DataFrame
forall a.
ForceNonSeekable
-> [Char] -> IOMode -> (FileBufferedOrSeekable -> IO a) -> IO a
withFileBufferedOrSeekable ForceNonSeekable
extraConfig [Char]
path IOMode
ReadMode ((FileBufferedOrSeekable -> IO DataFrame) -> IO DataFrame)
-> (FileBufferedOrSeekable -> IO DataFrame) -> IO DataFrame
forall a b. (a -> b) -> a -> b
$ \FileBufferedOrSeekable
file ->
ReaderIO FileBufferedOrSeekable DataFrame
-> FileBufferedOrSeekable -> IO DataFrame
forall r a. ReaderIO r a -> r -> IO a
runReaderIO (ParquetReadOptions -> ReaderIO FileBufferedOrSeekable DataFrame
forall (m :: * -> *).
(RandomAccess m, MonadIO m) =>
ParquetReadOptions -> m DataFrame
parseParquetWithOpts ParquetReadOptions
opts) FileBufferedOrSeekable
file
readParquetFiles :: FilePath -> IO DataFrame
readParquetFiles :: [Char] -> IO DataFrame
readParquetFiles = ParquetReadOptions -> [Char] -> IO DataFrame
readParquetFilesWithOpts ParquetReadOptions
defaultParquetReadOptions
readParquetFilesWithOpts :: ParquetReadOptions -> FilePath -> IO DataFrame
readParquetFilesWithOpts :: ParquetReadOptions -> [Char] -> IO DataFrame
readParquetFilesWithOpts ParquetReadOptions
opts [Char]
path = do
Bool
isDir <- [Char] -> IO Bool
doesDirectoryExist [Char]
path
let pat :: [Char]
pat = if Bool
isDir then [Char]
path [Char] -> ShowS
</> [Char]
"*.parquet" 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
case [[Char]]
files of
[] ->
[Char] -> IO DataFrame
forall a. HasCallStack => [Char] -> a
error ([Char] -> IO DataFrame) -> [Char] -> IO DataFrame
forall a b. (a -> b) -> a -> b
$
[Char]
"readParquetFiles: no parquet files found for " [Char] -> ShowS
forall a. [a] -> [a] -> [a]
++ [Char]
path
[[Char]]
_ -> do
let optsWithoutRowRange :: ParquetReadOptions
optsWithoutRowRange = ParquetReadOptions
opts{rowRange = Nothing}
[DataFrame]
dfs <- ([Char] -> IO DataFrame) -> [[Char]] -> IO [DataFrame]
forall (t :: * -> *) (m :: * -> *) a b.
(Traversable t, Monad m) =>
(a -> m b) -> t a -> m (t b)
forall (m :: * -> *) a b. Monad m => (a -> m b) -> [a] -> m [b]
mapM (ParquetReadOptions -> [Char] -> IO DataFrame
readParquetWithOpts ParquetReadOptions
optsWithoutRowRange) [[Char]]
files
DataFrame -> IO DataFrame
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (ParquetReadOptions -> DataFrame -> DataFrame
applyRowRange ParquetReadOptions
opts ([DataFrame] -> DataFrame
forall a. Monoid a => [a] -> a
mconcat [DataFrame]
dfs))
parseParquetWithOpts ::
(RandomAccess m, MonadIO m) =>
ParquetReadOptions ->
m DataFrame
{-# SPECIALIZE parseParquetWithOpts ::
ParquetReadOptions -> ReaderIO FileBufferedOrSeekable DataFrame
#-}
{-# INLINEABLE parseParquetWithOpts #-}
parseParquetWithOpts :: forall (m :: * -> *).
(RandomAccess m, MonadIO m) =>
ParquetReadOptions -> m DataFrame
parseParquetWithOpts ParquetReadOptions
opts = do
FileMetadata
metadata <- m FileMetadata
forall (m :: * -> *). RandomAccess m => m FileMetadata
parseFileMetadata
let schemaElems :: [SchemaElement]
schemaElems = Field 2 [SchemaElement] -> [SchemaElement]
forall (n :: Nat) a. KnownNat n => Field n a -> a
unField FileMetadata
metadata.schema
allNames :: [Text]
allNames = [SchemaElement] -> [Text]
getColumnNames (Int -> [SchemaElement] -> [SchemaElement]
forall a. Int -> [a] -> [a]
drop Int
1 [SchemaElement]
schemaElems)
leafNames :: [Text]
leafNames = [Text] -> [Text]
forall a. Eq a => [a] -> [a]
L.nub ((Text -> Text) -> [Text] -> [Text]
forall a b. (a -> b) -> [a] -> [b]
map ([Text] -> Text
forall a. HasCallStack => [a] -> a
last ([Text] -> Text) -> (Text -> [Text]) -> Text -> Text
forall b c a. (b -> c) -> (a -> b) -> a -> c
. HasCallStack => Text -> Text -> [Text]
Text -> Text -> [Text]
T.splitOn Text
".") [Text]
allNames)
predicateColumns :: [Text]
predicateColumns = [Text] -> (Expr Bool -> [Text]) -> Maybe (Expr Bool) -> [Text]
forall b a. b -> (a -> b) -> Maybe a -> b
maybe [] ([Text] -> [Text]
forall a. Eq a => [a] -> [a]
L.nub ([Text] -> [Text]) -> (Expr Bool -> [Text]) -> Expr Bool -> [Text]
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Expr Bool -> [Text]
forall a. Expr a -> [Text]
getColumns) (ParquetReadOptions -> Maybe (Expr Bool)
predicate ParquetReadOptions
opts)
selectedColumnsForRead :: Maybe [Text]
selectedColumnsForRead = case ParquetReadOptions -> Maybe [Text]
selectedColumns ParquetReadOptions
opts of
Maybe [Text]
Nothing -> Maybe [Text]
forall a. Maybe a
Nothing
Just [Text]
selected -> [Text] -> Maybe [Text]
forall a. a -> Maybe a
Just ([Text] -> [Text]
forall a. Eq a => [a] -> [a]
L.nub ([Text]
selected [Text] -> [Text] -> [Text]
forall a. [a] -> [a] -> [a]
++ [Text]
predicateColumns))
case Maybe [Text]
selectedColumnsForRead of
Maybe [Text]
Nothing -> () -> m ()
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ()
Just [Text]
requested ->
let missing :: [Text]
missing = [Text]
requested [Text] -> [Text] -> [Text]
forall a. Eq a => [a] -> [a] -> [a]
L.\\ [Text]
leafNames
in Bool -> m () -> m ()
forall (f :: * -> *). Applicative f => Bool -> f () -> f ()
unless ([Text] -> Bool
forall a. [a] -> Bool
forall (t :: * -> *) a. Foldable t => t a -> Bool
L.null [Text]
missing) (m () -> m ()) -> m () -> m ()
forall a b. (a -> b) -> a -> b
$
IO () -> m ()
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO () -> m ()) -> IO () -> m ()
forall a b. (a -> b) -> a -> b
$
DataFrameException -> IO ()
forall a e. Exception e => e -> a
throw
( [Text] -> Text -> [Text] -> DataFrameException
ColumnsNotFoundException
[Text]
missing
Text
"readParquetWithOpts"
[Text]
leafNames
)
let descriptions :: [ColumnDescription]
descriptions = [SchemaElement] -> [ColumnDescription]
generateColumnDescriptions [SchemaElement]
schemaElems
chunks :: [[ColumnChunk]]
chunks = FileMetadata -> [[ColumnChunk]]
columnChunksForAll FileMetadata
metadata
nCols :: Int
nCols = [[ColumnChunk]] -> Int
forall a. [a] -> Int
forall (t :: * -> *) a. Foldable t => t a -> Int
length [[ColumnChunk]]
chunks
nDescs :: Int
nDescs = [ColumnDescription] -> Int
forall a. [a] -> Int
forall (t :: * -> *) a. Foldable t => t a -> Int
length [ColumnDescription]
descriptions
Bool -> m () -> m ()
forall (f :: * -> *). Applicative f => Bool -> f () -> f ()
unless (Int
nCols Int -> Int -> Bool
forall a. Eq a => a -> a -> Bool
== Int
nDescs) (m () -> m ()) -> m () -> m ()
forall a b. (a -> b) -> a -> b
$
[Char] -> m ()
forall a. HasCallStack => [Char] -> a
error ([Char] -> m ()) -> [Char] -> m ()
forall a b. (a -> b) -> a -> b
$
[Char]
"Column count mismatch: got "
[Char] -> ShowS
forall a. Semigroup a => a -> a -> a
<> Int -> [Char]
forall a. Show a => a -> [Char]
show Int
nCols
[Char] -> ShowS
forall a. Semigroup a => a -> a -> a
<> [Char]
" columns but schema implied "
[Char] -> ShowS
forall a. Semigroup a => a -> a -> a
<> Int -> [Char]
forall a. Show a => a -> [Char]
show Int
nDescs
[Char] -> ShowS
forall a. Semigroup a => a -> a -> a
<> [Char]
" columns"
let topLevelRows :: Int
topLevelRows = Int64 -> Int
forall a b. (Integral a, Num b) => a -> b
fromIntegral (Int64 -> Int) -> (Field 3 Int64 -> Int64) -> Field 3 Int64 -> Int
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Field 3 Int64 -> Int64
forall (n :: Nat) a. KnownNat n => Field n a -> a
unField (Field 3 Int64 -> Int) -> Field 3 Int64 -> Int
forall a b. (a -> b) -> a -> b
$ FileMetadata
metadata.num_rows :: Int
rgRows :: Int
rgRows =
[Int] -> Int
forall a. Num a => [a] -> a
forall (t :: * -> *) a. (Foldable t, Num a) => t a -> a
sum ([Int] -> Int) -> [Int] -> Int
forall a b. (a -> b) -> a -> b
$ (RowGroup -> Int) -> [RowGroup] -> [Int]
forall a b. (a -> b) -> [a] -> [b]
map (Int64 -> Int
forall a b. (Integral a, Num b) => a -> b
fromIntegral (Int64 -> Int) -> (RowGroup -> Int64) -> RowGroup -> Int
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Field 3 Int64 -> Int64
forall (n :: Nat) a. KnownNat n => Field n a -> a
unField (Field 3 Int64 -> Int64)
-> (RowGroup -> Field 3 Int64) -> RowGroup -> Int64
forall b c a. (b -> c) -> (a -> b) -> a -> c
. RowGroup -> Field 3 Int64
rg_num_rows) (Field 4 [RowGroup] -> [RowGroup]
forall (n :: Nat) a. KnownNat n => Field n a -> a
unField FileMetadata
metadata.row_groups) ::
Int
vectorLength :: Int
vectorLength = if Int
topLevelRows Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
> Int
0 then Int
topLevelRows else Int
rgRows
let keep :: Text -> Bool
keep Text
name = case Maybe [Text]
selectedColumnsForRead of
Maybe [Text]
Nothing -> Bool
True
Just [Text]
req -> [Text] -> Text
forall a. HasCallStack => [a] -> a
last (HasCallStack => Text -> Text -> [Text]
Text -> Text -> [Text]
T.splitOn Text
"." Text
name) Text -> [Text] -> Bool
forall a. Eq a => a -> [a] -> Bool
forall (t :: * -> *) a. (Foldable t, Eq a) => a -> t a -> Bool
`elem` [Text]
req
kept :: [(Text, [ColumnChunk], ColumnDescription)]
kept = ((Text, [ColumnChunk], ColumnDescription) -> Bool)
-> [(Text, [ColumnChunk], ColumnDescription)]
-> [(Text, [ColumnChunk], ColumnDescription)]
forall a. (a -> Bool) -> [a] -> [a]
filter (\(Text
n, [ColumnChunk]
_, ColumnDescription
_) -> Text -> Bool
keep Text
n) ([Text]
-> [[ColumnChunk]]
-> [ColumnDescription]
-> [(Text, [ColumnChunk], ColumnDescription)]
forall a b c. [a] -> [b] -> [c] -> [(a, b, c)]
zip3 [Text]
allNames [[ColumnChunk]]
chunks [ColumnDescription]
descriptions)
keptNames :: [Text]
keptNames = [Text
n | (Text
n, [ColumnChunk]
_, ColumnDescription
_) <- [(Text, [ColumnChunk], ColumnDescription)]
kept]
keptChunks :: [[ColumnChunk]]
keptChunks = [[ColumnChunk]
c | (Text
_, [ColumnChunk]
c, ColumnDescription
_) <- [(Text, [ColumnChunk], ColumnDescription)]
kept]
keptDescs :: [ColumnDescription]
keptDescs = [ColumnDescription
d | (Text
_, [ColumnChunk]
_, ColumnDescription
d) <- [(Text, [ColumnChunk], ColumnDescription)]
kept]
[Column]
rawCols <- ([ColumnChunk] -> ColumnDescription -> m Column)
-> [[ColumnChunk]] -> [ColumnDescription] -> m [Column]
forall (m :: * -> *) a b c.
Applicative m =>
(a -> b -> m c) -> [a] -> [b] -> m [c]
zipWithM (Int -> [ColumnChunk] -> ColumnDescription -> m Column
forall (m :: * -> *).
(RandomAccess m, MonadIO m) =>
Int -> [ColumnChunk] -> ColumnDescription -> m Column
parseColumnChunks Int
vectorLength) [[ColumnChunk]]
keptChunks [ColumnDescription]
keptDescs
let finalCols :: [Column]
finalCols = (ColumnDescription -> Column -> Column)
-> [ColumnDescription] -> [Column] -> [Column]
forall a b c. (a -> b -> c) -> [a] -> [b] -> [c]
zipWith ColumnDescription -> Column -> Column
applyDescLogicalType [ColumnDescription]
keptDescs [Column]
rawCols
indices :: Map Text Int
indices = [(Text, Int)] -> Map Text Int
forall k a. Ord k => [(k, a)] -> Map k a
Map.fromList ([(Text, Int)] -> Map Text Int) -> [(Text, Int)] -> Map Text Int
forall a b. (a -> b) -> a -> b
$ [Text] -> [Int] -> [(Text, Int)]
forall a b. [a] -> [b] -> [(a, b)]
zip [Text]
keptNames [Int
0 ..]
dimensions :: (Int, Int)
dimensions = (Int
vectorLength, [Column] -> Int
forall a. [a] -> Int
forall (t :: * -> *) a. Foldable t => t a -> Int
length [Column]
finalCols)
let df :: DataFrame
df =
Vector Column
-> Map Text Int -> (Int, Int) -> Map Text UExpr -> DataFrame
DataFrame
(Int -> [Column] -> Vector Column
forall a. Int -> [a] -> Vector a
Vector.fromListN ([Column] -> Int
forall a. [a] -> Int
forall (t :: * -> *) a. Foldable t => t a -> Int
length [Column]
finalCols) [Column]
finalCols)
Map Text Int
indices
(Int, Int)
dimensions
Map Text UExpr
forall k a. Map k a
Map.empty
DataFrame -> m DataFrame
forall a. a -> m a
forall (m :: * -> *) a. Monad m => a -> m a
return (ParquetReadOptions -> DataFrame -> DataFrame
applyReadOptions ParquetReadOptions
opts DataFrame
df)
parseFileMetadata :: (RandomAccess m) => m FileMetadata
parseFileMetadata :: forall (m :: * -> *). RandomAccess m => m FileMetadata
parseFileMetadata = do
ByteString
footerBytes <- Int -> m ByteString
forall (m :: * -> *). RandomAccess m => Int -> m ByteString
readSuffix Int
8
let magic :: ByteString
magic = Int -> ByteString -> ByteString
BS.drop Int
4 ByteString
footerBytes
Bool -> m () -> m ()
forall (f :: * -> *). Applicative f => Bool -> f () -> f ()
when (ByteString
magic ByteString -> ByteString -> Bool
forall a. Eq a => a -> a -> Bool
/= ByteString
"PAR1") (m () -> m ()) -> m () -> m ()
forall a b. (a -> b) -> a -> b
$
[Char] -> m ()
forall a. HasCallStack => [Char] -> a
error
( [Char]
"Not a valid Parquet file: expected magic bytes \"PAR1\", got "
[Char] -> ShowS
forall a. [a] -> [a] -> [a]
++ ByteString -> [Char]
forall a. Show a => a -> [Char]
show ByteString
magic
)
let size :: Int
size = ByteString -> Int
getMetadataSize ByteString
footerBytes
ByteString
rawMetadata <- Int -> m ByteString
forall (m :: * -> *). RandomAccess m => Int -> m ByteString
readSuffix (Int
size Int -> Int -> Int
forall a. Num a => a -> a -> a
+ Int
8) m ByteString -> (ByteString -> ByteString) -> m ByteString
forall (f :: * -> *) a b. Functor f => f a -> (a -> b) -> f b
<&> Int -> ByteString -> ByteString
BS.take Int
size
case Protocol -> ByteString -> Either [Char] FileMetadata
forall a. Pinchable a => Protocol -> ByteString -> Either [Char] a
Pinch.decode Protocol
Pinch.compactProtocol ByteString
rawMetadata of
Left [Char]
e -> [Char] -> m FileMetadata
forall a. HasCallStack => [Char] -> a
error ([Char] -> m FileMetadata) -> [Char] -> m FileMetadata
forall a b. (a -> b) -> a -> b
$ [Char]
"Failed to parse Parquet metadata: " [Char] -> ShowS
forall a. [a] -> [a] -> [a]
++ ShowS
forall a. Show a => a -> [Char]
show [Char]
e
Right FileMetadata
metadata -> FileMetadata -> m FileMetadata
forall a. a -> m a
forall (m :: * -> *) a. Monad m => a -> m a
return FileMetadata
metadata
where
getMetadataSize :: ByteString -> Int
getMetadataSize ByteString
footer =
let sizes :: [Int]
sizes :: [Int]
sizes = (Int -> Int) -> [Int] -> [Int]
forall a b. (a -> b) -> [a] -> [b]
map (Word8 -> Int
forall a b. (Integral a, Num b) => a -> b
fromIntegral (Word8 -> Int) -> (Int -> Word8) -> Int -> Int
forall b c a. (b -> c) -> (a -> b) -> a -> c
. HasCallStack => ByteString -> Int -> Word8
ByteString -> Int -> Word8
BS.index ByteString
footer) [Int
0 .. Int
3]
in (Int -> Int -> Int) -> Int -> [Int] -> Int
forall b a. (b -> a -> b) -> b -> [a] -> b
forall (t :: * -> *) b a.
Foldable t =>
(b -> a -> b) -> b -> t a -> b
foldl' Int -> Int -> Int
forall a. Bits a => a -> a -> a
(.|.) Int
0 ([Int] -> Int) -> [Int] -> Int
forall a b. (a -> b) -> a -> b
$ (Int -> Int -> Int) -> [Int] -> [Int] -> [Int]
forall a b c. (a -> b -> c) -> [a] -> [b] -> [c]
zipWith Int -> Int -> Int
forall a. Bits a => a -> Int -> a
shiftL [Int]
sizes [Int
0, Int
8 .. Int
24]
readMetadataFromPath :: FilePath -> IO FileMetadata
readMetadataFromPath :: [Char] -> IO FileMetadata
readMetadataFromPath [Char]
path =
ForceNonSeekable
-> [Char]
-> IOMode
-> (FileBufferedOrSeekable -> IO FileMetadata)
-> IO FileMetadata
forall a.
ForceNonSeekable
-> [Char] -> IOMode -> (FileBufferedOrSeekable -> IO a) -> IO a
withFileBufferedOrSeekable ForceNonSeekable
forall a. Maybe a
Nothing [Char]
path IOMode
ReadMode ((FileBufferedOrSeekable -> IO FileMetadata) -> IO FileMetadata)
-> (FileBufferedOrSeekable -> IO FileMetadata) -> IO FileMetadata
forall a b. (a -> b) -> a -> b
$
ReaderIO FileBufferedOrSeekable FileMetadata
-> FileBufferedOrSeekable -> IO FileMetadata
forall r a. ReaderIO r a -> r -> IO a
runReaderIO ReaderIO FileBufferedOrSeekable FileMetadata
forall (m :: * -> *). RandomAccess m => m FileMetadata
parseFileMetadata
readMetadataFromHandle :: FileBufferedOrSeekable -> IO FileMetadata
readMetadataFromHandle :: FileBufferedOrSeekable -> IO FileMetadata
readMetadataFromHandle = ReaderIO FileBufferedOrSeekable FileMetadata
-> FileBufferedOrSeekable -> IO FileMetadata
forall r a. ReaderIO r a -> r -> IO a
runReaderIO ReaderIO FileBufferedOrSeekable FileMetadata
forall (m :: * -> *). RandomAccess m => m FileMetadata
parseFileMetadata
columnChunksForAll :: FileMetadata -> [[ColumnChunk]]
columnChunksForAll :: FileMetadata -> [[ColumnChunk]]
columnChunksForAll =
[[ColumnChunk]] -> [[ColumnChunk]]
forall a. [[a]] -> [[a]]
transpose ([[ColumnChunk]] -> [[ColumnChunk]])
-> (FileMetadata -> [[ColumnChunk]])
-> FileMetadata
-> [[ColumnChunk]]
forall b c a. (b -> c) -> (a -> b) -> a -> c
. (RowGroup -> [ColumnChunk]) -> [RowGroup] -> [[ColumnChunk]]
forall a b. (a -> b) -> [a] -> [b]
map (Field 1 [ColumnChunk] -> [ColumnChunk]
forall (n :: Nat) a. KnownNat n => Field n a -> a
unField (Field 1 [ColumnChunk] -> [ColumnChunk])
-> (RowGroup -> Field 1 [ColumnChunk]) -> RowGroup -> [ColumnChunk]
forall b c a. (b -> c) -> (a -> b) -> a -> c
. RowGroup -> Field 1 [ColumnChunk]
rg_columns) ([RowGroup] -> [[ColumnChunk]])
-> (FileMetadata -> [RowGroup]) -> FileMetadata -> [[ColumnChunk]]
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Field 4 [RowGroup] -> [RowGroup]
forall (n :: Nat) a. KnownNat n => Field n a -> a
unField (Field 4 [RowGroup] -> [RowGroup])
-> (FileMetadata -> Field 4 [RowGroup])
-> FileMetadata
-> [RowGroup]
forall b c a. (b -> c) -> (a -> b) -> a -> c
. FileMetadata -> Field 4 [RowGroup]
row_groups
parseColumnChunks ::
(RandomAccess m, MonadIO m) =>
Int ->
[ColumnChunk] ->
ColumnDescription ->
m Column
{-# SPECIALIZE parseColumnChunks ::
Int ->
[ColumnChunk] ->
ColumnDescription ->
ReaderIO FileBufferedOrSeekable Column
#-}
{-# INLINEABLE parseColumnChunks #-}
parseColumnChunks :: forall (m :: * -> *).
(RandomAccess m, MonadIO m) =>
Int -> [ColumnChunk] -> ColumnDescription -> m Column
parseColumnChunks Int
totalRows [ColumnChunk]
chunks ColumnDescription
description
| ColumnDescription
description.maxRepetitionLevel Int32 -> Int32 -> Bool
forall a. Eq a => a -> a -> Bool
== Int32
0 Bool -> Bool -> Bool
&& ColumnDescription
description.maxDefinitionLevel Int32 -> Int32 -> Bool
forall a. Eq a => a -> a -> Bool
== Int32
0 =
Int -> ColumnDescription -> [ColumnChunk] -> m Column
forall (m :: * -> *).
(RandomAccess m, MonadIO m) =>
Int -> ColumnDescription -> [ColumnChunk] -> m Column
getNonNullableColumn Int
totalRows ColumnDescription
description [ColumnChunk]
chunks
| ColumnDescription
description.maxRepetitionLevel Int32 -> Int32 -> Bool
forall a. Eq a => a -> a -> Bool
== Int32
0 =
Int -> ColumnDescription -> [ColumnChunk] -> m Column
forall (m :: * -> *).
(RandomAccess m, MonadIO m) =>
Int -> ColumnDescription -> [ColumnChunk] -> m Column
getNullableColumn Int
totalRows ColumnDescription
description [ColumnChunk]
chunks
| Bool
otherwise =
ColumnDescription -> [ColumnChunk] -> m Column
forall (m :: * -> *).
(RandomAccess m, MonadIO m) =>
ColumnDescription -> [ColumnChunk] -> m Column
getRepeatedColumn ColumnDescription
description [ColumnChunk]
chunks
{-# INLINEABLE getNonNullableColumn #-}
getNonNullableColumn ::
forall m.
(RandomAccess m, MonadIO m) =>
Int ->
ColumnDescription ->
[ColumnChunk] ->
m Column
getNonNullableColumn :: forall (m :: * -> *).
(RandomAccess m, MonadIO m) =>
Int -> ColumnDescription -> [ColumnChunk] -> m Column
getNonNullableColumn Int
totalRows ColumnDescription
description [ColumnChunk]
chunks =
case ColumnDescription
description.colElementType of
Just (BOOLEAN Enumeration 0
_) -> UnboxedPageDecoder Bool -> m Column
forall a.
(Columnable a, Unbox a) =>
UnboxedPageDecoder a -> m Column
unboxedGo UnboxedPageDecoder Bool
boolDecoder
Just (INT32 Enumeration 1
_) -> UnboxedPageDecoder Int32 -> m Column
forall a.
(Columnable a, Unbox a) =>
UnboxedPageDecoder a -> m Column
unboxedGo UnboxedPageDecoder Int32
int32Decoder
Just (INT64 Enumeration 2
_) -> UnboxedPageDecoder Int64 -> m Column
forall a.
(Columnable a, Unbox a) =>
UnboxedPageDecoder a -> m Column
unboxedGo UnboxedPageDecoder Int64
int64Decoder
Just (INT96 Enumeration 3
_) -> PageDecoder UTCTime -> m Column
forall a. Columnable a => PageDecoder a -> m Column
go PageDecoder UTCTime
int96Decoder
Just (FLOAT Enumeration 4
_) -> UnboxedPageDecoder Float -> m Column
forall a.
(Columnable a, Unbox a) =>
UnboxedPageDecoder a -> m Column
unboxedGo UnboxedPageDecoder Float
floatDecoder
Just (DOUBLE Enumeration 5
_) -> UnboxedPageDecoder Double -> m Column
forall a.
(Columnable a, Unbox a) =>
UnboxedPageDecoder a -> m Column
unboxedGo UnboxedPageDecoder Double
doubleDecoder
Just (BYTE_ARRAY Enumeration 6
_) -> m Column
goPackedText
Just (FIXED_LEN_BYTE_ARRAY Enumeration 7
_) -> case ColumnDescription
description.typeLength of
Maybe Int32
Nothing -> [Char] -> m Column
forall a. HasCallStack => [Char] -> a
error [Char]
"FIXED_LEN_BYTE_ARRAY requires type_length to be set"
Just Int32
tl -> PageDecoder Text -> m Column
forall a. Columnable a => PageDecoder a -> m Column
go (Int -> PageDecoder Text
fixedLenByteArrayDecoder (Int32 -> Int
forall a b. (Integral a, Num b) => a -> b
fromIntegral Int32
tl))
Maybe ThriftType
Nothing -> [Char] -> m Column
forall a. HasCallStack => [Char] -> a
error [Char]
"Column has no Parquet type"
where
go ::
forall a.
(Columnable a) =>
PageDecoder a ->
m Column
go :: forall a. Columnable a => PageDecoder a -> m Column
go PageDecoder a
decoder =
Int -> PageFold m Vector a -> m Column
forall (m :: * -> *) a.
(MonadIO m, Columnable a) =>
Int -> PageFold m Vector a -> m Column
foldNonNullable Int
totalRows (ColumnDescription
-> PageDecoder a
-> [ColumnChunk]
-> (acc -> (Vector a, Vector Int, Vector Int) -> m acc)
-> acc
-> m acc
forall (m :: * -> *) (v :: * -> *) a acc.
(RandomAccess m, MonadIO m, Vector v a) =>
ColumnDescription
-> (Maybe DictVals -> Encoding -> Int -> ByteString -> v a)
-> [ColumnChunk]
-> (acc -> (v a, Vector Int, Vector Int) -> m acc)
-> acc
-> m acc
foldColumnPagesM ColumnDescription
description PageDecoder a
decoder [ColumnChunk]
chunks)
goPackedText :: m Column
goPackedText :: m Column
goPackedText = do
TextBuilder RealWorld
builder <- IO (TextBuilder RealWorld) -> m (TextBuilder RealWorld)
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO (TextBuilder RealWorld) -> m (TextBuilder RealWorld))
-> IO (TextBuilder RealWorld) -> m (TextBuilder RealWorld)
forall a b. (a -> b) -> a -> b
$ ST RealWorld (TextBuilder RealWorld) -> IO (TextBuilder RealWorld)
forall a. ST RealWorld a -> IO a
stToIO (Int -> Int -> ST RealWorld (TextBuilder RealWorld)
forall s. Int -> Int -> ST s (TextBuilder s)
newTextBuilder Int
totalRows (Int
totalRows Int -> Int -> Int
forall a. Num a => a -> a -> a
* Int
8))
()
_ <-
ColumnDescription
-> [ColumnChunk] -> (() -> RawPage -> m ()) -> () -> m ()
forall (m :: * -> *) acc.
(RandomAccess m, MonadIO m) =>
ColumnDescription
-> [ColumnChunk] -> (acc -> RawPage -> m acc) -> acc -> m acc
foldColumnDataPagesM
ColumnDescription
description
[ColumnChunk]
chunks
( \() (Maybe DictVals
dict, Encoding
enc, Int
nPresent, ByteString
valBytes, Vector Int
_, Vector Int
_) ->
IO () -> m ()
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (TextBuilder RealWorld
-> Maybe DictVals -> Encoding -> Int -> ByteString -> IO ()
appendStringPageIO TextBuilder RealWorld
builder Maybe DictVals
dict Encoding
enc Int
nPresent ByteString
valBytes)
)
()
TextChunk
chunk <- IO TextChunk -> m TextChunk
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO TextChunk -> m TextChunk) -> IO TextChunk -> m TextChunk
forall a b. (a -> b) -> a -> b
$ ST RealWorld TextChunk -> IO TextChunk
forall a. ST RealWorld a -> IO a
stToIO (TextBuilder RealWorld -> ST RealWorld TextChunk
forall s. TextBuilder s -> ST s TextChunk
freezeTextChunk TextBuilder RealWorld
builder)
Column -> m Column
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ([TextChunk] -> Column
mergeTextChunks [TextChunk
chunk])
unboxedGo ::
forall a.
(Columnable a, VU.Unbox a) =>
UnboxedPageDecoder a ->
m Column
unboxedGo :: forall a.
(Columnable a, Unbox a) =>
UnboxedPageDecoder a -> m Column
unboxedGo UnboxedPageDecoder a
decoder =
Int -> PageFold m Vector a -> m Column
forall (m :: * -> *) a.
(MonadIO m, Columnable a, Unbox a) =>
Int -> PageFold m Vector a -> m Column
foldNonNullableUnboxed Int
totalRows (ColumnDescription
-> UnboxedPageDecoder a
-> [ColumnChunk]
-> (acc -> (Vector a, Vector Int, Vector Int) -> m acc)
-> acc
-> m acc
forall (m :: * -> *) (v :: * -> *) a acc.
(RandomAccess m, MonadIO m, Vector v a) =>
ColumnDescription
-> (Maybe DictVals -> Encoding -> Int -> ByteString -> v a)
-> [ColumnChunk]
-> (acc -> (v a, Vector Int, Vector Int) -> m acc)
-> acc
-> m acc
foldColumnPagesM ColumnDescription
description UnboxedPageDecoder a
decoder [ColumnChunk]
chunks)
{-# INLINEABLE getNullableColumn #-}
getNullableColumn ::
forall m.
(RandomAccess m, MonadIO m) =>
Int ->
ColumnDescription ->
[ColumnChunk] ->
m Column
getNullableColumn :: forall (m :: * -> *).
(RandomAccess m, MonadIO m) =>
Int -> ColumnDescription -> [ColumnChunk] -> m Column
getNullableColumn Int
totalRows ColumnDescription
description [ColumnChunk]
chunks =
case ColumnDescription
description.colElementType of
Just (BOOLEAN Enumeration 0
_) -> UnboxedPageDecoder Bool -> m Column
forall a.
(Columnable a, Unbox a) =>
UnboxedPageDecoder a -> m Column
unboxedGo UnboxedPageDecoder Bool
boolDecoder
Just (INT32 Enumeration 1
_) -> UnboxedPageDecoder Int32 -> m Column
forall a.
(Columnable a, Unbox a) =>
UnboxedPageDecoder a -> m Column
unboxedGo UnboxedPageDecoder Int32
int32Decoder
Just (INT64 Enumeration 2
_) -> UnboxedPageDecoder Int64 -> m Column
forall a.
(Columnable a, Unbox a) =>
UnboxedPageDecoder a -> m Column
unboxedGo UnboxedPageDecoder Int64
int64Decoder
Just (INT96 Enumeration 3
_) -> PageDecoder UTCTime -> m Column
forall a. Columnable a => PageDecoder a -> m Column
go PageDecoder UTCTime
int96Decoder
Just (FLOAT Enumeration 4
_) -> UnboxedPageDecoder Float -> m Column
forall a.
(Columnable a, Unbox a) =>
UnboxedPageDecoder a -> m Column
unboxedGo UnboxedPageDecoder Float
floatDecoder
Just (DOUBLE Enumeration 5
_) -> UnboxedPageDecoder Double -> m Column
forall a.
(Columnable a, Unbox a) =>
UnboxedPageDecoder a -> m Column
unboxedGo UnboxedPageDecoder Double
doubleDecoder
Just (BYTE_ARRAY Enumeration 6
_) -> m Column
goPackedTextNullable
Just (FIXED_LEN_BYTE_ARRAY Enumeration 7
_) -> case ColumnDescription
description.typeLength of
Maybe Int32
Nothing -> [Char] -> m Column
forall a. HasCallStack => [Char] -> a
error [Char]
"FIXED_LEN_BYTE_ARRAY requires type_length to be set"
Just Int32
tl -> PageDecoder Text -> m Column
forall a. Columnable a => PageDecoder a -> m Column
go (Int -> PageDecoder Text
fixedLenByteArrayDecoder (Int32 -> Int
forall a b. (Integral a, Num b) => a -> b
fromIntegral Int32
tl))
Maybe ThriftType
Nothing -> [Char] -> m Column
forall a. HasCallStack => [Char] -> a
error [Char]
"Column has no Parquet type"
where
maxDef :: Int
maxDef :: Int
maxDef = Int32 -> Int
forall a b. (Integral a, Num b) => a -> b
fromIntegral ColumnDescription
description.maxDefinitionLevel
go ::
forall a.
(Columnable a) =>
PageDecoder a ->
m Column
go :: forall a. Columnable a => PageDecoder a -> m Column
go PageDecoder a
decoder =
Int -> Int -> PageFold m Vector a -> m Column
forall (m :: * -> *) a.
(MonadIO m, Columnable a) =>
Int -> Int -> PageFold m Vector a -> m Column
foldNullable Int
maxDef Int
totalRows (ColumnDescription
-> PageDecoder a
-> [ColumnChunk]
-> (acc -> (Vector a, Vector Int, Vector Int) -> m acc)
-> acc
-> m acc
forall (m :: * -> *) (v :: * -> *) a acc.
(RandomAccess m, MonadIO m, Vector v a) =>
ColumnDescription
-> (Maybe DictVals -> Encoding -> Int -> ByteString -> v a)
-> [ColumnChunk]
-> (acc -> (v a, Vector Int, Vector Int) -> m acc)
-> acc
-> m acc
foldColumnPagesM ColumnDescription
description PageDecoder a
decoder [ColumnChunk]
chunks)
goPackedTextNullable :: m Column
goPackedTextNullable :: m Column
goPackedTextNullable = do
TextBuilder RealWorld
builder <- IO (TextBuilder RealWorld) -> m (TextBuilder RealWorld)
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO (TextBuilder RealWorld) -> m (TextBuilder RealWorld))
-> IO (TextBuilder RealWorld) -> m (TextBuilder RealWorld)
forall a b. (a -> b) -> a -> b
$ ST RealWorld (TextBuilder RealWorld) -> IO (TextBuilder RealWorld)
forall a. ST RealWorld a -> IO a
stToIO (Int -> Int -> ST RealWorld (TextBuilder RealWorld)
forall s. Int -> Int -> ST s (TextBuilder s)
newTextBuilder Int
totalRows (Int
totalRows Int -> Int -> Int
forall a. Num a => a -> a -> a
* Int
8))
()
_ <-
ColumnDescription
-> [ColumnChunk] -> (() -> RawPage -> m ()) -> () -> m ()
forall (m :: * -> *) acc.
(RandomAccess m, MonadIO m) =>
ColumnDescription
-> [ColumnChunk] -> (acc -> RawPage -> m acc) -> acc -> m acc
foldColumnDataPagesM
ColumnDescription
description
[ColumnChunk]
chunks
( \() (Maybe DictVals
dict, Encoding
enc, Int
nPresent, ByteString
valBytes, Vector Int
defs, Vector Int
_) ->
IO () -> m ()
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO
(TextBuilder RealWorld
-> Int
-> Maybe DictVals
-> Encoding
-> Int
-> ByteString
-> Vector Int
-> IO ()
appendNullableStringPageIO TextBuilder RealWorld
builder Int
maxDef Maybe DictVals
dict Encoding
enc Int
nPresent ByteString
valBytes Vector Int
defs)
)
()
TextChunk
chunk <- IO TextChunk -> m TextChunk
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO TextChunk -> m TextChunk) -> IO TextChunk -> m TextChunk
forall a b. (a -> b) -> a -> b
$ ST RealWorld TextChunk -> IO TextChunk
forall a. ST RealWorld a -> IO a
stToIO (TextBuilder RealWorld -> ST RealWorld TextChunk
forall s. TextBuilder s -> ST s TextChunk
freezeTextChunk TextBuilder RealWorld
builder)
Column -> m Column
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ([TextChunk] -> Column
mergeTextChunks [TextChunk
chunk])
unboxedGo ::
forall a.
(Columnable a, VU.Unbox a) =>
UnboxedPageDecoder a ->
m Column
unboxedGo :: forall a.
(Columnable a, Unbox a) =>
UnboxedPageDecoder a -> m Column
unboxedGo UnboxedPageDecoder a
decoder =
Int -> Int -> PageFold m Vector a -> m Column
forall (m :: * -> *) a.
(MonadIO m, Columnable a, Unbox a) =>
Int -> Int -> PageFold m Vector a -> m Column
foldNullableUnboxed
Int
maxDef
Int
totalRows
(ColumnDescription
-> UnboxedPageDecoder a
-> [ColumnChunk]
-> (acc -> (Vector a, Vector Int, Vector Int) -> m acc)
-> acc
-> m acc
forall (m :: * -> *) (v :: * -> *) a acc.
(RandomAccess m, MonadIO m, Vector v a) =>
ColumnDescription
-> (Maybe DictVals -> Encoding -> Int -> ByteString -> v a)
-> [ColumnChunk]
-> (acc -> (v a, Vector Int, Vector Int) -> m acc)
-> acc
-> m acc
foldColumnPagesM ColumnDescription
description UnboxedPageDecoder a
decoder [ColumnChunk]
chunks)
{-# INLINEABLE getRepeatedColumn #-}
getRepeatedColumn ::
forall m.
(RandomAccess m, MonadIO m) =>
ColumnDescription ->
[ColumnChunk] ->
m Column
getRepeatedColumn :: forall (m :: * -> *).
(RandomAccess m, MonadIO m) =>
ColumnDescription -> [ColumnChunk] -> m Column
getRepeatedColumn ColumnDescription
description [ColumnChunk]
chunks =
case ColumnDescription
description.colElementType of
Just (BOOLEAN Enumeration 0
_) -> UnboxedPageDecoder Bool -> m Column
forall a.
(Unbox a, Columnable a, Columnable (Maybe [Maybe a]),
Columnable (Maybe [Maybe [Maybe a]]),
Columnable (Maybe [Maybe [Maybe [Maybe a]]])) =>
UnboxedPageDecoder a -> m Column
unboxedGo UnboxedPageDecoder Bool
boolDecoder
Just (INT32 Enumeration 1
_) -> UnboxedPageDecoder Int32 -> m Column
forall a.
(Unbox a, Columnable a, Columnable (Maybe [Maybe a]),
Columnable (Maybe [Maybe [Maybe a]]),
Columnable (Maybe [Maybe [Maybe [Maybe a]]])) =>
UnboxedPageDecoder a -> m Column
unboxedGo UnboxedPageDecoder Int32
int32Decoder
Just (INT64 Enumeration 2
_) -> UnboxedPageDecoder Int64 -> m Column
forall a.
(Unbox a, Columnable a, Columnable (Maybe [Maybe a]),
Columnable (Maybe [Maybe [Maybe a]]),
Columnable (Maybe [Maybe [Maybe [Maybe a]]])) =>
UnboxedPageDecoder a -> m Column
unboxedGo UnboxedPageDecoder Int64
int64Decoder
Just (INT96 Enumeration 3
_) -> PageDecoder UTCTime -> m Column
forall a.
(Columnable a, Columnable (Maybe [Maybe a]),
Columnable (Maybe [Maybe [Maybe a]]),
Columnable (Maybe [Maybe [Maybe [Maybe a]]])) =>
PageDecoder a -> m Column
go PageDecoder UTCTime
int96Decoder
Just (FLOAT Enumeration 4
_) -> UnboxedPageDecoder Float -> m Column
forall a.
(Unbox a, Columnable a, Columnable (Maybe [Maybe a]),
Columnable (Maybe [Maybe [Maybe a]]),
Columnable (Maybe [Maybe [Maybe [Maybe a]]])) =>
UnboxedPageDecoder a -> m Column
unboxedGo UnboxedPageDecoder Float
floatDecoder
Just (DOUBLE Enumeration 5
_) -> UnboxedPageDecoder Double -> m Column
forall a.
(Unbox a, Columnable a, Columnable (Maybe [Maybe a]),
Columnable (Maybe [Maybe [Maybe a]]),
Columnable (Maybe [Maybe [Maybe [Maybe a]]])) =>
UnboxedPageDecoder a -> m Column
unboxedGo UnboxedPageDecoder Double
doubleDecoder
Just (BYTE_ARRAY Enumeration 6
_) -> PageDecoder Text -> m Column
forall a.
(Columnable a, Columnable (Maybe [Maybe a]),
Columnable (Maybe [Maybe [Maybe a]]),
Columnable (Maybe [Maybe [Maybe [Maybe a]]])) =>
PageDecoder a -> m Column
go PageDecoder Text
byteArrayDecoder
Just (FIXED_LEN_BYTE_ARRAY Enumeration 7
_) -> case ColumnDescription
description.typeLength of
Maybe Int32
Nothing -> [Char] -> m Column
forall a. HasCallStack => [Char] -> a
error [Char]
"FIXED_LEN_BYTE_ARRAY requires type_length to be set"
Just Int32
tl -> PageDecoder Text -> m Column
forall a.
(Columnable a, Columnable (Maybe [Maybe a]),
Columnable (Maybe [Maybe [Maybe a]]),
Columnable (Maybe [Maybe [Maybe [Maybe a]]])) =>
PageDecoder a -> m Column
go (Int -> PageDecoder Text
fixedLenByteArrayDecoder (Int32 -> Int
forall a b. (Integral a, Num b) => a -> b
fromIntegral Int32
tl))
Maybe ThriftType
Nothing -> [Char] -> m Column
forall a. HasCallStack => [Char] -> a
error [Char]
"Column has no Parquet type"
where
maxRep :: Int
maxRep :: Int
maxRep = Int32 -> Int
forall a b. (Integral a, Num b) => a -> b
fromIntegral ColumnDescription
description.maxRepetitionLevel
maxDef :: Int
maxDef :: Int
maxDef = Int32 -> Int
forall a b. (Integral a, Num b) => a -> b
fromIntegral ColumnDescription
description.maxDefinitionLevel
go ::
forall a.
( Columnable a
, Columnable (Maybe [Maybe a])
, Columnable (Maybe [Maybe [Maybe a]])
, Columnable (Maybe [Maybe [Maybe [Maybe a]]])
) =>
PageDecoder a ->
m Column
go :: forall a.
(Columnable a, Columnable (Maybe [Maybe a]),
Columnable (Maybe [Maybe [Maybe a]]),
Columnable (Maybe [Maybe [Maybe [Maybe a]]])) =>
PageDecoder a -> m Column
go PageDecoder a
decoder =
Int -> Int -> PageFold m Vector a -> m Column
forall (m :: * -> *) a.
(MonadIO m, Columnable a, Columnable (Maybe [Maybe a]),
Columnable (Maybe [Maybe [Maybe a]]),
Columnable (Maybe [Maybe [Maybe [Maybe a]]])) =>
Int -> Int -> PageFold m Vector a -> m Column
foldRepeated Int
maxRep Int
maxDef (ColumnDescription
-> PageDecoder a
-> [ColumnChunk]
-> (acc -> (Vector a, Vector Int, Vector Int) -> m acc)
-> acc
-> m acc
forall (m :: * -> *) (v :: * -> *) a acc.
(RandomAccess m, MonadIO m, Vector v a) =>
ColumnDescription
-> (Maybe DictVals -> Encoding -> Int -> ByteString -> v a)
-> [ColumnChunk]
-> (acc -> (v a, Vector Int, Vector Int) -> m acc)
-> acc
-> m acc
foldColumnPagesM ColumnDescription
description PageDecoder a
decoder [ColumnChunk]
chunks)
unboxedGo ::
forall a.
( VU.Unbox a
, Columnable a
, Columnable (Maybe [Maybe a])
, Columnable (Maybe [Maybe [Maybe a]])
, Columnable (Maybe [Maybe [Maybe [Maybe a]]])
) =>
UnboxedPageDecoder a ->
m Column
unboxedGo :: forall a.
(Unbox a, Columnable a, Columnable (Maybe [Maybe a]),
Columnable (Maybe [Maybe [Maybe a]]),
Columnable (Maybe [Maybe [Maybe [Maybe a]]])) =>
UnboxedPageDecoder a -> m Column
unboxedGo UnboxedPageDecoder a
decoder =
Int -> Int -> PageFold m Vector a -> m Column
forall (m :: * -> *) a.
(MonadIO m, Columnable a, Unbox a, Columnable (Maybe [Maybe a]),
Columnable (Maybe [Maybe [Maybe a]]),
Columnable (Maybe [Maybe [Maybe [Maybe a]]])) =>
Int -> Int -> PageFold m Vector a -> m Column
foldRepeatedUnboxed Int
maxRep Int
maxDef (ColumnDescription
-> UnboxedPageDecoder a
-> [ColumnChunk]
-> (acc -> (Vector a, Vector Int, Vector Int) -> m acc)
-> acc
-> m acc
forall (m :: * -> *) (v :: * -> *) a acc.
(RandomAccess m, MonadIO m, Vector v a) =>
ColumnDescription
-> (Maybe DictVals -> Encoding -> Int -> ByteString -> v a)
-> [ColumnChunk]
-> (acc -> (v a, Vector Int, Vector Int) -> m acc)
-> acc
-> m acc
foldColumnPagesM ColumnDescription
description UnboxedPageDecoder a
decoder [ColumnChunk]
chunks)
applyRowRange :: ParquetReadOptions -> DataFrame -> DataFrame
applyRowRange :: ParquetReadOptions -> DataFrame -> DataFrame
applyRowRange ParquetReadOptions
opts DataFrame
df =
DataFrame
-> ((Int, Int) -> DataFrame) -> Maybe (Int, Int) -> DataFrame
forall b a. b -> (a -> b) -> Maybe a -> b
maybe DataFrame
df ((Int, Int) -> DataFrame -> DataFrame
`DS.range` DataFrame
df) (ParquetReadOptions -> Maybe (Int, Int)
rowRange ParquetReadOptions
opts)
applySelectedColumns :: ParquetReadOptions -> DataFrame -> DataFrame
applySelectedColumns :: ParquetReadOptions -> DataFrame -> DataFrame
applySelectedColumns ParquetReadOptions
opts DataFrame
df =
DataFrame -> ([Text] -> DataFrame) -> Maybe [Text] -> DataFrame
forall b a. b -> (a -> b) -> Maybe a -> b
maybe DataFrame
df ([Text] -> DataFrame -> DataFrame
`DS.select` DataFrame
df) (ParquetReadOptions -> Maybe [Text]
selectedColumns ParquetReadOptions
opts)
applyPredicate :: ParquetReadOptions -> DataFrame -> DataFrame
applyPredicate :: ParquetReadOptions -> DataFrame -> DataFrame
applyPredicate ParquetReadOptions
opts DataFrame
df =
DataFrame
-> (Expr Bool -> DataFrame) -> Maybe (Expr Bool) -> DataFrame
forall b a. b -> (a -> b) -> Maybe a -> b
maybe DataFrame
df (Expr Bool -> DataFrame -> DataFrame
`DS.filterWhere` DataFrame
df) (ParquetReadOptions -> Maybe (Expr Bool)
predicate ParquetReadOptions
opts)
applySafeRead :: ParquetReadOptions -> DataFrame -> DataFrame
applySafeRead :: ParquetReadOptions -> DataFrame -> DataFrame
applySafeRead ParquetReadOptions
opts DataFrame
df
| ParquetReadOptions -> Bool
safeColumns ParquetReadOptions
opts = DataFrame
df{columns = Vector.map DI.ensureOptional (columns df)}
| Bool
otherwise = DataFrame
df
applyReadOptions :: ParquetReadOptions -> DataFrame -> DataFrame
applyReadOptions :: ParquetReadOptions -> DataFrame -> DataFrame
applyReadOptions ParquetReadOptions
opts =
ParquetReadOptions -> DataFrame -> DataFrame
applySafeRead ParquetReadOptions
opts
(DataFrame -> DataFrame)
-> (DataFrame -> DataFrame) -> DataFrame -> DataFrame
forall b c a. (b -> c) -> (a -> b) -> a -> c
. ParquetReadOptions -> DataFrame -> DataFrame
applyRowRange ParquetReadOptions
opts
(DataFrame -> DataFrame)
-> (DataFrame -> DataFrame) -> DataFrame -> DataFrame
forall b c a. (b -> c) -> (a -> b) -> a -> c
. ParquetReadOptions -> DataFrame -> DataFrame
applySelectedColumns ParquetReadOptions
opts
(DataFrame -> DataFrame)
-> (DataFrame -> DataFrame) -> DataFrame -> DataFrame
forall b c a. (b -> c) -> (a -> b) -> a -> c
. ParquetReadOptions -> DataFrame -> DataFrame
applyPredicate ParquetReadOptions
opts
applyDescLogicalType :: ColumnDescription -> DI.Column -> DI.Column
applyDescLogicalType :: ColumnDescription -> Column -> Column
applyDescLogicalType ColumnDescription
desc = Maybe LogicalType -> Column -> Column
applyLogicalType (ColumnDescription -> Maybe LogicalType
colLogicalType ColumnDescription
desc)
applyLogicalType :: Maybe LogicalType -> DI.Column -> DI.Column
applyLogicalType :: Maybe LogicalType -> Column -> Column
applyLogicalType (Just (LT_TIMESTAMP Field 8 TimestampType
f)) Column
col =
let ts :: TimestampType
ts = Field 8 TimestampType -> TimestampType
forall (n :: Nat) a. KnownNat n => Field n a -> a
unField Field 8 TimestampType
f
unit :: TimeUnit
unit = Field 2 TimeUnit -> TimeUnit
forall (n :: Nat) a. KnownNat n => Field n a -> a
unField TimestampType
ts.timestamp_unit
conv :: Int64 -> UTCTime
conv = case TimeUnit
unit of
MILLIS Field 1 MilliSeconds
_ -> Int64 -> Integer -> Int64 -> UTCTime
epochToUTCTime Int64
1_000 Integer
1_000_000_000
MICROS Field 2 MicroSeconds
_ -> Int64 -> Integer -> Int64 -> UTCTime
epochToUTCTime Int64
1_000_000 Integer
1_000_000
NANOS Field 3 NanoSeconds
_ -> Int64 -> Integer -> Int64 -> UTCTime
epochToUTCTime Int64
1_000_000_000 Integer
1_000
in Column -> Either DataFrameException Column -> Column
forall b a. b -> Either a b -> b
fromRight Column
col (Either DataFrameException Column -> Column)
-> Either DataFrameException Column -> Column
forall a b. (a -> b) -> a -> b
$ (Int64 -> UTCTime) -> Column -> Either DataFrameException Column
forall b c.
(Columnable b, Columnable c) =>
(b -> c) -> Column -> Either DataFrameException Column
DI.mapColumn Int64 -> UTCTime
conv Column
col
applyLogicalType (Just (LT_DECIMAL Field 5 DecimalType
f)) Column
col =
let dt :: DecimalType
dt = Field 5 DecimalType -> DecimalType
forall (n :: Nat) a. KnownNat n => Field n a -> a
unField Field 5 DecimalType
f
scale :: Int32
scale = Field 1 Int32 -> Int32
forall (n :: Nat) a. KnownNat n => Field n a -> a
unField DecimalType
dt.decimal_scale
precision :: Int32
precision = Field 2 Int32 -> Int32
forall (n :: Nat) a. KnownNat n => Field n a -> a
unField DecimalType
dt.decimal_precision
in if Int32
precision Int32 -> Int32 -> Bool
forall a. Ord a => a -> a -> Bool
<= Int32
9
then case forall a (v :: * -> *).
(Vector v a, Columnable a) =>
Column -> Either DataFrameException (v a)
DI.toVector @Int32 @VU.Vector Column
col of
Right Vector Int32
xs ->
Vector Double -> Column
forall a. (Columnable a, Unbox a) => Vector a -> Column
DI.fromUnboxedVector (Vector Double -> Column) -> Vector Double -> Column
forall a b. (a -> b) -> a -> b
$
(Int32 -> Double) -> Vector Int32 -> Vector Double
forall a b. (Unbox a, Unbox b) => (a -> b) -> Vector a -> Vector b
VU.map (\Int32
raw -> forall a b. (Integral a, Num b) => a -> b
fromIntegral @Int32 @Double Int32
raw Double -> Double -> Double
forall a. Fractional a => a -> a -> a
/ Double
10 Double -> Int32 -> Double
forall a b. (Num a, Integral b) => a -> b -> a
^ Int32
scale) Vector Int32
xs
Left DataFrameException
_ -> Column
col
else
if Int32
precision Int32 -> Int32 -> Bool
forall a. Ord a => a -> a -> Bool
<= Int32
18
then case forall a (v :: * -> *).
(Vector v a, Columnable a) =>
Column -> Either DataFrameException (v a)
DI.toVector @Int64 @VU.Vector Column
col of
Right Vector Int64
xs ->
Vector Double -> Column
forall a. (Columnable a, Unbox a) => Vector a -> Column
DI.fromUnboxedVector (Vector Double -> Column) -> Vector Double -> Column
forall a b. (a -> b) -> a -> b
$
(Int64 -> Double) -> Vector Int64 -> Vector Double
forall a b. (Unbox a, Unbox b) => (a -> b) -> Vector a -> Vector b
VU.map (\Int64
raw -> forall a b. (Integral a, Num b) => a -> b
fromIntegral @Int64 @Double Int64
raw Double -> Double -> Double
forall a. Fractional a => a -> a -> a
/ Double
10 Double -> Int32 -> Double
forall a b. (Num a, Integral b) => a -> b -> a
^ Int32
scale) Vector Int64
xs
Left DataFrameException
_ -> Column
col
else Column
col
applyLogicalType Maybe LogicalType
_ Column
col = Column
col
epochToUTCTime :: Int64 -> Integer -> Int64 -> UTCTime
epochToUTCTime :: Int64 -> Integer -> Int64 -> UTCTime
epochToUTCTime Int64
ticksPerSecond Integer
psPerTick Int64
v =
let (Int64
s, Int64
subTicks) = Int64
v Int64 -> Int64 -> (Int64, Int64)
forall a. Integral a => a -> a -> (a, a)
`divMod` Int64
ticksPerSecond
(Int64
days, Int64
sInDay) = Int64
s Int64 -> Int64 -> (Int64, Int64)
forall a. Integral a => a -> a -> (a, a)
`divMod` Int64
86_400
ps :: Integer
ps = Int64 -> Integer
forall a b. (Integral a, Num b) => a -> b
fromIntegral Int64
sInDay Integer -> Integer -> Integer
forall a. Num a => a -> a -> a
* Integer
1_000_000_000_000 Integer -> Integer -> Integer
forall a. Num a => a -> a -> a
+ Int64 -> Integer
forall a b. (Integral a, Num b) => a -> b
fromIntegral Int64
subTicks Integer -> Integer -> Integer
forall a. Num a => a -> a -> a
* Integer
psPerTick
in Day -> DiffTime -> UTCTime
UTCTime
(Integer -> Day
ModifiedJulianDay (Integer
40_587 Integer -> Integer -> Integer
forall a. Num a => a -> a -> a
+ Int64 -> Integer
forall a b. (Integral a, Num b) => a -> b
fromIntegral Int64
days))
(Integer -> DiffTime
picosecondsToDiffTime Integer
ps)