Class: Karafka::Pro::Processing::ConsumerGroups::SubscriptionGroupsCoordinator
- Inherits:
-
Object
- Object
- Karafka::Pro::Processing::ConsumerGroups::SubscriptionGroupsCoordinator
- Includes:
- Singleton
- Defined in:
- lib/karafka/pro/processing/consumer_groups/subscription_groups_coordinator.rb
Overview
Uses the jobs queue API to lock (pause) and unlock (resume) operations of a given subscription group. It is abstracted away from jobs queue on this layer because we do not want to introduce jobs queue as a concept to the consumers layer
Instance Method Summary collapse
- #pause(subscription_group, lock_id = nil) ⇒ Object
- #resume(subscription_group, lock_id = nil) ⇒ Object
Instance Method Details
#pause(subscription_group, lock_id = nil) ⇒ Object
45 46 47 48 49 50 51 |
# File 'lib/karafka/pro/processing/consumer_groups/subscription_groups_coordinator.rb', line 45 def pause(subscription_group, lock_id = nil, **) jobs_queue.lock_async( subscription_group.id, lock_id, ** ) end |
#resume(subscription_group, lock_id = nil) ⇒ Object
56 57 58 |
# File 'lib/karafka/pro/processing/consumer_groups/subscription_groups_coordinator.rb', line 56 def resume(subscription_group, lock_id = nil) jobs_queue.unlock_async(subscription_group.id, lock_id) end |