massiv 0.3.3.0 → 0.3.4.0
raw patch · 10 files changed
+275/−21 lines, 10 filesdep +mersenne-random-pure64dep +splitmixdep ~schedulerPVP: major bump suggested
API removals or changes: PVP suggests a major version bump
Dependencies added: mersenne-random-pure64, splitmix
Dependency ranges changed: scheduler
API changes (from Hackage documentation)
+ Data.Massiv.Array: forWS :: forall r ix b r' a s m. (Source r' ix a, Mutable r ix b, MonadUnliftIO m, PrimMonad m) => WorkerStates s -> Array r' ix a -> (a -> s -> m b) -> m (Array r ix b)
+ Data.Massiv.Array: iforWS :: forall r ix b r' a s m. (Source r' ix a, Mutable r ix b, MonadUnliftIO m, PrimMonad m) => WorkerStates s -> Array r' ix a -> (ix -> a -> s -> m b) -> m (Array r ix b)
+ Data.Massiv.Array: imapWS :: forall r ix b r' a s m. (Source r' ix a, Mutable r ix b, MonadUnliftIO m, PrimMonad m) => WorkerStates s -> (ix -> a -> s -> m b) -> Array r' ix a -> m (Array r ix b)
+ Data.Massiv.Array: mapWS :: forall r ix b r' a s m. (Source r' ix a, Mutable r ix b, MonadUnliftIO m, PrimMonad m) => WorkerStates s -> (a -> s -> m b) -> Array r' ix a -> m (Array r ix b)
+ Data.Massiv.Array: randomArrayS :: forall r ix e g. Mutable r ix e => g -> Sz ix -> (g -> (e, g)) -> (g, Array r ix e)
+ Data.Massiv.Array: randomArrayWS :: forall r ix e g m. (Mutable r ix e, MonadUnliftIO m, PrimMonad m) => WorkerStates g -> Sz ix -> (g -> m e) -> m (Array r ix e)
+ Data.Massiv.Array.Mutable: generateArrayLinearWS :: forall r ix e s m. (Mutable r ix e, MonadUnliftIO m, PrimMonad m) => WorkerStates s -> Sz ix -> (Int -> s -> m e) -> m (Array r ix e)
+ Data.Massiv.Array.Mutable: generateArrayWS :: forall r ix e s m. (Mutable r ix e, MonadUnliftIO m, PrimMonad m) => WorkerStates s -> Sz ix -> (ix -> s -> m e) -> m (Array r ix e)
+ Data.Massiv.Core: data WorkerStates s
+ Data.Massiv.Core: initWorkerStates :: MonadIO m => Comp -> (WorkerId -> m s) -> m (WorkerStates s)
+ Data.Massiv.Core.Index: splitLinearlyWithStatefulM_ :: Monad m => SchedulerWS s m () -> Int -> (Int -> s -> m b) -> (Int -> b -> m c) -> m ()
- Data.Massiv.Array: randomArray :: Index ix => g -> (g -> (g, g)) -> (g -> (e, g)) -> Comp -> Sz ix -> Array DL ix e
+ Data.Massiv.Array: randomArray :: forall ix e g. Index ix => g -> (g -> (g, g)) -> (g -> (e, g)) -> Comp -> Sz ix -> Array DL ix e
- Data.Massiv.Core: pattern Par :: () => () => Comp
+ Data.Massiv.Core: pattern Par' :: () => () => Comp
Files
- CHANGELOG.md +11/−0
- massiv.cabal +6/−3
- src/Data/Massiv/Array/Delayed/Windowed.hs +4/−4
- src/Data/Massiv/Array/Mutable.hs +38/−0
- src/Data/Massiv/Array/Ops/Construct.hs +105/−8
- src/Data/Massiv/Array/Ops/Map.hs +60/−1
- src/Data/Massiv/Core.hs +4/−1
- src/Data/Massiv/Core/Common.hs +3/−1
- src/Data/Massiv/Core/Iterator.hs +26/−1
- tests/Data/Massiv/Array/Ops/MapSpec.hs +18/−2
CHANGELOG.md view
@@ -1,3 +1,14 @@+# 0.3.4++* Use the the new stateful workers feature of `schdeuler-1.4.0`+* Addition of:+ * `randomArrayS`+ * `randomArrayWS`+ * `generateArrayWS`+ * `generateArrayLinearWS`+ * `mapWS`, `forWS`, `imapWS` and `iforWS`+ * and `splitLinearlyWithStatefulM_`+ # 0.3.3 * Fix type signature for `createArray`.
massiv.cabal view
@@ -1,5 +1,5 @@ name: massiv-version: 0.3.3.0+version: 0.3.4.0 synopsis: Massiv (Массив) is an Array Library. description: Multi-dimensional Arrays with fusion, stencils and parallel computation. homepage: https://github.com/lehins/massiv@@ -75,7 +75,7 @@ , data-default-class , deepseq , exceptions- , scheduler >= 1.1.0+ , scheduler >= 1.4.0 , primitive , unliftio-core , vector@@ -122,6 +122,7 @@ , deepseq , massiv , hspec+ , scheduler , QuickCheck , unliftio , vector@@ -143,8 +144,10 @@ , doctest >=0.15 , QuickCheck , massiv- , template-haskell+ , mersenne-random-pure64 , random+ , splitmix >= 0.0.1+ , template-haskell default-language: Haskell2010 source-repository head
src/Data/Massiv/Array/Delayed/Windowed.hs view
@@ -332,13 +332,13 @@ loadWindowIx2 :: Monad m => Int -> (Ix2 -> m ()) -> Ix2 -> m ()-loadWindowIx2 numWorkers loadWindow (it :. ib) = do- let !(chunkHeight, slackHeight) = (ib - it) `quotRem` numWorkers- loopM_ 0 (< numWorkers) (+ 1) $ \ !wid ->+loadWindowIx2 nWorkers loadWindow (it :. ib) = do+ let !(chunkHeight, slackHeight) = (ib - it) `quotRem` nWorkers+ loopM_ 0 (< nWorkers) (+ 1) $ \ !wid -> let !it' = wid * chunkHeight + it in loadWindow (it' :. (it' + chunkHeight)) when (slackHeight > 0) $- let !itSlack = numWorkers * chunkHeight + it+ let !itSlack = nWorkers * chunkHeight + it in loadWindow (itSlack :. (itSlack + slackHeight)) {-# INLINE loadWindowIx2 #-}
src/Data/Massiv/Array/Mutable.hs view
@@ -48,6 +48,9 @@ , generateArrayLinear , generateArrayS , generateArrayLinearS+ -- *** Stateful worker threads+ , generateArrayWS+ , generateArrayLinearWS -- *** Unfold , unfoldrPrimM_ , iunfoldrPrimM_@@ -505,6 +508,41 @@ -> m (Array r ix e) generateArrayLinear comp sz f = makeMArrayLinear comp sz f >>= unsafeFreeze comp {-# INLINE generateArrayLinear #-}+++-- | Same as `generateArrayWS`, but use linear indexing instead.+--+-- @since 0.3.4+generateArrayLinearWS ::+ forall r ix e s m. (Mutable r ix e, MonadUnliftIO m, PrimMonad m)+ => WorkerStates s+ -> Sz ix+ -> (Int -> s -> m e)+ -> m (Array r ix e)+generateArrayLinearWS states sz make = do+ marr <- unsafeNew sz+ withSchedulerWS_ states $ \schedulerWS ->+ splitLinearlyWithStatefulM_+ schedulerWS+ (totalElem sz)+ make+ (unsafeLinearWrite marr)+ unsafeFreeze (workerStatesComp states) marr+{-# INLINE generateArrayLinearWS #-}++-- | Use per worker thread state while generating elements of the array. Very useful for+-- things that are not thread safe.+--+-- @since 0.3.4+generateArrayWS ::+ forall r ix e s m. (Mutable r ix e, MonadUnliftIO m, PrimMonad m)+ => WorkerStates s+ -> Sz ix+ -> (ix -> s -> m e)+ -> m (Array r ix e)+generateArrayWS states sz make =+ generateArrayLinearWS states sz (\ix -> make (fromLinearIndex sz ix))+{-# INLINE generateArrayWS #-} -- | Sequentially unfold an array from the left.
src/Data/Massiv/Array/Ops/Construct.hs view
@@ -36,6 +36,8 @@ , iunfoldrS_ -- *** Random , randomArray+ , randomArrayS+ , randomArrayWS -- *** Applicative , makeArrayA , makeArrayAR@@ -65,6 +67,7 @@ import Control.Monad.ST import Data.Massiv.Array.Delayed.Pull import Data.Massiv.Array.Delayed.Push+import Data.Massiv.Array.Mutable import Data.Massiv.Core.Common import Prelude as P hiding (enumFromTo, replicate) @@ -246,16 +249,26 @@ -- | Create an array with random values by using a pure splittable random number generator--- such as the one provided by [random](https://www.stackage.org/package/random) or--- [splitmix](https://www.stackage.org/package/splitmix) packages.+-- such as one provided by either [splitmix](https://www.stackage.org/package/splitmix) or+-- [random](https://www.stackage.org/package/random) packages. If you don't have a+-- splittable generator consider using `randomArrayS` or `randomArrayIO` instead. -- -- Because of the pure nature of the generator and its splitability we are not only able -- to parallelize the random value generation, but also guarantee that it will be--- deterministic, granted none of the arguments are changed.+-- deterministic, granted none of the arguments have changed. -- -- ==== __Examples__ -- -- >>> import Data.Massiv.Array+-- >>> import System.Random.SplitMix as SplitMix+-- >>> gen = SplitMix.mkSMGen 217+-- >>> randomArray gen SplitMix.splitSMGen SplitMix.nextDouble (ParN 2) (Sz2 2 3) :: Array DL Ix2 Double+-- Array DL (ParN 2) (Sz (2 :. 3))+-- [ [ 0.7383156058619669, 0.39904053166835896, 0.5617584038393628 ]+-- , [ 0.7218718218678238, 0.7006722805067258, 0.7225894731396042 ]+-- ]+--+-- >>> import Data.Massiv.Array -- >>> import System.Random as System -- >>> gen = System.mkStdGen 217 -- >>> randomArray gen System.split System.random (ParN 2) (Sz2 2 3) :: Array DL Ix2 Double@@ -266,11 +279,11 @@ -- -- @since 0.3.3 randomArray ::- Index ix+ forall ix e g. Index ix => g -- ^ Initial random value generator -> (g -> (g, g)) -- ^ A function that can split a generator in two independent generators- -> (g -> (e, g))+ -> (g -> (e, g)) -- TODO: move this argument after the sz, to keep ot consistent. -- ^ A function that produces a random value and the next generator -> Comp -- ^ Computation strategy. -> Sz ix -- ^ Resulting size of the array.@@ -285,18 +298,102 @@ pure genII' genForSlack <- loopM startAt (< slackStartAt) (+ chunkLength) gen $ \start genI -> do- let (genI0, genI1) = splitGen genI+ let (genI0, genI1) =+ if numWorkers scheduler == 1+ then (genI, genI)+ else splitGen genI scheduleWork_ scheduler $ void $ loopM start (< (start + chunkLength)) (+ 1) genI0 writeRandom pure genI1 when (slackStartAt < totalLength + startAt) $ scheduleWork_ scheduler $- void $- loopM slackStartAt (< totalLength + startAt) (+ 1) genForSlack writeRandom+ void $ loopM slackStartAt (< totalLength + startAt) (+ 1) genForSlack writeRandom where !totalLength = totalElem sz {-# INLINE randomArray #-} +-- | Similar to `randomArray` but performs generation sequentially, which means it doesn't+-- require splitability property. Another consequence is that it returns the new generator+-- together with /manifest/ array of random values.+--+-- ==== __Examples__+--+-- >>> import Data.Massiv.Array+-- >>> import System.Random.SplitMix as SplitMix+-- >>> gen = SplitMix.mkSMGen 217+-- >>> snd $ randomArrayS gen (Sz2 2 3) SplitMix.nextDouble :: Array P Ix2 Double+-- Array P Seq (Sz (2 :. 3))+-- [ [ 0.8878273949359751, 0.11290807610140963, 0.7383156058619669 ]+-- , [ 0.39904053166835896, 0.5617584038393628, 0.16248374266020216 ]+-- ]+--+-- >>> import Data.Massiv.Array+-- >>> import System.Random.Mersenne.Pure64 as MT+-- >>> gen = MT.pureMT 217+-- >>> snd $ randomArrayS gen (Sz2 2 3) MT.randomDouble :: Array P Ix2 Double+-- Array P Seq (Sz (2 :. 3))+-- [ [ 0.5504018416543631, 0.22504666452851707, 0.4480480867867128 ]+-- , [ 0.7139711572975297, 0.49401087853770953, 0.9397201599368645 ]+-- ]+--+-- >>> import Data.Massiv.Array+-- >>> import System.Random as System+-- >>> gen = System.mkStdGen 217+-- >>> snd $ randomArrayS gen (Sz2 2 3) System.random :: Array P Ix2 Double+-- Array P Seq (Sz (2 :. 3))+-- [ [ 0.7972230393466304, 0.4485860543300083, 0.257773196880671 ]+-- , [ 0.19115043859955794, 0.33784788936970034, 3.479381605706322e-2 ]+-- ]+--+-- @since 0.3.4+randomArrayS ::+ forall r ix e g. Mutable r ix e+ => g -- ^ Initial random value generator+ -> Sz ix -- ^ Resulting size of the array.+ -> (g -> (e, g))+ -- ^ A function that produces a random value and the next generator+ -> (g, Array r ix e)+randomArrayS gen sz nextRandom =+ runST $ unfoldrPrimM Seq sz (pure . nextRandom) gen+{-# INLINE randomArrayS #-}++-- | This is a stateful approach of generating random values. If your generator is pure+-- and splittable, it is better to use `randomArray` instead, which will give you a pure,+-- deterministic and parallelizable generation of arrays. On the other hand, if your+-- generator is not thread safe, which is most likely the case, instead of using some sort+-- of global mutex, `WorkerStates` allows you to keep track of individual state per worker+-- (thread), which fits parallelization of random value generation perfectly. All that+-- needs to be done is generators need to be initialized once per worker and then they can+-- be reused as many times as necessary.+--+-- ==== __Examples__+--+-- In the example below we take a stateful random generator from+-- [wmc-random](https://www.stackage.org/package/mwc-random), which is not thread safe,+-- and safely parallelize it by giving each thread it's own generator:+--+-- > λ> import Data.Massiv.Array+-- > λ> import System.Random.MWC (createSystemRandom, uniformR)+-- > λ> import System.Random.MWC.Distributions (standard)+-- > λ> gens <- initWorkerStates Par (\_ -> createSystemRandom)+-- > λ> randomArrayWS gens (Sz2 2 3) standard :: IO (Array P Ix2 Double)+-- > Array P Par (Sz (2 :. 3))+-- > [ [ -0.9066144845415213, 0.5264323240310042, -1.320943607597422 ]+-- > , [ -0.6837929005619592, -0.3041255565826211, 6.53353089112833e-2 ]+-- > ]+-- > λ> randomArrayWS gens (Sz1 10) (uniformR (0, 9)) :: IO (Array P Ix1 Int)+-- > Array P Par (Sz1 10)+-- > [ 3, 6, 1, 2, 1, 7, 6, 0, 8, 8 ]+--+-- @since 0.3.4+randomArrayWS ::+ forall r ix e g m. (Mutable r ix e, MonadUnliftIO m, PrimMonad m)+ => WorkerStates g -- ^ Use `initWorkerStates` to initialize you per thread generators+ -> Sz ix -- ^ Resulting size of the array+ -> (g -> m e) -- ^ Generate the value using the per thread generator.+ -> m (Array r ix e)+randomArrayWS states sz genRandom = generateArrayLinearWS states sz (const genRandom)+{-# INLINE randomArrayWS #-} infix 4 ..., ..:
src/Data/Massiv/Array/Ops/Map.hs view
@@ -44,12 +44,16 @@ , iforM_ -- *** Parallelizable , mapIO+ , mapWS , mapIO_ , imapIO+ , imapWS , imapIO_ , forIO+ , forWS , forIO_ , iforIO+ , iforWS , iforIO_ , imapSchedulerM_ , iforSchedulerM_@@ -579,6 +583,7 @@ (\i -> void . action (fromLinearIndex sz i)) {-# INLINE imapSchedulerM_ #-} + -- | Same as `imapM_`, but will use the supplied scheduler. -- -- @since 0.3.1@@ -588,7 +593,6 @@ {-# INLINE iforSchedulerM_ #-} - -- | Same as `mapIO` but map an index aware action instead. -- -- @since 0.2.6@@ -610,6 +614,61 @@ -> m (Array r ix b) forIO = flip mapIO {-# INLINE forIO #-}++++-- | Same as `imapIO`, but ignores the inner computation strategy and uses stateful+-- workers during computation instead. Use `initWorkerStates` for the `WorkerStates`+-- initialization.+--+-- @since 0.3.4+imapWS ::+ forall r ix b r' a s m. (Source r' ix a, Mutable r ix b, MonadUnliftIO m, PrimMonad m)+ => WorkerStates s+ -> (ix -> a -> s -> m b)+ -> Array r' ix a+ -> m (Array r ix b)+imapWS states f arr = generateArrayWS states (size arr) (\ix s -> f ix (unsafeIndex arr ix) s)+{-# INLINE imapWS #-}++-- | Same as `imapWS`, but without the index.+--+-- @since 0.3.4+mapWS ::+ forall r ix b r' a s m. (Source r' ix a, Mutable r ix b, MonadUnliftIO m, PrimMonad m)+ => WorkerStates s+ -> (a -> s -> m b)+ -> Array r' ix a+ -> m (Array r ix b)+mapWS states f = imapWS states (\ _ -> f)+{-# INLINE mapWS #-}+++-- | Same as `imapWS`, but with source array and mapping action arguments flipped.+--+-- @since 0.3.4+iforWS ::+ forall r ix b r' a s m. (Source r' ix a, Mutable r ix b, MonadUnliftIO m, PrimMonad m)+ => WorkerStates s+ -> Array r' ix a+ -> (ix -> a -> s -> m b)+ -> m (Array r ix b)+iforWS states f arr = imapWS states arr f+{-# INLINE iforWS #-}++-- | Same as `iforWS`, but without the index.+--+-- @since 0.3.4+forWS ::+ forall r ix b r' a s m. (Source r' ix a, Mutable r ix b, MonadUnliftIO m, PrimMonad m)+ => WorkerStates s+ -> Array r' ix a+ -> (a -> s -> m b)+ -> m (Array r ix b)+forWS states arr f = imapWS states (\ _ -> f) arr+{-# INLINE forWS #-}++ -- | Same as `mapIO_` but with arguments flipped. --
src/Data/Massiv/Core.hs view
@@ -27,7 +27,9 @@ , L(..) , LN , ListItem- , Comp(Seq, Par, ParOn, ParN)+ , Comp(Seq, Par, Par', ParOn, ParN)+ , WorkerStates+ , initWorkerStates , module Data.Massiv.Core.Index -- * Exceptions , MonadThrow(..)@@ -45,6 +47,7 @@ ) where import Control.Exception (Exception(..), SomeException)+import Control.Scheduler (WorkerStates, initWorkerStates) import Data.Massiv.Core.Common import Data.Massiv.Core.Index import Data.Massiv.Core.List
src/Data/Massiv/Core/Common.hs view
@@ -33,6 +33,7 @@ , numWorkers , scheduleWork , scheduleWork_+ , WorkerStates , unsafeRead , unsafeWrite , unsafeLinearModify@@ -81,7 +82,8 @@ import Control.Monad.Catch (MonadThrow(..)) import Control.Monad.IO.Unlift (MonadIO(liftIO), MonadUnliftIO) import Control.Monad.Primitive-import Control.Scheduler (Comp(..), Scheduler, numWorkers, scheduleWork, scheduleWork_)+import Control.Scheduler (Comp(..), Scheduler, WorkerStates, numWorkers,+ scheduleWork, scheduleWork_) import Data.Massiv.Core.Exception import Data.Massiv.Core.Index import Data.Typeable
src/Data/Massiv/Core/Iterator.hs view
@@ -17,9 +17,10 @@ , splitLinearlyWith_ , splitLinearlyWithM_ , splitLinearlyWithStartAtM_+ , splitLinearlyWithStatefulM_ ) where -import Control.Scheduler (Scheduler, numWorkers, scheduleWork_)+import Control.Scheduler -- | Efficient loop with an accumulator --@@ -146,3 +147,27 @@ scheduleWork_ scheduler $ loopM_ (slackStart + startAt) (< (totalLength + startAt)) (+ 1) $ \ !k -> make k >>= write k {-# INLINE splitLinearlyWithStartAtM_ #-}++++-- | Interator that can be used to split computation jobs, while using a stateful scheduler.+--+-- @since 0.3.4+splitLinearlyWithStatefulM_ ::+ Monad m+ => SchedulerWS s m ()+ -> Int -- ^ Total linear length+ -> (Int -> s -> m b) -- ^ Element producing action+ -> (Int -> b -> m c) -- ^ Element storing action+ -> m ()+splitLinearlyWithStatefulM_ schedulerWS totalLength make store =+ let nWorkers = numWorkers (unwrapSchedulerWS schedulerWS)+ in splitLinearly nWorkers totalLength $ \chunkLength slackStart -> do+ loopM_ 0 (< slackStart) (+ chunkLength) $ \ !start ->+ scheduleWorkState_ schedulerWS $ \s ->+ loopM_ start (< (start + chunkLength)) (+ 1) $ \ !k ->+ make k s >>= store k+ scheduleWorkState_ schedulerWS $ \s ->+ loopM_ slackStart (< totalLength) (+ 1) $ \ !k ->+ make k s >>= store k+{-# INLINE splitLinearlyWithStatefulM_ #-}
tests/Data/Massiv/Array/Ops/MapSpec.hs view
@@ -5,13 +5,13 @@ {-# LANGUAGE TypeApplications #-} module Data.Massiv.Array.Ops.MapSpec (spec) where +import Data.IORef import Control.Monad.ST import Data.Foldable as F import Data.Massiv.Array.Unsafe import Data.Massiv.CoreArbitrary as A import Prelude as P--+import Control.Scheduler.Internal prop_zipUnzip :: (Index ix, Show (Array D ix Int))@@ -82,6 +82,8 @@ it "zipFlip3" $ property $ prop_zipFlip3 @ix describe "Traversing" $ do it "itraverseA" $ property $ prop_itraverseA @ix+ describe "StatefulMapping" $ do+ it "mapWS" $ property $ prop_MapWS @ix spec :: Spec spec = do@@ -107,3 +109,17 @@ zipWithIndex :: forall r ix e . Source r ix e => Array r ix e -> Array D ix (ix, e) zipWithIndex arr = A.zip (range Seq zeroIndex (unSz (size arr))) arr {-# INLINE zipWithIndex #-}+++prop_MapWS :: (Show (Array U ix Int), Index ix) => Array U ix Int -> Property+prop_MapWS arr =+ monadicIO $+ run $ do+ states <- initWorkerStates (getComp arr) (\_ -> newIORef 0)+ arr' <-+ forWS states arr $ \e ref -> do+ acc <- readIORef ref+ writeIORef ref (acc + e)+ pure e+ accsArr <- A.mapM @P readIORef (evalArray Seq (_workerStatesArray states))+ pure (A.sum arr' === A.sum accsArr .&&. arr === arr')