Optimización de la transferencia de datos en cargas de trabajo de inferencia de IA/ML por lotes

es una optimización de la transferencia de datos en cargas de trabajo de IA/ML donde demostramos el uso de los sistemas NVIDIA Nsight™ (nsys) para estudiar y resolver los cuellos de botella comunes en la carga de datos: situaciones en las que la GPU está inactiva mientras espera datos de entrada de la CPU. En esta publicación centramos nuestra atención en los datos que viajan en la dirección opuesta, desde el dispositivo GPU hasta el host de la CPU. Más específicamente, abordamos cargas de trabajo de inferencia de IA/ML donde el tamaño de la salida que devuelve el modelo es relativamente alto. Los ejemplos comunes incluyen: 1) ejecutar un modelo de segmentación de escena (etiquetado por píxel) en lotes de imágenes de alta resolución y 2) capturar incorporaciones de características de alta dimensión de secuencias de entrada utilizando un modelo codificador (por ejemplo, para crear una base de datos vectorial). Ambos ejemplos implican ejecutar un modelo en un lote de entrada y luego copiar el tensor de salida de la GPU a la CPU para procesamiento, almacenamiento y/o comunicación a través de la red adicionales.

Las copias de memoria de GPU a CPU de la salida del modelo suelen recibir mucha menos atención en los tutoriales de optimización que las copias de CPU a GPU que alimentan el modelo (por ejemplo, consulte aquí). Pero su impacto potencial en la eficiencia del modelo y los costos de ejecución puede ser igualmente perjudicial. Además, si bien las optimizaciones para la carga de datos de CPU a GPU están bien documentadas y son fáciles de implementar, optimizar la copia de datos en la dirección opuesta requiere un poco más de trabajo manual.

En esta publicación aplicaremos la misma estrategia que usamos en nuestra publicación anterior: definiremos un modelo de juguete y usaremos nsys Profiler para identificar y resolver cuellos de botella en el rendimiento. Ejecutaremos nuestros experimentos en una instancia Amazon EC2 g6e.2xlarge (con una GPU NVIDIA L40S) que ejecuta una AMI de aprendizaje profundo de AWS (Ubuntu 24.04) con PyTorch (2.8), perfilador nsys-cli (versión 2025.6.1) y la biblioteca NVIDIA Tools Extension (NVTX).

Descargos de responsabilidad

El código que compartiremos tiene fines demostrativos; no confíe en su corrección u optimidad. No interprete nuestro uso de ninguna biblioteca, herramienta o plataforma como un respaldo a su uso. El impacto de las optimizaciones que cubriremos puede variar mucho según los detalles del modelo y el entorno de ejecución. Asegúrese de evaluar su efecto en su propio caso de uso antes de integrar su uso.

Muchas gracias a Yitzhak Levi y Gilad Wasserman por sus contribuciones a esta publicación.

Un modelo de juguete PyTorch

Presentamos un script de inferencia por lotes que realiza la segmentación de imágenes en un conjunto de datos sintéticos utilizando un modelo DeepLabV3 con una red troncal ResNet-50. Las salidas del modelo se copian a la CPU para su posterior procesamiento y almacenamiento. Envolvemos las diferentes partes del paso de inferencia con anotaciones nvtx codificadas por colores:

