Fix PulsarMQRunner shutdown.

This commit is contained in:
John Chilton
2019-04-19 12:44:53 -04:00
committed by Nate Coraor
parent b0f0fdc171
commit ca55cab671
2 changed files with 8 additions and 1 deletions
+2
View File
@@ -752,6 +752,8 @@ class PulsarMQJobRunner(PulsarJobRunner):
def _monitor(self):
# This is a message queue driven runner, don't monitor
# just setup required callback.
self._init_noop_monitor()
self.client_manager.ensure_has_status_update_callback(self.__async_update)
self.client_manager.ensure_has_ack_consumers()
+6 -1
View File
@@ -30,6 +30,10 @@ class Monitors(object):
self._start = start
register_postfork_function(self.start_monitoring)
def _init_noop_monitor(self):
self.sleeper = None
self.monitor_join = False
def start_monitoring(self):
if self._start:
self.monitor_thread.start()
@@ -42,7 +46,8 @@ class Monitors(object):
def shutdown_monitor(self):
self.stop_monitoring()
self.sleeper.wake()
if self.sleeper is not None:
self.sleeper.wake()
if self.monitor_join:
log.debug("Joining monitor thread")
self.monitor_thread.join(self.monitor_join_sleep)