Estoy tratando de procesar un archivo (cada línea es un documento json). El tamaño del archivo puede llegar hasta cientos de mbs a gb. Así que escribí un código generador para obtener cada documento línea por línea del archivo.
def jl_file_iterator(file): with codecs.open(file, 'r', 'utf-8') as f: for line in f: document = json.loads(line) yield documentMi sistema tiene 4 núcleos, por lo que me gustaría procesar 4 líneas del archivo en paralelo. Actualmente tengo este código que toma 4 líneas a la vez y llama al código para procesamiento paralelo
threads = 4 files, i = [], 1 for jl in jl_file_iterator(input_path): files.append(jl) if i % (threads) == 0: # pool.map(processFile, files) parallelProcess(files, o) files = [] i += 1 if files: parallelProcess(files, o) files = []Este es mi código donde ocurre el procesamiento real
def parallelProcess(files, outfile): processes = [] for i in range(len(files)): p = Process(target=processFile, args=(files[i],)) processes.append(p) p.start() for i in range(len(files)): processes[i].join() def processFile(doc): extractors = {} ... do some processing on doc o.write(json.dumps(doc) + '\n')Como puede ver, espero que las 4 líneas terminen de procesarse antes de enviar los siguientes 4 archivos para procesar. Pero lo que me gustaría hacer es tan pronto como un proceso termine de procesar el archivo, quiero comenzar la siguiente línea para asignarla al procesador liberado. ¿Cómo puedo hacer eso?
PD: El problema es que, dado que es un generador, no puedo cargar todos los archivos y usar algo como un mapa para ejecutar los procesos.
Gracias por tu ayuda
Como dijo @pvg en un comentario, una cola (limitada) es la forma natural de mediar entre un productor y consumidores con diferentes velocidades, asegurando que todos estén lo más ocupados posible pero sin dejar que el productor se adelante.
Aquí hay un ejemplo ejecutable autónomo. La cola está restringida a un tamaño máximo igual al número de procesos de trabajo. Si los consumidores corren mucho más rápido que el productor, podría tener sentido dejar que la cola sea más grande que eso.
En su caso específico, probablemente tendría sentido pasar líneas a los consumidores y dejarles hacer la parte document = json.loads(line) en paralelo.
import multiprocessing as mp NCORE = 4 def process(q, iolock): from time import sleep while True: stuff = q.get() if stuff is None: break with iolock: print("processing", stuff) sleep(stuff) if __name__ == '__main__': q = mp.Queue(maxsize=NCORE) iolock = mp.Lock() pool = mp.Pool(NCORE, initializer=process, initargs=(q, iolock)) for stuff in range(20): q.put(stuff) # blocks until q below its max size with iolock: print("queued", stuff) for _ in range(NCORE): # tell workers we're done q.put(None) pool.close() pool.join()Así que terminé ejecutando esto con éxito. Creando fragmentos de líneas de mi archivo y ejecutando las líneas en paralelo. Publicarlo aquí para que pueda ser útil para alguien en el futuro.
def run_parallel(self, processes=4): processes = int(processes) pool = mp.Pool(processes) try: pool = mp.Pool(processes) jobs = [] # run for chunks of files for chunkStart,chunkSize in self.chunkify(input_path): jobs.append(pool.apply_async(self.process_wrapper,(chunkStart,chunkSize))) for job in jobs: job.get() pool.close() except Exception as e: print e def process_wrapper(self, chunkStart, chunkSize): with open(self.input_file) as f: f.seek(chunkStart) lines = f.read(chunkSize).splitlines() for line in lines: document = json.loads(line) self.process_file(document) # Splitting data into chunks for parallel processing def chunkify(self, filename, size=1024*1024): fileEnd = os.path.getsize(filename) with open(filename,'r') as f: chunkEnd = f.tell() while True: chunkStart = chunkEnd f.seek(size,1) f.readline() chunkEnd = f.tell() yield chunkStart, chunkEnd - chunkStart if chunkEnd > fileEnd: breakLa respuesta de Tim Peters es genial.
Pero mi caso específico era ligeramente diferente y tuve que modificar su respuesta para que se ajustara a mi necesidad. Haciendo referencia aquí.
Esto responde a la pregunta de @CpILL en los comentarios.
En mi caso, usé una cadena de generadores (para crear una tubería).
Entre esta cadena de generadores, uno de ellos estaba haciendo cálculos pesados, ralentizando toda la tubería.
Algo como esto :
def fast_generator1(): for line in file: yield line def slow_generator(lines): for line in lines: yield heavy_processing(line) def fast_generator2(): for line in lines: yield fast_func(line) if __name__ == "__main__": lines = fast_generator1() lines = slow_generator(lines) lines = fast_generator2(lines) for line in lines: print(line) Para hacerlo más rápido, tenemos que ejecutar el generador lento con múltiples procesos.
El código modificado se parece a:
import multiprocessing as mp NCORE = 4 def fast_generator1(): for line in file: yield line def slow_generator(lines): def gen_to_queue(input_q, lines): # This function simply consume our generator and write it to the input queue for line in lines: input_q.put(line) for _ in range(NCORE): # Once generator is consumed, send end-signal input_q.put(None) def process(input_q, output_q): while True: line = input_q.get() if line is None: output_q.put(None) break output_q.put(heavy_processing(line)) input_q = mp.Queue(maxsize=NCORE * 2) output_q = mp.Queue(maxsize=NCORE * 2) # Here we need 3 groups of worker : # * One that will consume the input generator and put it into a queue. It will be `gen_pool`. It's ok to have only 1 process doing this, since this is a very light task # * One that do the main processing. It will be `pool`. # * One that read the results and yield it back, to keep it as a generator. The main thread will do it. gen_pool = mp.Pool(1, initializer=gen_to_queue, initargs=(input_q, lines)) pool = mp.Pool(NCORE, initializer=process, initargs=(input_q, output_q)) finished_workers = 0 while True: line = output_q.get() if line is None: finished_workers += 1 if finished_workers == NCORE: break else: yield line def fast_generator2(): for line in lines: yield fast_func(line) if __name__ == "__main__": lines = fast_generator1() lines = slow_generator(lines) lines = fast_generator2(lines) for line in lines: print(line)Con esta implementación, tenemos un generador multiproceso: se usa exactamente como otros generadores (como en el primer ejemplo de esta respuesta), pero todos los cálculos pesados se realizan usando multiprocesamiento, ¡acelerándolo!