Fix _mp_consume queue variable name conflict
This commit is contained in:
@@ -25,7 +25,7 @@ log = logging.getLogger(__name__)
|
|||||||
Events = namedtuple("Events", ["start", "pause", "exit"])
|
Events = namedtuple("Events", ["start", "pause", "exit"])
|
||||||
|
|
||||||
|
|
||||||
def _mp_consume(client, group, topic, queue, size, events, **consumer_options):
|
def _mp_consume(client, group, topic, message_queue, size, events, **consumer_options):
|
||||||
"""
|
"""
|
||||||
A child process worker which consumes messages based on the
|
A child process worker which consumes messages based on the
|
||||||
notifications given by the controller process
|
notifications given by the controller process
|
||||||
@@ -69,7 +69,7 @@ def _mp_consume(client, group, topic, queue, size, events, **consumer_options):
|
|||||||
if message:
|
if message:
|
||||||
while True:
|
while True:
|
||||||
try:
|
try:
|
||||||
queue.put(message, timeout=FULL_QUEUE_WAIT_TIME_SECONDS)
|
message_queue.put(message, timeout=FULL_QUEUE_WAIT_TIME_SECONDS)
|
||||||
break
|
break
|
||||||
except queue.Full:
|
except queue.Full:
|
||||||
if events.exit.is_set(): break
|
if events.exit.is_set(): break
|
||||||
|
|||||||
Reference in New Issue
Block a user