importar tiempo, antorcha, nvtx de torch.utils.data importar conjunto de datos, cargador de datos de torch.cuda importar perfilador de torchvision.models.segmentation importar deeplabv3_resnet50 DEVICE = "cuda" WARMUP_STEPS = 10 PROFILE_STEPS = 3 COOLDOWN_STEPS = 1 TOTAL_STEPS = WARMUP_STEPS + PROFILE_STEPS + COOLDOWN_STEPS BATCH_SIZE = 64 TOTAL_SAMPLES = TOTAL_STEPS * BATCH_SIZE IMG_SIZE = 512 N_CLASSES = 21 NUM_WORKERS = 8 ASYNC_DATALOAD = True # Un conjunto de datos sintético con imágenes aleatorias clase FakeDataset(Dataset): def __len__(self): return TOTAL_SAMPLES def __getitem__(self, index): img = torch.randn((3, IMG_SIZE, IMG_SIZE)) return img # clase de utilidad para precargar datos a la clase GPU DataPrefetcher: def __init__(self, loader): self.loader = iter(loader) self.stream = torch.cuda.Stream() self.next_batch = Ninguno self.preload() def preload(self): try: data = next(self.loader) with torch.cuda.stream(self.stream): next_data = data.to(DEVICE, non_blocking=ASYNC_DATALOAD) self.next_batch = next_data excepto: self.next_batch = Ninguno def __iter__(self): devolver self def __next__(self): torch.cuda.current_stream().wait_stream(self.stream) datos = self.next_batch self.preload() devolver modelo de datos = deeplabv3_resnet50(weights_backbone=None).to(DEVICE).eval() data_loader = DataLoader( FakeDataset(), lote_size=BATCH_SIZE, num_workers=NUM_WORKERS, pin_memory=ASYNC_DATALOAD ) data_iter = DataPrefetcher(data_loader) def sincronizar_todos(): torch.cuda.synchronize() def to_cpu(salida): devuelve salida.cpu() def proceso_salida(batch_id, logits): # realiza un posprocesamiento en la salida con open('/dev/null', 'wb') como f: f.write(logits.numpy().tobytes()) con torch.inference_mode(): for i in range(TOTAL_STEPS): if i == WARMUP_STEPS: sincronizar_all() start_time = time.perf_counter() perfiler.start() elif i == WARMUP_STEPS + PROFILE_STEPS: sincronizar_todos() perfiler.stop() end_time = time.perf_counter() con nvtx.annotate(f"Batch {i}", color="blue"): con nvtx.annotate("get lote", color="red"): lote = next(data_iter) con nvtx.annotate("compute", color="green"): salida = modelo(lote) con nvtx.annotate("copiar a CPU", color="amarillo"): salida_cpu = to_cpu(salida['out']) con nvtx.annotate("proceso de salida", color="cian"): salida_de_proceso(i, CPU_de_salida) tiempo_total = tiempo_final – tiempo_inicio rendimiento = PROFILE_STEPS / total_time print(f"Rendimiento: {rendimiento:.2f} pasos/seg")

Tenga en cuenta la inclusión de todas las optimizaciones de carga de datos de CPU a GPU discutidas en nuestra publicación anterior.

Ejecutamos el siguiente comando para capturar un seguimiento del perfil nsys:

perfil nsys –capture-range=cudaProfilerApi –trace=cuda,nvtx,osrt –output=baseline python batch_infer.py

Esto da como resultado un archivo de seguimiento baseline.nsys-rep que copiamos a nuestra máquina de desarrollo para su análisis.

Para medir el rendimiento de la inferencia, aumentamos el número de pasos a 100. El rendimiento promedio de nuestro experimento de referencia es de 0,45 pasos por segundo. En las siguientes secciones utilizaremos los seguimientos del perfil nsys para mejorar gradualmente este resultado.

Análisis de desempeño de referencia

La siguiente imagen muestra el seguimiento del perfil nsys de nuestro experimento de referencia:

Seguimiento del perfilador de sistemas Baseline Nsight (por autor)

En la sección de GPU vemos el siguiente patrón recurrente:

Un bloque de cálculo del kernel (en azul claro) que se ejecuta durante aproximadamente 520 milisegundos. Un pequeño bloque de copia de memoria del host al dispositivo (en verde) que se ejecuta en paralelo al cálculo del kernel. Esta simultaneidad se logró utilizando las optimizaciones discutidas en nuestra publicación anterior. Un bloque de copia de memoria del dispositivo al host (en rojo) que se ejecuta durante aproximadamente 750 milisegundos. Un largo período (~940 milisegundos) de tiempo de inactividad de la GPU (espacio en blanco) entre cada dos pasos.

