hw-kafka-client 2.4.0 → 2.4.1
raw patch · 3 files changed
+18/−7 lines, 3 files
Files
- hw-kafka-client.cabal +1/−1
- src/Kafka/Consumer.hs +1/−1
- src/Kafka/Consumer/Callbacks.hs +16/−5
hw-kafka-client.cabal view
@@ -1,5 +1,5 @@ name: hw-kafka-client-version: 2.4.0+version: 2.4.1 homepage: https://github.com/haskell-works/hw-kafka-client bug-reports: https://github.com/haskell-works/hw-kafka-client/issues license: MIT
src/Kafka/Consumer.hs view
@@ -90,7 +90,7 @@ -> Timeout -- ^ the timeout, in milliseconds -> m (Either KafkaError (ConsumerRecord (Maybe BS.ByteString) (Maybe BS.ByteString))) -- ^ Left on error or timeout, right for success pollMessage c@(KafkaConsumer _ (KafkaConf _ qr _)) (Timeout ms) = liftIO $ do--- pollConsumerEvents c Nothing+ pollConsumerEvents c Nothing mbq <- readIORef qr case mbq of Nothing -> return . Left $ KafkaBadSpecification "Messages queue is not configured, internal error, fatal."
src/Kafka/Consumer/Callbacks.hs view
@@ -35,7 +35,7 @@ where realCb k err pl = do k' <- newForeignPtr_ k- pls <- fromNativeTopicPartitionList' pl+ pls <- newForeignPtr_ pl setRebalanceCallback callback (KafkaConsumer (Kafka k') kc) (KafkaResponseError err) pls -- | Sets a callback that is called when rebalance is needed.@@ -67,17 +67,19 @@ setRebalanceCallback :: (KafkaConsumer -> RebalanceEvent -> IO ()) -> KafkaConsumer -> KafkaError- -> [TopicPartition] -> IO ()-setRebalanceCallback f k e ps =+ -> RdKafkaTopicPartitionListTPtr -> IO ()+setRebalanceCallback f k e pls = do+ ps <- fromNativeTopicPartitionList'' pls let assignment = (tpTopicName &&& tpPartition) <$> ps- in case e of++ case e of KafkaResponseError RdKafkaRespErrAssignPartitions -> do mbq <- getRdMsgQueue $ getKafkaConf k case mbq of Nothing -> pure () Just mq -> forM_ ps (\tp -> redirectPartitionQueue (getKafka k) (tpTopicName tp) (tpPartition tp) mq) f k (RebalanceBeforeAssign assignment)- void $ assign k ps+ void $ assign' k pls -- pass as pointer to avoid possible serialisation issues f k (RebalanceAssign assignment) KafkaResponseError RdKafkaRespErrRevokePartitions -> do f k (RebalanceBeforeRevoke assignment)@@ -94,3 +96,12 @@ else toNativeTopicPartitionList ps er = KafkaResponseError <$> (pl >>= rdKafkaAssign k) in kafkaErrorToMaybe <$> er++-- | Assigns specified partitions to a current consumer.+-- Assigning an empty list means unassigning from all partitions that are currently assigned.+assign' :: KafkaConsumer -> RdKafkaTopicPartitionListTPtr -> IO (Maybe KafkaError)+assign' (KafkaConsumer (Kafka k) _) pls = do+ let er = KafkaResponseError <$> rdKafkaAssign k pls+ m <- kafkaErrorToMaybe <$> er+ return m+