{-# 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))

-- Options -----------------------------------------------------------------

{- | Options for reading Parquet data.

These options are applied in this order:

1. predicate filtering
2. column projection
3. row range
4. safe column promotion

Column selection for @selectedColumns@ uses leaf column names only.
-}
data ParquetReadOptions = ParquetReadOptions
    { ParquetReadOptions -> Maybe [Text]
selectedColumns :: Maybe [T.Text]
    {- ^ Columns to keep in the final dataframe. If set, only these columns are returned.
    Predicate-referenced columns are read automatically when needed and projected out after filtering.
    -}
    , ParquetReadOptions -> Maybe (Expr Bool)
predicate :: Maybe (Expr Bool)
    -- ^ Optional row filter expression applied before projection.
    , ParquetReadOptions -> Maybe (Int, Int)
rowRange :: Maybe (Int, Int)
    -- ^ Optional row slice @(start, end)@ with start-inclusive/end-exclusive semantics.
    , ParquetReadOptions -> Bool
safeColumns :: Bool
    -- ^ When True, every column is promoted to OptionalColumn after read, regardless of nullability in the schema.
    }
    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)

{- | Default Parquet read options.

Equivalent to:

@
ParquetReadOptions
    { selectedColumns = Nothing
    , predicate = Nothing
    , rowRange = Nothing
    , safeColumns = False
    }
@
-}
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
        }

-- Public API --------------------------------------------------------------

{- | Read a parquet file from path and load it into a dataframe.

==== __Example__
@
ghci> D.readParquet ".\/data\/mtcars.parquet"
@
-}
readParquet :: FilePath -> IO DataFrame
readParquet :: [Char] -> IO DataFrame
readParquet = ParquetReadOptions -> [Char] -> IO DataFrame
readParquetWithOpts ParquetReadOptions
defaultParquetReadOptions

{- | Read a Parquet file using explicit read options.

==== __Example__
@
ghci> D.readParquetWithOpts
ghci|   (D.defaultParquetReadOptions{D.selectedColumns = Just ["id"], D.rowRange = Just (0, 10)})
ghci|   "./tests/data/alltypes_plain.parquet"
@

When @selectedColumns@ is set and @predicate@ references other columns, those predicate columns
are auto-included for decoding, then projected back to the requested output columns.
-}
readParquetWithOpts :: ParquetReadOptions -> FilePath -> IO DataFrame
readParquetWithOpts :: ParquetReadOptions -> [Char] -> IO DataFrame
readParquetWithOpts = ForceNonSeekable -> ParquetReadOptions -> [Char] -> IO DataFrame
_readParquetWithOpts ForceNonSeekable
forall a. Maybe a
Nothing

-- | Internal entry point used by tests to force non-seekable mode.
_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

{- | Read Parquet files from a directory or glob path.

This is equivalent to calling 'readParquetFilesWithOpts' with 'defaultParquetReadOptions'.
-}
readParquetFiles :: FilePath -> IO DataFrame
readParquetFiles :: [Char] -> IO DataFrame
readParquetFiles = ParquetReadOptions -> [Char] -> IO DataFrame
readParquetFilesWithOpts ParquetReadOptions
defaultParquetReadOptions

{- | Read multiple Parquet files (directory or glob) using explicit options.

If @path@ is a directory, all non-directory entries are read.
If @path@ is a glob, matching files are read.

For multi-file reads, @rowRange@ is applied once after concatenation (global range semantics).

==== __Example__
@
ghci> D.readParquetFilesWithOpts
ghci|   (D.defaultParquetReadOptions{D.selectedColumns = Just ["id"], D.rowRange = Just (0, 5)})
ghci|   "./tests/data/alltypes_plain*.parquet"
@
-}
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))

-- Core parsing pipeline ---------------------------------------------------

{- | Parse a Parquet file via the 'RandomAccess' handle, applying all
read options. This is the central parsing entry point used by
'_readParquetWithOpts'.
-}
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))

    -- TODO: When rowRange is set, compute cumulative row offsets from
    -- rg_num_rows in each RowGroup and skip any group whose row interval does
    -- not overlap the requested range, avoiding all decoding for those groups.

    -- TODO: When predicate is set, inspect cmd_statistics min/max values for
    -- predicate-referenced columns in each RowGroup and skip groups where
    -- statistics prove the predicate cannot be satisfied.

    -- Validate selected columns
    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"

    -- Some files omit the top-level num_rows field; fall back to summing row-group counts.
    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

    -- Column-projection pushdown: decode only the columns needed for the
    -- requested output plus any predicate, instead of decoding every column
    -- and dropping the unwanted ones afterward. A 'Nothing' selection keeps
    -- all columns, so full reads are unchanged.
    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)

{- | Parse the file-level Thrift metadata from the Parquet file footer.
Validates the trailing 4-byte magic marker (\"PAR1\") before decoding.
-}
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]

-- | Read the file metadata from a Parquet file at the given path.
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

-- | Read only the file metadata from an open 'FileBufferedOrSeekable' handle.
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

-- | Collect column chunks per column (transposed across all row groups).
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

-- | Dispatch a column's chunks to the correct decoder path.
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

-- | Decode a required (non-nullable, non-repeated) column.
{-# 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)

    -- Decode a non-nullable BYTE_ARRAY (UTF-8) column straight into a single
    -- shared byte buffer + offsets ('PackedText'), instead of a boxed vector
    -- of per-row 'Text'. Each page's decoded 'Text' values (which share the
    -- chunk dictionary for dictionary-encoded pages) are appended by memcpy
    -- into one builder across all pages/chunks, then frozen once. This is the
    -- same representation the fast CSV reader uses and matches Arrow's string
    -- layout: no retained per-row 'Text' headers, no eager UTF-8 validation.
    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)

-- | Decode an optional (nullable) column.
{-# 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)

    -- Nullable BYTE_ARRAY (UTF-8): decode straight into a 'PackedText' (shared
    -- byte buffer + offsets + validity bitmap) via the text builder, walking
    -- def-levels to interleave nulls. Avoids the boxed @Vector Text@ the
    -- generic 'foldNullable' path would build.
    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)

-- | Decode a repeated (list/nested) column.
{-# 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)

-- Options application -----------------------------------------------------

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

-- Logical type conversion -------------------------------------------------

{- | Apply a column-description's logical type annotation to convert raw
decoded values (e.g. millisecond integers → 'UTCTime').
-}
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
        -- (ticks per second, picoseconds per tick) for each unit. Convert at
        -- native precision: a millisecond grid keeps ms, a nanosecond grid
        -- keeps ns. (The old code multiplied by @1_000_000 `div` divisor@,
        -- which truncated NANOS to 0 and collapsed every value to the epoch.)
        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

{- | Convert an epoch timestamp expressed as @ticksPerSecond@ ticks/second
(each tick = @psPerTick@ picoseconds) to 'UTCTime', at full precision.

Splits into days + within-day picoseconds with integer 'divMod' (which floors,
so the split is correct for pre-epoch negative values too), and never forms
picoseconds-since-epoch — that would overflow 'Int64' for modern dates — only
the bounded within-day picosecond count. 40587 is the Modified Julian Day of
the Unix epoch (1970-01-01).
-}
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)