Al observar la barra NVTX de la sección de CPU, podemos ver que el espacio en blanco se alinea perfectamente con el bloque de “salida del proceso” (en cian). En nuestra implementación inicial, tanto la ejecución del modelo como la función de almacenamiento de salida se ejecutan en el mismo proceso único de manera secuencial. Esto genera un tiempo de inactividad significativo en la GPU, ya que la CPU espera a que regrese la función de almacenamiento antes de alimentar a la GPU con el siguiente lote.

Optimización 1: procesamiento de resultados de varios trabajadores

El primer paso que damos es ejecutar la función de almacenamiento de salida en procesos de trabajo paralelos. Dimos un paso similar en nuestra publicación anterior cuando trasladamos la secuencia de preparación del lote de entrada a trabajadores dedicados. Sin embargo, mientras que allí pudimos automatizar la carga de datos multiproceso simplemente configurando el argumento num_workers de la clase DataLoader en un valor distinto de cero, aplicar el procesamiento de salida de múltiples trabajadores requiere una implementación manual. Aquí elegimos una solución simple con fines demostrativos. Esto debe personalizarse según sus necesidades y preferencias de diseño.

Multiprocesamiento PyTorch

Implementamos una estrategia productor-consumidor utilizando el paquete de multiprocesamiento integrado de PyTorch, torch.multiprocessing. Definimos una cola para almacenar lotes de salida y varios trabajadores consumidores que procesan los lotes en la cola. Modificamos nuestro bucle de inferencia para colocar los buffers de salida en la cola de salida. También actualizamos la utilidad sincronizar_all() para vaciar la cola y agregar una secuencia de limpieza al final del script.

El siguiente bloque de código contiene nuestra implementación inicial. Como veremos en las siguientes secciones, esto requerirá algunos ajustes para alcanzar el máximo rendimiento.

importar torch.multiprocessing como mp POSTPROC_WORKERS = 8 # sintonizar para un rendimiento óptimo output_queue = mp.JoinableQueue(maxsize=POSTPROC_WORKERS) def output_worker(in_q): while True: item = in_q.get() si el elemento es Ninguno: romper # señal para cerrar lote_id, lote_preds = elemento proceso_salida(batch_id, lote_preds) in_q.task_done() procesos =[]para _ en rango (POSTPROC_WORKERS): p = mp.Process(target=output_worker, args=(output_queue,)) p.start()processs.append(p) def sincronizar_all(): torch.cuda.synchronize()output_queue.join() # drenaje de cola con torch.inference_mode(): para i en rango(TOTAL_STEPS): si i == WARMUP_STEPS: sincronizar_todo() tiempo_inicial = time.perf_counter() perfilador.start() elif i == PASOS_CALENTAMIENTOS + PASOS_PERFIL: sincronizar_todo() perfilador.stop() tiempo_final = tiempo.perf_counter() con nvtx.annotate(f"Batch {i}", color="azul"): con nvtx.annotate("obtener lote", color="rojo"): lote = siguiente(data_iter) con nvtx.annotate("compute", color="green"): salida = modelo(batch) con nvtx.annotate("copiar a CPU", color="amarillo"): salida_cpu = to_cpu(salida['out']) con nvtx.annotate("cola de salida", color="cyan"): salida_queue.put((i, salida_cpu)) total_time = end_time – start_time rendimiento = PROFILE_STEPS / tiempo_total print(f"Rendimiento: {rendimiento:.2f} pasos/seg") # limpieza para _ en el rango(POSTPROC_WORKERS): salida_queue.put(Ninguno)

La optimización del procesamiento de salida de múltiples trabajadores da como resultado un rendimiento de 0,71 pasos por segundo, un aumento del 58 % con respecto a nuestros resultados de referencia.

Al volver a ejecutar el comando nsys se obtiene el siguiente seguimiento del perfil:

Cronología del perfilador de sistemas Nsight para múltiples trabajadores (por autor)

