Como compartir grandes Datasets entre procesos
sin perder la salud mental
Juan Francisco Huete Verdejo
Planteamiento del Problema
Step 1
df_A
Step 2
df_A
df_B
Step 3
df_C
df_D
Step 4
df_A
df_B
df_E
df_F
df_A
PIPELINE
3 s
15 s
12 s
20 s
Planteamiento del Problema
Step 1
df_A
Step 2
df_A
df_B
Step 3
df_C
df_D
Step 4
df_A
df_B
df_E
df_F
df_A
PIPELINE
3 s
15 s
12 s
20 s
Planteamiento del Problema
Step 1
df_A
Step 2
df_A
df_B
Step 3
df_C
df_D
Step 4
df_A
df_B
df_E
df_F
df_A
PIPELINE
3 s
15 s
12 s
20 s
15 s
Planteamiento del Problema
Step 1
df_A
Step 2
df_A
df_B
Step 3
df_C
df_D
Step 4
df_A
df_B
df_E
df_F
df_A
PIPELINE
3 s
15 s
12 s
20 s
15 s
Planteamiento del Problema
Celery worker
Step 1
df_A
Celery worker
Step 2
df_A
df_B
Celery worker
Celery worker
Step 3
df_C
df_D
Step 4
df_A
df_B
df_E
df_F
df_A
PIPELINE
3 s
15 s
12 s
20 s
Comienzan los errores
Celery worker
Step 1
df_A
Celery worker
Step 2
df_A
df_B
Celery worker
Celery worker
Step 3
df_C
df_D
Step 4
df_A
df_B
df_E
df_F
df_A
PIPELINE
3 s
15 s
12 s
20 s
Error
25 s
Usamos pickle y pasamos los df por
argumento a los workers
Comienzan los errores
Celery worker
Step 1
df_A
Celery worker
Step 2
df_A
df_B
Celery worker
Celery worker
Step 3
df_C
df_D
Step 4
df_A
df_B
df_E
df_F
df_A
PIPELINE
3 s
15 s
12 s
20 s
Error
25 s
Usamos pickle y pasamos los df por
argumento a los workers
Comienzan los errores
Celery worker
Step 1
df_A
Celery worker
Step 2
df_A
df_B
Celery worker
Celery worker
Step 3
df_C
df_D
Step 4
df_A
df_B
df_E
df_F
df_A
PIPELINE
3 s
15 s
12 s
20 s
50s
60 s
Guardamos los df en disco y leemos en los workers
Comienzan los errores
Algunas Alternativas
/ Redis
/ PyArrow + Plasma
/ Vaex
/ Vaex +S3
Redis
/ Base de datos en memoria
/ Clave valor
Inconvenientes
/ Limitación del tamaño de los valores a 500MB
/ Serialización usando pickle es lenta respecto a otras serializaciones
/ Muy fácil de configurar
/ Muy versátil
/ Para pequeñas ETL con df pequeños (< 500MB) es una solución rápida y simple
Ventajas
Redis
Celery worker
Step 1
df_A
Celery worker
Step 2
df_A
df_B
Celery worker
Step 3
df_C
df_D
df_E
df_F
df_A
key_A
key_B
Redis
key_A: <df_A>
key_A: <df_A>
key_A
Redis
Transferencia
Lectura
Benchmark
Pyarrow + Plasma
/ Arrow es un formato de datos en memoria en columnas
/ Facilita la conversión de los datos entre diferentes lenguajes
/ Plasma es un almacén de objetos arrow en memoria
Inconvenientes
/ Necesita un filesystem
/ Solo compatible con linux y MacOS
/ Serialización muy rápida
/ Transmisión muy rápida
/ Compatibilidad con df de muchos lenguajes
Ventajas
Pyarrow + Plasma
Instalación de Plasma
Plasma viene de serie en el paquete de pyarrow.
$ pip install pyarrow
Demonio de Plasma
$ plasma_store -m 8000000000 -s /tmp/plasma
-m Número de bytes reservados en memoria
-s Path del socket
Pyarrow + Plasma
Celery worker
Step 1
df_A
Celery worker
Step 2
df_A
df_B
Celery worker
Step 3
df_C
df_D
df_E
df_F
df_A
ObjectId_A
ObjectId_B
PLASMA
ObjectId_A: <df_A>
ObjectId_A: <df_A>
ObjectId_A
Pyarrow + Plasma
Transferencia
Lectura
Benchmark
Pyarrow + Plasma
Transferencia
Pyarrow + Plasma
Transferencia
Lectura
Benchmark
Pyarrow + Plasma
Transferencia
Lectura
Benchmark
MIS
DIESES
Vaex
/ Mapeo de datos en memoria
/ Lectura ultra rápida
/ Alternativa a pandas
/ Cacheo
Inconvenientes
/ Api similar a pandas
/ Alternativa a pandas
/ Operaciones muy rápidas (puede calcular medias, sumas, conteos de columnas de 10⁹ filas por segundo)
/ Integración con cloud
Ventajas
/ No está tan madura como pandas
/ No hay tanta comunidad como pandas
Vaex
Celery worker
Step 1
df_A
Celery worker
Step 2
df_A
df_B
Celery worker
Step 3
df_C
df_D
df_E
df_F
df_A
df_A.hdf5
HDD
df_A.hdf5
df_A.hdf5
df_B.hdf5
Vaex
Transferencia
Lectura
Benchmark
Vaex + S3
/ Potencia IO de Vaex
/ Orientado a la nube
Inconvenientes
/ Permite computación distribuida
/ Aprovechamos la concurrencia de s3
Ventajas
/ Vaex no está maduro
/ Más lento que otras opciones
Vaex + S3
Celery worker
Step 1
df_A
Celery worker
Step 2
df_A
df_B
Celery worker
Step 3
df_C
df_D
df_E
df_F
df_A
df_A.chunk1.hdf5
df_A.chunk2.hdf5
df_A.chunk3.hdf5
S3
df_A.chunk1.hdf5
df_A.chunk1.hdf5
df_A.chunk2.hdf5
df_A.chunk3.hdf5
df_A.chunk2.hdf5
df_A.chunk3.hdf5
df_B.chunk1.hdf5
df_B.chunk2.hdf5
df_B.chunk3.hdf5
Vaex + S3
Benchmark
Código de este ejemplo en el repositorio Git:
https://github.com/jfhuete/pycones_2021_compartir_grandes_datasets_entre_procesos
Buscamos Gente!
No lo dudes y contacta conmigo: jf.huete@bluetab.net o por Discord jfhuete#3999