Tengo un problema que ya se puede resolver con multiprocessing.Pool pero la solución no es muy óptima. Es decir, lo que tengo es un conjunto bastante pequeño de entradas, cada una de las cuales se asigna a un gran conjunto de datos. Si bien puedo usar imap_unordered con una función que devuelve una lista, esto está lejos de ser eficiente, porque cada uno de los grandes conjuntos de datos debe devolverse como una lista.
Mi función podría devolverlos como un generador para una latencia más baja, pero no puedo devolver un generador de un subproceso.
Un ejemplo ficticio:
import time import multiprocessing def generate(x): for j in range(x, x + 10): yield j time.sleep(1) def wrapper(x): return list(generate(x)) with multiprocessing.Pool(10) as pool: for ready in pool.imap_unordered(wrapper, range(0, 100, 10)): for item in set(ready): # to show that order does not matter: print(item)El problema es que, si bien la ejecución completa ahora toma solo la décima parte del tiempo de ejecución secuencial, aún debo esperar 10 segundos para obtener el primer resultado, que podría estar disponible de inmediato al:
def generate(x): for j in range(x, x + 10): yield j time.sleep(1) for ready in map(generate, range(0, 100, 10): for item in set(ready): print(item)Que imprimirá el primer elemento sin demora, pero tarda 100 segundos en ejecutarse.
Lo que no puedo hacer es subdividir aún más el problema, los generadores en los subprocesos deben ser evaluados perezosamente por el consumidor:
def generate(x): for j in range(x, x + 10): yield j time.sleep(1) with multiprocessing.Pool(10) as pool: for item in pool.??flatmap_unordered??(generate, range(0, 100, 10)): print(item)que imprimiría el primer elemento de inmediato, ¡pero solo demora ~ 10 segundos en ejecutarse!
¿Cómo podría lograr eso?
Parece que no hay una forma integrada para que un grupo Pool de forma incremental los elementos generados. Sin embargo, es razonablemente sencillo escribir su propio ayudante de "mapa plano".
La idea general es tener un contenedor en los procesos del grupo que ejecuta el iterador y empuja cada elemento individual a una cola. En el proceso principal, solo hay un bucle simple que obtiene y yield cada elemento.
import functools import multiprocessing def flatmap(pool: multiprocessing.Pool, func, iterable, chunksize=None): """A flattening, unordered equivalent of Pool.map()""" # use a queue to stream individual results from processes queue = multiprocessing.Manager().Queue() # reuse task management and mapping of Pool pool.map_async( functools.partial(_flat_mappper, queue, func), iterable, chunksize, # callback: push a signal that everything is done lambda _: queue.put(None), lambda err: queue.put((None, err)) ) # yield each result as it becomes available while True: item = queue.get() if item is None: break result, err = item if err is None: yield result else: raise err def _flat_mappper(queue: multiprocessing.Queue, func, *args): """Helper to run `(*args) -> iter` and stream results to a queue""" for item in func(*args): queue.put((item, None)) Si lo desea, se podría parchear el tipo Pool para tener un flatmap como método en lugar de como función.
El asistente flatmap se puede usar directamente para acumular resultados en los generadores. Para el caso del ejemplo, termina en un poco más de 10 segundos.
import time def generate(x): for j in range(x, x + 10): yield j time.sleep(1) if __name__ == "__main__": with multiprocessing.Pool(10) as pool: for item in flatmap(pool, generate, range(0, 100, 10)): print(item)La respuesta de MisterMiyagi puede estar estancada en Linux debido a la forma en que se serializa "Cola" cuando está dentro de un objeto "parcial". (No puedo entender por qué).
Este es el mismo código, con un truco para pasar la cola al asignador en el lado del parámetro, y no incrustado en el objeto invocable.
El cerebro se siente en cámara lenta hoy, lo siento, no puedo entender exactamente qué está sucediendo; de todos modos, esta variante del código funcionó aquí:
import functools import multiprocessing import secrets from itertools import repeat import time def flatmap(pool: multiprocessing.Pool, func, iterable, chunksize=None): """A flattening, unordered equivalent of Pool.map()""" # use a queue to stream individual results from processes queue = multiprocessing.Manager().Queue() sentinel = secrets.token_bytes() # reuse task management and mapping of Pool pool.map_async( functools.partial(_flat_mappper, func), zip(iterable,repeat(queue)), chunksize, # callback: push a signal that everything is done lambda _: queue.put(sentinel), lambda *e: print("argh!", e), ) # yield each result as it becomes available while True: item = queue.get() if item == sentinel: break yield item def _flat_mappper(func, arg): """Helper to run `(*args) -> iter` and stream results to a queue""" data, queue = arg for item in func(data): queue.put(item) def generate(x): for j in range(x, x + 10): yield j time.sleep(1) if __name__ == "__main__": with multiprocessing.Pool(10) as pool: for item in flatmap(pool, generate, range(0, 100, 10)): print(item)Esta es una instancia en la que creo que el enfoque más simple sería implementar su propio grupo de multiprocesamiento desde los procesos daemon y las instancias de multiprocessing.Queue (que tienen más rendimiento que las colas administradas devueltas por las llamadas a multiprocessing.Manager().Queue() :
import time import multiprocessing # We just need something distinct from a value # generated by generate. In this case nothing fancy is required: SENTINEL = None def generate(x): for j in range(x, x + 10): yield j time.sleep(1) def worker(in_q, out_q): while True: x = in_q.get() for j in generate(x): out_q.put(j) # put sentinel to queue to show task is complete: out_q.put(SENTINEL) # In case we are running under Windows: if __name__ == '__main__': POOL_SIZE = 10 in_q = multiprocessing.Queue() out_q = multiprocessing.Queue() for _ in range(POOL_SIZE): multiprocessing.Process(target=worker, args=(in_q, out_q), daemon=True).start() t = time.time() # Submit the 10 tasks: for x in range(0, 100, 10): in_q.put(x) # We have submitted 10 tasks, so when we have seen 10 # sentinels, we know we have processed all the results sentinels_seen = 0 results = [] while sentinels_seen < 10: return_value = out_q.get() if return_value is SENTINEL: sentinels_seen += 1 else: # Process return value: results.append(return_value) if len(results) == 1: # First result: print('Elapsed time to first result:', time.time() - t) print('Total elapsed time:', time.time() - t)Huellas dactilares:
Elapsed time to first result: 0.004938602447509766 Total elapsed time: 10.025248765945435Probablemente podría usar multiprocessing.pipes y enviar el resultado a través de las conexiones para obtener los primeros elementos de inmediato:
import time import multiprocessing def generate_and_send(x, conn): for j in range(x, x + 10): conn.send(j) time.sleep(1) conn.send("POISON") print(f'program begins at {time.ctime()}') pipes = [multiprocessing.Pipe() for _ in range(10)] child_conns = [y for (x, y) in pipes] parent_conns = [x for (x, y) in pipes] processes = [multiprocessing.Process(target=generate_and_send, args=(*el,)) for el in zip(range(0, 100, 10), child_conns)] for p in processes: p.start() first_printed = False while True: if parent_conns: for par_conn in parent_conns: recvd_val = par_conn.recv() if recvd_val == "POISON": par_conn.close() parent_conns.remove(par_conn) else: if not first_printed: print(f'first item printed at {time.ctime()}') first_printed = True print(recvd_val) else: break print(f'program ends at {time.ctime()}')Como señaló el OP, los objetos pasados entre Procesos deben ser pickleables. El objetivo declarado es procesar objetos de datos tan pronto como estén disponibles. Por lo tanto, los Procesos secundarios deben devolver los objetos al Proceso principal lo antes posible. Esto establece restricciones sobre qué herramientas se pueden utilizar.
Los procesos secundarios pueden usar una cola de multiprocesamiento para transmitir los objetos individuales de regreso al proceso principal. Las funciones de trabajador en los Procesos secundarios no necesitan devolver nada. Parte de la maquinaria en los métodos de la biblioteca estándar está diseñada para recopilar y administrar valores devueltos; no será útil. Sin embargo, aún es necesario inicializar los Procesos secundarios y enviarles tareas. También es útil administrar el flujo del programa utilizando los administradores de contexto de biblioteca estándar.
Los métodos de multiprocessing.Pool funcionarán, pero prefiero una solución usando ProcessPoolExecutor. Mi impl1() a continuación usa Pool, e impl2() usa el Ejecutor. Ambas soluciones son cortas y, creo, fáciles de entender y generalizar. En ambos casos, un objeto multiprocessing.Queue debe inicializarse y pasarse al Proceso secundario en su función de inicialización.
Dado que el orden en el que se devuelven los datos al proceso principal no es determinista, es necesario tener alguna forma de clasificar qué datos provienen de qué tarea. Utilizo un esquema de indexación de enteros simple y recopilo todos los resultados en un solo diccionario. Esto se administra en un subproceso secundario, mientras que el subproceso principal espera a que finalicen todas las tareas enviadas (el contexto ProcessPoolExecutor no se cerrará hasta que se completen todas las tareas enviadas). Cuando eso sucede, publico un objeto centinela en la cola para apagarlo. Me uno al hilo secundario para esperar a que se procesen todos los objetos. En este punto, el diccionario que contiene los elementos de datos recopilados está completo.
El código es multiproceso pero la estructura es bastante simple. Las funciones del Proceso principal están bien separadas en configuración, envío de tareas y manejo de datos. No hay que preocuparse por las condiciones de carrera.
import time import multiprocessing from concurrent.futures import ProcessPoolExecutor import threading from collections import defaultdict # Process code _q = None def initialize(q): global _q # pylint: disable=global-statement _q = q def generate(x): for j in range(x, x + 10): yield j time.sleep(1) def wrapper(x): for j in generate(x): _q.put((x, j)) # collection thread def collect(q, results): for a, b in iter(q.get, (None, None)): print(a, b) results[a].append(b) def impl1(q): with multiprocessing.Pool(10, initializer=initialize, initargs=(q,)) as pool: for _ in pool.imap_unordered(wrapper, range(0, 100, 10)): pass def impl2(q): with ProcessPoolExecutor(10, initializer=initialize, initargs=(q,)) as ex: for n in range(0, 100, 10): ex.submit(wrapper, n) def main(): results = defaultdict(list) q = multiprocessing.Queue() thr = threading.Thread(target=collect, args=(q, results)) thr.start() # impl1(q) impl2(q) q.put((None, None)) thr.join() print(results) if __name__ == "__main__": main()