Podemos ver que el tamaño del bloque de espacios en blanco se ha reducido considerablemente (de ~940 milisegundos a ~50). Si acercáramos el espacio en blanco restante, lo encontraríamos alineado con una operación "munmap". En nuestra publicación anterior, el mismo hallazgo informó nuestra optimización de copia de datos asincrónica. Pero esta vez damos un paso intermedio de optimización de la memoria en forma de un grupo de buffers preasignados.

Optimización 2: preasignación del grupo de búfer

Para reducir la sobrecarga de asignar y administrar un nuevo tensor de CPU en cada iteración, inicializamos un grupo de tensores preasignados en la memoria compartida y definimos una segunda cola para administrar su uso.

Nuestro código actualizado aparece a continuación:

forma = (BATCH_SIZE, N_CLASSES, IMG_SIZE, IMG_SIZE) buffer_pool = [torch.empty(shape).share_memory_() for _ in range(POSTPROC_WORKERS)] buf_queue = mp.Queue() for i in range(POSTPROC_WORKERS): buf_queue.put(i) def output_worker(buffer_pool, in_q, buf_q): while Verdadero: item = in_q.get() si el elemento es Ninguno: romper # señal para cerrar batch_id, buf_id = item process_output(batch_id, buffer_pool[buf_id]) buf_q.put(buf_id) in_q.task_done() procesos =[]para _ en rango(POSTPROC_WORKERS): p = mp.Process(target=output_worker, args=(buffer_pool,output_queue,buf_queue)) p.start() procesos.append(p) def to_cpu(salida): buf_id = buf_queue.get() salida_cpu = buffer_pool[buf_id] salida_cpu.copy_(salida) return salida_cpu, buf_id con torch.inference_mode(): para i en el rango(TOTAL_STEPS): if i == WARMUP_STEPS: sincronizar_todos() start_time = time.perf_counter() perfilador.start() elif i == WARMUP_STEPS + PROFILE_STEPS: sincronizar_all() perfilador.stop() end_time = time.perf_counter() con nvtx.annotate(f"Batch {i}", color="blue"): con nvtx.annotate("obtener lote", color="red"): lote = siguiente(data_iter) con nvtx.annotate("compute", color="green"): salida = modelo(batch) con nvtx.annotate("copiar a CPU", color="amarillo"): salida_cpu, buf_id = to_cpu(salida['salida']) con nvtx.annotate("salida de cola", color="cian"): salida_queue.put((i, buf_id))

Después de estos cambios, el rendimiento de inferencia salta a 1,51, una aceleración de más del doble con respecto a nuestro resultado anterior.

El nuevo seguimiento del perfil aparece a continuación:

Cronología del perfilador de sistemas Nsight de Buffer Pool (por autor)

No solo el espacio en blanco prácticamente desapareció, sino que la operación de la memoria CUDA DtoH (en rojo) se redujo de ~750 milisegundos a ~110. Presumiblemente, la gran copia de datos de GPU a CPU implicó bastante sobrecarga de administración de memoria que eliminamos mediante la implementación de un grupo de búfer dedicado.

A pesar de la mejora considerable, si nos acercamos encontraremos que quedan alrededor de ~0,5 milisegundos de espacio en blanco causado por la sincronicidad del comando de copia de GPU a CPU; mientras la copia no se haya completado, la CPU no activa el cálculo del kernel del siguiente lote.

Optimización 3: copia de datos asincrónica

Nuestra tercera optimización es cambiar la copia del dispositivo al host para que sea asíncrona. Como antes, encontraremos que implementar este cambio es más difícil que en la dirección de CPU a GPU.

El primer paso es pasar non_blocking=True al comando de copia de GPU a CPU.

def to_cpu(salida): buf_id = buf_queue.get() salida_cpu = buffer_pool[buf_id] salida_cpu.copy_(salida, non_blocking=True) return salida_cpu, buf_id

Sin embargo, como vimos en nuestra publicación anterior, este cambio no tendrá un impacto significativo a menos que modifiquemos nuestros tensores para usar memoria fija:

