Empresas
Empleos
  • Sobre nosotros
  • Soluciones
    • Publicación de vacantes
      Publica tu vacante y recibe candidatos calificados en 48h.
    • Evaluación de candidatos
      500+ pruebas técnicas y psicológicas, más anti-fraude.
    • Headhunting
      Búsqueda ejecutiva a la medida de principio a fin.
    • Nómina + EOR
      Dispersión de nómina y EOR en más de 15 países de LATAM.
  • Precios
  • Empleos

0

296
Vistas
imap_unordered, pero con un generador aplanado perezoso

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?

over 4 years ago · Santiago Trujillo
6 Respuestas
Responde la pregunta

0

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)
over 4 years ago · Santiago Trujillo Denunciar

0

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.025248765945435
over 4 years ago · Santiago Trujillo Denunciar

0

No me gusta la idea de pasar por alto las colas internas de Pool superponiéndolas con colas externas. En mi opinión, conduce a muchas más partes móviles y una complejidad innecesaria y fácilmente terminas creando condiciones de carrera difíciles de detectar. Pool solo es bastante complejo bajo el capó e inflar la cantidad de código que se ejecuta, incluso con más tuberías en la parte superior es algo que preferiría evitar (KISS). Está usando Pool solo por su efecto secundario de administrar los procesos de trabajo y, en todo caso, solo lo consideraría para un código único, no para una construcción de sistema para estabilidad o posiblemente necesidades en evolución.

Para darle una comparación del argumento de la complejidad... Pool emplea tres subprocesos de trabajo solo para administrar trabajadores y canalizar datos de un lado a otro. Junto con el subproceso principal, esto forma cuatro subprocesos en el proceso principal. La versión Non-Pool que se proporciona a continuación, por otro lado, tiene dos subprocesos ( multiprocessing.Queue . Queue inicia un subproceso de alimentación en el primer uso de .put() ). Las ~60 líneas de la solución que no es Pool se comparan con las ~900 líneas de multiprocessing.pool.py solo. Una buena parte de este último se ejecutará de todos modos, solo para mezclar None en lugar de los resultados reales.

Pool es excelente para el caso de uso frecuente de funciones de procesamiento de tareas en su conjunto , la recuperación de resultados de subtareas de los generadores simplemente no se ajusta aquí.


Uso multiprocessing.Pool

Ahora bien, si está decidido a seguir adelante con este enfoque de todos modos, al menos no use Manager para ello. Casi no hay ninguna razón para buscar el Manager en un solo nodo y la única necesidad que se me ocurre para la que realmente necesita usar las colas del administrador es cuando tiene que enviar una referencia de cola a un proceso que ya está en funcionamiento . Dado que el initializer() de Pool le permite pasar argumentos ya en el inicio del proceso de trabajo, las colas del administrador tampoco son necesarias.

Debe tener en cuenta que cada interacción entre el proceso principal y el secundario a través de los Manager da como resultado un desvío a través de un proceso de administrador adicional, lo que aumenta la latencia en más IPC, cambios de contexto y vaciados de caché. También tiene el potencial de resultar en una huella de memoria considerablemente mayor.

 import time import multiprocessing as mp # def generate(x): ... # take from question def init_queue(queue): globals()['queue'] = queue def wrapper(x): q = queue for item in generate(x): q.put(item) if __name__ == '__main__': POISON = 'POISON' queue = mp.SimpleQueue() with mp.Pool(processes=4, initializer=init_queue, initargs=(queue,)) as pool: pool.map_async( func=wrapper, iterable=range(0, 100, 10), chunksize=1, callback=lambda _: queue.put(POISON) ) for res in iter(queue.get, POISON): print(res)

Uso multiprocessing.Process

Ahora, la alternativa que preferiría sobre el uso Pool en este caso es construir su propio pequeño grupo especializado, con multiprocessing.Process y algo de multiprocessing-queue. Sí, es un poco más de código para escribir, pero la cantidad de código que realmente se ejecuta se reduce considerablemente en comparación con cualquier solución que involucre multiprocessing.Pool . Esto deja menos espacio para errores sutiles, además está más abierto a condiciones cambiantes y utiliza menos recursos del sistema.

 import time import multiprocessing as mp from itertools import chain DONE = 'DONE' POISON = 'POISON' def _worker(func, inqueue, outqueue): for chunk in iter(inqueue.get, POISON): for res in func(chunk): outqueue.put(res) outqueue.put(DONE) outqueue.put(POISON) def _init_pool(n_workers, func): """Initialize worker-processes and queues.""" inqueue, outqueue = mp.Queue(), mp.Queue() pool = [ mp.Process(target=_worker, args=(func, inqueue, outqueue)) for _ in range(n_workers) ] for p in pool: p.start() return pool, inqueue, outqueue def iflatmap(n_workers, func, iterable): """Yield results from subprocesses unordered and immediately.""" iterable = chain(iterable, [POISON] * n_workers) pool, inqueue, outqueue = _init_pool(n_workers, func) for _ in pool: inqueue.put(next(iterable)) while n_workers: res = outqueue.get() if res == DONE: # there's a free worker now inqueue.put(next(iterable)) elif res == POISON: # a worker has shut down n_workers -= 1 else: yield res for p in pool: p.join()

El ejemplo aquí usa solo cuatro trabajadores para mostrar que esta solución (también) no depende de tener la misma cantidad de trabajadores y tareas. Dado que tampoco necesita conocer la longitud de la entrada iterable para su flujo de control, esto permite proporcionar un generador de longitud desconocida como iterable. Si lo prefiere con clase, puede envolver la lógica anterior en una clase "Pool" en su lugar.

 def generate(x): for j in range(x, x + 10): yield j time.sleep(1) if __name__ == '__main__': for res in iflatmap(n_workers=4, func=generate, iterable=range(0, 100, 10)): print(res)

Para personas que no están familiarizadas con el uso de iter(object, sentinel) : documentos y algunos razonamientos aquí

over 4 years ago · Santiago Trujillo Denunciar

0

Probablemente 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()}')
over 4 years ago · Santiago Trujillo Denunciar

0

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 vuelta 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()
over 4 years ago · Santiago Trujillo Denunciar

0

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)
over 4 years ago · Santiago Trujillo Denunciar
Responde la pregunta
Encuentra empleos remotos

¡Descubre la nueva forma de encontrar empleo!

Top de empleos
Top categorías de empleo
Empresas
Publicar vacante Precios Comercial
Legal
Términos y condiciones Política de privacidad
© 2026 PeakU Inc. All Rights Reserved.
Andres GPT
Recomiéndame algunas ofertas
Necesito ayuda