hw-kafka-conduit 1.1.2 → 1.1.4
raw patch · 3 files changed
+61/−13 lines, 3 files
Files
- hw-kafka-conduit.cabal +13/−4
- src/Kafka/Conduit/Combinators.hs +20/−8
- src/Kafka/Conduit/Sink.hs +28/−1
hw-kafka-conduit.cabal view
@@ -1,8 +1,8 @@ name: hw-kafka-conduit-version: 1.1.2-synopsis: Conduit bindings for kafka-client-homepage: https://github.com/haskell-works/hw-kafka-client-conduit-bug-reports: https://github.com/haskell-works/hw-kafka-client-conduit/issues+version: 1.1.4+synopsis: Conduit bindings for hw-kafka-client+homepage: https://github.com/haskell-works/hw-kafka-conduit+bug-reports: https://github.com/haskell-works/hw-kafka-conduit/issues license: MIT license-file: LICENSE author: Alexey Raga <alexey.raga@gmail.com>@@ -14,6 +14,11 @@ cabal-version: >=1.10 description: Conduit bindings for hw-kafka-client +flag examples+ default: False+ manual: True+ description: Also compile examples+ source-repository head type: git location: https://github.com/haskell-works/hw-kafka-conduit@@ -29,6 +34,10 @@ , containers , resourcet , hw-kafka-conduit+ if flag(examples)+ buildable: True+ else+ buildable: False library hs-source-dirs: src
src/Kafka/Conduit/Combinators.hs view
@@ -1,14 +1,16 @@ module Kafka.Conduit.Combinators- ( batchByOrFlush+ ( BatchSize(..)+ , batchByOrFlush+ , batchByOrFlushEither , foldYield , throwLeft , throwLeftSatisfy ) where -import Control.Exception-import Control.Monad-import Control.Monad.Catch-import Data.Conduit+import Control.Exception+import Control.Monad+import Control.Monad.Catch+import Data.Conduit -- | Throws the left part of a value in a 'MonadThrow' context throwLeft :: (MonadThrow m, Exception e) => Conduit (Either e i) m i@@ -19,7 +21,7 @@ throwLeftSatisfy :: (MonadThrow m, Exception e) => (e -> Bool) -> Conduit (Either e i) m (Either e i) throwLeftSatisfy p = awaitForever awaitHandle where awaitHandle (Left e) | p e = throwM e- awaitHandle v = yield v+ awaitHandle v = yield v -- | Create a conduit that folds with the function f over its input i with its -- internal state s and emits outputs [o], then finally emits outputs [o] from@@ -34,10 +36,20 @@ foldYield f g s' Nothing -> forM_ (g s) yield -batchByOrFlush :: Monad m => Int -> Conduit (Maybe a) m [a]-batchByOrFlush n = foldYield folder finish (0 :: Int, [])+newtype BatchSize = BatchSize Int deriving (Show, Eq, Ord)++batchByOrFlush :: Monad m => BatchSize -> Conduit (Maybe a) m [a]+batchByOrFlush (BatchSize n) = foldYield folder finish (0 :: Int, []) where folder Nothing (_, xs) = ((0 , []), [reverse xs ]) folder (Just a) (i, xs) | (i + 1) >= n = ((0 , []), [reverse (a:xs)]) folder (Just a) (i, xs) = ((i + 1, a:xs), [])+ finish (_, xs) = [reverse xs]++batchByOrFlushEither :: Monad m => BatchSize -> Conduit (Either e a) m [a]+batchByOrFlushEither (BatchSize n) = foldYield folder finish (0 :: Int, [])+ where+ folder (Left _) (_, xs) = ((0 , []), [reverse xs ])+ folder (Right a) (i, xs) | i + 1 >= n = ((0 , []), [reverse (a:xs)])+ folder (Right a) (i, xs) = ((i + 1, a:xs ), []) finish (_, xs) = [reverse xs]
src/Kafka/Conduit/Sink.hs view
@@ -1,12 +1,14 @@ module Kafka.Conduit.Sink ( module X-, kafkaSink, kafkaSinkAutoClose, kafkaSinkNoClose+, kafkaSink, kafkaSinkAutoClose, kafkaSinkNoClose, kafkaBatchSinkNoClose+, commitOffsetsSink, flushThenCommitSink ) where import Control.Monad.IO.Class import Control.Monad (void) import Control.Monad.Trans.Resource import Data.Conduit+import Kafka.Consumer as X import Kafka.Producer as X import Kafka.Conduit.Combinators as X @@ -45,6 +47,23 @@ Nothing -> go Just err -> return (Just err) +-- | Creates a batching Sink for a given `KafkaProducer`.+-- The producer will NOT be closed automatically.+kafkaBatchSinkNoClose :: MonadIO m+ => KafkaProducer+ -> Sink [ProducerRecord] m [(ProducerRecord, KafkaError)]+kafkaBatchSinkNoClose prod = go+ where+ go = do+ mbMsg <- await+ case mbMsg of+ Nothing -> return []+ Just msgs -> do+ res <- produceMessageBatch prod msgs+ case res of+ [] -> go+ xs -> return xs+ -- | Creates a kafka producer for given properties and returns a Sink. -- -- This method of creating a Sink represents a simple case@@ -69,3 +88,11 @@ Just msg -> do res <- produceMessage prod msg runHandler $ maybe (Right prod) Left res++commitOffsetsSink :: MonadIO m => KafkaConsumer -> Sink i m ()+commitOffsetsSink = awaitForever . const . commitAllOffsets OffsetCommit++flushThenCommitSink :: MonadIO m => KafkaConsumer -> KafkaProducer -> Sink i m ()+flushThenCommitSink c p = awaitForever . const $ do+ flushProducer p+ commitAllOffsets OffsetCommit c