dataframe-fastcsv 1.4.1.0 → 1.4.1.1
raw patch · 12 files changed
+45/−86 lines, 12 filesdep ~dataframe-coredep ~dataframe-operationsPVP ok
version bump matches the API change (PVP)
Dependency ranges changed: dataframe-core, dataframe-operations
API changes (from Hackage documentation)
Files
- dataframe-fastcsv.cabal +5/−6
- src/DataFrame/IO/CSV/Fast.hs +0/−1
- src/DataFrame/IO/CSV/Fast/Columns.hs +1/−1
- src/DataFrame/IO/CSV/Fast/IndexPar.hs +1/−1
- src/DataFrame/IO/CSV/Fast/Parallel.hs +5/−5
- src/DataFrame/IO/CSV/Fast/Passes.hs +1/−1
- src/DataFrame/IO/CSV/Fast/TextMerge.hs +9/−5
- src/DataFrame/IO/CSV/Fast/Workers.hs +0/−59
- tests/Operations/Projection.hs +1/−1
- tests/Operations/ReadCsv.hs +13/−2
- tests/Operations/TypedExtraction.hs +5/−1
- tests/Properties/Csv.hs +4/−3
dataframe-fastcsv.cabal view
@@ -1,6 +1,6 @@ cabal-version: 3.4 name: dataframe-fastcsv-version: 1.4.1.0+version: 1.4.1.1 synopsis: SIMD-accelerated CSV reader for the dataframe library. description: A fast, SIMD-accelerated CSV/TSV reader using memory-mapped I/O@@ -53,13 +53,12 @@ DataFrame.IO.CSV.Fast.Passes DataFrame.IO.CSV.Fast.Slice DataFrame.IO.CSV.Fast.TextMerge- DataFrame.IO.CSV.Fast.Workers build-depends: base >= 4 && < 5, bytestring >= 0.11 && < 0.14, containers >= 0.6.7 && < 0.10,- dataframe-core >= 2.4 && < 2.5,+ dataframe-core >= 2.5 && < 2.6, dataframe-csv >= 2.3 && < 2.4,- dataframe-operations >= 2.4 && < 2.5,+ dataframe-operations >= 2.5 && < 2.6, dataframe-parsing >= 2.2 && < 2.3, dataframe-parsing >= 2.2 && < 2.3, mmap >= 0.5.8 && < 0.7,@@ -91,10 +90,10 @@ Properties.Csv build-depends: base >= 4 && < 5, containers >= 0.6.7 && < 0.10,- dataframe-core >= 2.4 && < 2.5,+ dataframe-core >= 2.5 && < 2.6, dataframe-csv >= 2.3 && < 2.4, dataframe-fastcsv,- dataframe-operations >= 2.4 && < 2.5,+ dataframe-operations >= 2.5 && < 2.6, dataframe-parsing >= 2.2 && < 2.3, directory >= 1.3.0.0 && < 2, HUnit >= 1.6 && < 1.8,
src/DataFrame/IO/CSV/Fast.hs view
@@ -42,7 +42,6 @@ CsvParseError (..), comma, getDelimiterIndices,- tab, ) import DataFrame.Internal.DataFrame (DataFrame) import DataFrame.Schema (Schema (..))
src/DataFrame/IO/CSV/Fast/Columns.hs view
@@ -44,7 +44,7 @@ import DataFrame.IO.CSV.Fast.Passes import DataFrame.IO.CSV.Fast.Slice (extractField) import DataFrame.Internal.Column (Column, ensureOptional, fromVector)-import DataFrame.Internal.DictEncode (dictCompactColumn)+import DataFrame.Internal.Column.Encode (dictCompactColumn) import DataFrame.Operations.Typing ( ParseOptions (..), ParsingAssumption (..),
src/DataFrame/IO/CSV/Fast/IndexPar.hs view
@@ -29,7 +29,7 @@ import Control.Concurrent (getNumCapabilities) import DataFrame.IO.CSV.Fast.Index (byteStringView, lf, quote, scalarScan)-import DataFrame.IO.CSV.Fast.Workers (pooledRun)+import DataFrame.Internal.Control.Concurrent (pooledRun) foreign import capi "process_csv.h count_delimiters_chunk" count_delimiters_chunk ::
src/DataFrame/IO/CSV/Fast/Parallel.hs view
@@ -4,7 +4,7 @@ into contiguous, newline-aligned chunks (alignment is by construction — chunks are row ranges over the quote-resolved delimiter index), each worker runs every column's 'ColumnPlan' over its chunk with right-sized-builders, and the merger splices the results ('mergeColumns' /+builders, and the merger splices the results ('concatColumns' / 'mergeTextChunksPar'). Inference classifies its sample once, globally, before fan-out; when chunks resolve different types the merger re-parses only the narrower-typed chunks.@@ -31,10 +31,10 @@ import DataFrame.IO.CSV.Fast.Columns import DataFrame.IO.CSV.Fast.Passes (NullSpec, PassCol (..)) import DataFrame.IO.CSV.Fast.TextMerge (mergeTextChunksPar)-import DataFrame.IO.CSV.Fast.Workers (pooledRun) import DataFrame.Internal.Column (Column, forceColumn)-import DataFrame.Internal.ColumnBuilder (mergeColumns)-import DataFrame.Internal.DictEncode (dictCompactColumn)+import DataFrame.Internal.Column.Builder (concatColumns)+import DataFrame.Internal.Column.Encode (dictCompactColumn)+import DataFrame.Internal.Control.Concurrent (pooledRun) import DataFrame.Operations.Typing (SafeReadMode) -- | Below this input size the fan-out overhead outweighs the parallelism.@@ -66,7 +66,7 @@ mergePassCols :: Int -> [PassCol] -> IO Column mergePassCols width ps@(PassText{} : _) = mergeTextChunksPar width [tc | PassText tc <- ps]-mergePassCols _ ps = pure $! mergeColumns [c | PassFull c <- ps]+mergePassCols _ ps = pure $! concatColumns [c | PassFull c <- ps] {- | Build every column, fanning the row range out over @nChunks@ chunks on a capability-wide worker pool (sequentially when @nChunks <= 1@). Same
src/DataFrame/IO/CSV/Fast/Passes.hs view
@@ -44,7 +44,7 @@ import DataFrame.IO.CSV.Fast.Index (quote) import DataFrame.IO.CSV.Fast.Slice import DataFrame.Internal.Column (Column, Columnable, fromVector)-import DataFrame.Internal.ColumnBuilder (+import DataFrame.Internal.Column.Builder ( ColumnBuilder (..), TextChunk, appendNum,
src/DataFrame/IO/CSV/Fast/TextMerge.hs view
@@ -18,15 +18,19 @@ import Control.Monad.ST (stToIO) import Data.Int (Int32)-import DataFrame.IO.CSV.Fast.Workers (pooledRun) import DataFrame.Internal.Column (Column (..))-import DataFrame.Internal.ColumnMerge (+import DataFrame.Internal.Column.Bitmap (Validity (Validity))+import DataFrame.Internal.Column.Merge ( TextChunk (..),+ concatValidity, mergeTextChunks,- spliceBitmaps, tcRows, )-import DataFrame.Internal.PackedText (mkPackedContiguous, mkPackedContiguous32)+import DataFrame.Internal.Control.Concurrent (pooledRun)+import DataFrame.Internal.Data.PackedText (+ mkPackedContiguous,+ mkPackedContiguous32,+ ) {- | Merge text chunks with @width@-way parallel byte copies + offset rebase, then wrap the shared buffer as 'PackedText'. Single chunks take@@ -73,5 +77,5 @@ mkPackedContiguous <$> stToIO (A.unsafeFreeze marr) <*> VU.unsafeFreeze offsMV- let !bm = spliceBitmaps [(tcBitmap c, tcRows c) | c <- cs]+ let !bm = concatValidity [Validity (tcBitmap c) (tcRows c) | c <- cs] pure (PackedText bm packed)
− src/DataFrame/IO/CSV/Fast/Workers.hs
@@ -1,59 +0,0 @@-{- | Thread fan-out primitives shared by the parallel scan, the chunk-parse and the parallel merge: plain 'forkIO' workers joined through-'MVar's (no sparks), with a counter-based pool for finer-grained chunks.--}-module DataFrame.IO.CSV.Fast.Workers (- forkJoin,- pooledRun,-) where--import qualified Data.Vector as V-import qualified Data.Vector.Mutable as VM--import Control.Concurrent (forkIO)-import Control.Concurrent.MVar (newEmptyMVar, putMVar, takeMVar)-import Control.Exception (ErrorCall (..), SomeException, throwIO, try)-import Control.Monad (when)-import Data.IORef (atomicModifyIORef', newIORef)---- | Run each action in its own thread; rethrow the first failure in order.-forkJoin :: [IO a] -> IO [a]-forkJoin actions = do- vars <- mapM spawn actions- results <- mapM takeMVar vars- either (throwIO :: SomeException -> IO [a]) pure (sequence results)- where- spawn act = do- var <- newEmptyMVar- _ <- forkIO (try act >>= putMVar var)- pure var--{- | Run the actions on a pool of @width@ threads (work-stealing via a-shared counter), so finer-grained chunks balance load without running-every chunk's builders concurrently. Results keep their input order.--Each action slot is cleared before the action runs, so data captured by a-completed closure is unreachable as soon as it finishes — callers rely on-this to release per-column chunk payloads during the parallel merge.--}-pooledRun :: Int -> [IO a] -> IO [a]-pooledRun width actions- | width >= n = forkJoin actions- | otherwise = do- next <- newIORef 0- out <- VM.unsafeNew n- acts <- VM.unsafeNew n- sequence_ [VM.unsafeWrite acts i a | (i, a) <- zip [0 ..] actions]- let worker = do- i <- atomicModifyIORef' next (\j -> (j + 1, j))- when (i < n) $ do- act <- VM.unsafeRead acts i- VM.unsafeWrite acts i consumed- r <- act- VM.write out i r- worker- _ <- forkJoin (replicate width worker)- V.toList <$> V.freeze out- where- n = length actions- consumed = throwIO (ErrorCall "pooledRun: slot already consumed")
tests/Operations/Projection.hs view
@@ -10,7 +10,7 @@ import qualified Data.Map as M import qualified Data.Text as T-import qualified Data.Text.IO as TIO+import qualified Data.Text.IO.Utf8 as TIO import Control.Exception (SomeException, evaluate, try) import Data.List (isInfixOf)
tests/Operations/ReadCsv.hs view
@@ -26,8 +26,9 @@ ReadOptions (..), defaultReadOptions, )-import DataFrame.Internal.Column (Column (..), bitmapTestBit)+import DataFrame.Internal.Column (Column (..)) import qualified DataFrame.Internal.Column as DI+import DataFrame.Internal.Column.Bitmap (bitmapTestBit) import DataFrame.Internal.DataFrame ( DataFrame (..), columnIndices,@@ -37,7 +38,14 @@ ) import DataFrame.Schema (Schema (..), SchemaType (..)) import System.Directory (removeFile)-import System.IO (IOMode (..), withFile)+import System.IO (+ IOMode (..),+ hSetEncoding,+ hSetNewlineMode,+ noNewlineTranslation,+ utf8,+ withFile,+ ) import Test.HUnit import Type.Reflection (typeRep) @@ -59,6 +67,8 @@ prettyPrintSeparated :: Char -> FilePath -> DataFrame -> IO () prettyPrintSeparated sep filepath df = withFile filepath WriteMode $ \handle -> do+ hSetEncoding handle utf8+ hSetNewlineMode handle noNewlineTranslation let (rows, _) = dataframeDimensions df let headers = map fst (L.sortBy (compare `on` snd) (M.toList (columnIndices df))) TIO.hPutStrLn@@ -85,6 +95,7 @@ where go :: Int -> Column -> [T.Text] -> [T.Text] go idx c@(PackedText _ _) acc = go idx (DI.materializePacked c) acc+ go idx c@(MergedColumn _ _) acc = go idx (DI.materializeMerged c) acc go _ (BoxedColumn bm (c :: V.Vector a)) acc = case c V.!? i of Just e -> escapeField sep textRep : acc where
tests/Operations/TypedExtraction.hs view
@@ -12,7 +12,11 @@ import qualified Data.Proxy as P import qualified Data.Text as T import qualified Data.Text.Encoding as TE-import qualified Data.Text.IO as TIO++-- UTF-8 byte-mode IO: the plain Data.Text.IO writer honours the+-- handle's text mode, which on Windows turns \n into \r\n and+-- corrupts inputs meant for byte-level parsers.+import qualified Data.Text.IO.Utf8 as TIO import Control.Exception (ErrorCall, evaluate, try) import Data.Time (Day)
tests/Properties/Csv.hs view
@@ -20,7 +20,7 @@ import qualified Data.Map as M import qualified Data.Text as T import qualified Data.Text.Encoding as TE-import qualified Data.Text.IO as TIO+import qualified Data.Text.IO.Utf8 as TIO import qualified Data.Vector as V import DataFrame.IO.CSV (defaultReadOptions)@@ -124,13 +124,14 @@ Nothing -> Nothing UnboxedColumn{} -> Nothing PackedText{} -> Nothing+ MergedColumn{} -> Nothing {- | Run a property in @IO@ against a generated CSV text, cleaning up the temp file afterwards no matter what. -} withCsvFile :: String -> T.Text -> (FilePath -> IO a) -> IO a withCsvFile label body action = do- let path = "/tmp/fastcsv_prop_" <> label <> ".csv"+ let path = "./tests/data/unstable_csv/fastcsv_prop_" <> label <> ".csv" TIO.writeFile path body r <- action path removeFile path@@ -209,7 +210,7 @@ T.intercalate "," (map (encodeCell ',' . unCell) cells) csv = "v\n" <> plainRow <> ",\"dangling\n" result <- run $ do- let path = "/tmp/fastcsv_prop_unclosed.csv"+ let path = "./tests/data/unstable_csv/fastcsv_prop_unclosed.csv" TIO.writeFile path csv r <- try @CsvParseError (D.fastReadCsv path) removeFile path