forma = (BATCH_SIZE, N_CLASSES, IMG_SIZE, IMG_SIZE) buffer_pool = [torch.empty(shape, pin_memory=True).share_memory_() para _ en el rango (POSTPROC_WORKERS)]

Fundamentalmente, si aplicamos solo estos dos cambios a nuestro script, el rendimiento aumentaría pero la salida podría dañarse (por ejemplo, ver aquí). Necesitamos un mecanismo basado en eventos para identificar cada vez que se completa una copia de GPU a CPU para que podamos continuar con el procesamiento de datos de salida. (Tenga en cuenta que esto no fue necesario al realizar la copia de CPU a GPU asíncrona. Debido a que un único flujo de GPU procesa los comandos de forma secuencial, el cálculo del núcleo solo comienza cuando la copia se ha completado. La sincronización solo fue necesaria al introducir un segundo flujo).

Para implementar el mecanismo de notificación, definimos un grupo de eventos CUDA y una cola adicional para gestionar su uso. Además, definimos un subproceso de escucha para monitorear el estado de los eventos en la cola y completar la cola de salida una vez que se completen las copias.

importar subprocesos, cola event_pool = [torch.cuda.Event() for _ in range(POSTPROC_WORKERS)] event_queue = queue.Queue() def event_monitor(event_pool, event_queue, output_queue): while True: item = event_queue.get() si el elemento es Ninguno: romper batch_id, buf_idx = item event_pool[buf_idx].synchronize() salida_queue.put((batch_id, buf_idx)) event_queue.task_done() monitor = threading.Thread(target=event_monitor, args=(event_pool, event_queue, salida_queue)) monitor.start()

La secuencia de inferencia actualizada consta de los siguientes pasos:

Obtenga un lote de entrada que se capturó previamente en la GPU. Ejecute el modelo en el lote de entrada para obtener un tensor de salida en la GPU. Solicite un búfer de CPU vacante de la cola de búfer y utilícelo para activar una copia de datos asincrónica. Configure un evento para que se active cuando se complete la copia y envíe el evento a la cola de eventos. El hilo del monitor espera a que se active el evento y luego envía el tensor de salida a la cola de salida para su procesamiento. Un subproceso de trabajo extrae el tensor de salida de la cola y lo guarda en el disco. Luego libera el búfer a la cola del búfer.

El código actualizado aparece a continuación.

def sincronizar_todo(): torch.cuda.synchronize() event_queue.join() salida_queue.join() con torch.inference_mode(): para i en rango(TOTAL_STEPS): if i == WARMUP_STEPS: sincronizar_todos() hora_inicio = time.perf_counter() perfilador.start() elif i == PASOS_CALENTAMIENTO + PASOS_PERFIL: sincronizar_todos() perfilador.stop() hora_final = time.perf_counter() con nvtx.annotate(f"Batch {i}", color="blue"): con nvtx.annotate("get lote", color="red"): lote = next(data_iter) con nvtx.annotate("compute", color="green"): salida = modelo(batch) con nvtx.annotate("copiar a CPU", color="amarillo"): salida_cpu, buf_id = to_cpu(output['out']) con nvtx.annotate("queue CUDA event", color="cyan"): event_pool[buf_id].record() event_queue.put((i, buf_id)) total_time = end_time – start_time rendimiento = PROFILE_STEPS / total_time print(f"Rendimiento: {rendimiento:.2f} pasos/seg") # limpieza event_queue.put(Ninguno) para _ en el rango(POSTPROC_WORKERS): salida_queue.put(Ninguno)

El rendimiento resultante es de 1,55 pasos por segundo.

El nuevo seguimiento del perfil aparece a continuación:

Cronología de Nsight Systems Profiler de transferencia de datos asíncrona (por autor)

En la fila NVTX de la sección CPU podemos ver todas las operaciones en el bucle de inferencia agrupadas en el lado izquierdo, lo que implica que todas se ejecutaron de forma inmediata y asincrónica. También vemos las llamadas de sincronización de eventos (en verde claro) ejecutándose en el hilo del monitor dedicado. En la sección GPU vemos que el cálculo del kernel comienza inmediatamente después de que se completa la copia del dispositivo al host.

