It is first time I play with parallel computing seriously.
I am using multiprocessing module in python and I am running into this problem:
A queue consumer run in a different process then queue producer, the former should wait the latter to finish its job, before stop iterating over the queue. Sometimes the consumer is faster then producer and the queue stays empty. If I don't put any condition the program won't stop.
In the sample code I use the wildcard PRODUCER_IS_OVER to example what I need.
following code sketch the problem:
def save_data(save_que, file_):
### Coroutine instantiation
PRODUCER_IS_OVER = False
empty = False
### Queue consumer
while not(empty and PRODUCER_IS_OVER):
try:
data = save_que.get()
print("saving data",data)
except:
empty = save_que.empty()
print(empty)
pass
#PRODUCER_IS_OVER = get_condition()
print ("All data saved")
return
def get_condition():
###NameError: global name 'PRODUCER_IS_OVER' is not defined
if PRODUCER_IS_OVER:
return True
else:
return False
def produce_data(save_que):
for _ in range(5):
time.sleep(random.randint(1,5))
data = random.randint(1,10)
print("sending data", data)
save_que.put(data)
### Main function here
import random
import time
from multiprocessing import Queue, Manager, Process
manager = Manager()
save_que = manager.Queue()
file_ = "file"
save_p = Process(target= save_data, args=(save_que, file_))
save_p.start()
PRODUCER_IS_OVER = False
produce_data(save_que)
PRODUCER_IS_OVER = True
save_p.join()
produce_data takes variable time and I want the save_p process to start BEFORE populate the queue, in order to consume the queue while is filled.
I think there are workaround to communicate when to stop iteration, but I want to know whether exist a proper way to do it.
I tried both multiprocessing.Pipe and .Lock, but I don't know how implement correctly and efficiently.
SOLVED: is it the best way?
following code implement STOPMESSAGE in the Q, works fine, I can refine it with a class, QMsg, in case the language supports only static types.
def save_data(save_que, file_):
# Coroutine instantiation
PRODUCER_IS_OVER = False
empty = False
# Queue consumer
while not(empty and PRODUCER_IS_OVER):
data = save_que.get()
empty = save_que.empty()
print("saving data", data)
if data == "STOP":
PRODUCER_IS_OVER = True
print("All data saved")
return
def get_condition():
# NameError: global name 'PRODUCER_IS_OVER' is not defined
if PRODUCER_IS_OVER:
return True
else:
return False
def produce_data(save_que):
for _ in range(5):
time.sleep(random.randint(1, 5))
data = random.randint(1, 10)
print("sending data", data)
save_que.put(data)
save_que.put("STOP")
# Main function here
import random
import time
from multiprocessing import Queue, Manager, Process
manager = Manager()
save_que = manager.Queue()
file_ = "file"
save_p = Process(target=save_data, args=(save_que, file_))
save_p.start()
PRODUCER_IS_OVER = False
produce_data(save_que)
PRODUCER_IS_OVER = True
save_p.join()
But this cannot work in case the queue is produced by several separated process: who is going to send the ALT message in that case?
another solution is to store the processes index in a list and execute:
def some_alive():
for p in processes:
if p.is_alive():
return True
return False
But multiprocessing supports .is_alive method only in the parent process, which is limiting in my case.