Nuestra optimización final se centrará en mejorar la paralelización del kernel y las operaciones de memoria en la GPU.

Optimización 4: canalización mediante flujos CUDA

Como en nuestra publicación anterior, deseamos aprovechar los motores independientes para la copia de memoria (el DMA) y el cálculo del kernel (los SM). Hacemos esto asignando la copia de memoria a una secuencia CUDA dedicada:

egress_stream = torch.cuda.Stream() con torch.inference_mode(): para i en rango(TOTAL_STEPS): if i == WARMUP_STEPS: sincronizar_todos() start_time = time.perf_counter() perfilador.start() elif i == WARMUP_STEPS + PROFILE_STEPS: sincronizar_todos() perfiler.stop() end_time = time.perf_counter() con nvtx.annotate(f"Batch {i}", color="blue"): con nvtx.annotate("get batch", color="red"): lote = next(data_iter) con nvtx.annotate("compute", color="green"): salida = model(batch) # en una secuencia separada con torch.cuda.stream(egress_stream): # espera a que la secuencia predeterminada complete el cálculo egress_stream.wait_stream(torch.cuda.default_stream()) con nvtx.annotate("copiar a CPU", color="amarillo"): salida_cpu, buf_id = to_cpu(salida['out']) con nvtx.annotate("cola de evento CUDA", color="cian"): event_pool[buf_id].record(egress_stream) event_queue.put((i, buf_id))

Esto da como resultado un rendimiento de 1,85 pasos por segundo, una mejora adicional del 19,3 % con respecto a nuestro experimento anterior.

El seguimiento del perfil final aparece a continuación:

Cronología del perfilador de sistemas Pipelined Nsight (por autor)

En la sección de GPU vemos un bloque continuo de cálculo del kernel (en azul claro) con el host al dispositivo (en verde claro) y el dispositivo al host (en violeta) ejecutándose en paralelo. Nuestro bucle de inferencia ahora está vinculado a la computación, lo que implica que hemos agotado todas las oportunidades prácticas para la optimización de la transferencia de datos.

Resultados

Resumimos nuestros resultados en la siguiente tabla:

Resultados del experimento (por autor)

Mediante el uso de nsys Profiler pudimos aumentar la eficiencia en más de 4 veces. Naturalmente, el impacto de las optimizaciones que analizamos variará según los detalles del modelo y el entorno de ejecución.

Resumen

Con esto concluye la segunda parte de nuestra serie de publicaciones sobre el tema de la optimización de la transferencia de datos en cargas de trabajo de IA/ML. La primera parte se centró en las copias de host a dispositivo y la segunda parte en las copias de dispositivo a host. Cuando se implementa de manera ingenua, la transferencia de datos en cualquier dirección puede provocar importantes cuellos de botella en el rendimiento, lo que resulta en una escasez de GPU y un aumento de los costos de tiempo de ejecución. Utilizando el perfilador de Nsight Systems, demostramos cómo identificar y resolver estos cuellos de botella y aumentar la eficiencia del tiempo de ejecución.

Aunque la optimización de ambas direcciones implicó pasos similares, los detalles de implementación fueron muy diferentes. Si bien la optimización de la transferencia de datos de CPU a GPU está bien respaldada por las API de carga de datos de PyTorch y requirió cambios relativamente pequeños en el ciclo de ejecución, optimizar la dirección de GPU a CPU requirió un poco más de ingeniería de software. Es importante destacar que las soluciones que presentamos en esta publicación fueron elegidas con fines demostrativos. Su propia solución puede diferir considerablemente según las necesidades de su proyecto y sus preferencias de diseño.

Habiendo cubierto las copias de datos de CPU a GPU y de GPU a CPU, centramos nuestra atención en las transacciones de GPU a GPU: Estén atentos a una publicación futura sobre el tema de la optimización de la transferencia de datos entre GPU en cargas de trabajo de entrenamiento distribuidas.