Tengo el siguiente código que ejecuta dos capacitaciones de TensorFlow en paralelo con trabajadores de Dask implementados en contenedores de Docker.
Necesito lanzar dos procesos, usando el mismo cliente dask, donde cada uno entrenará sus respectivos modelos con N trabajadores.
Para ello hago lo siguiente:
joblib.delayed para generar los dos procesos.with joblib.parallel_backend('dask'): para ejecutar la lógica de ajuste/entrenamiento. Cada proceso de formación desencadena N dask trabajadores.El problema es que no sé si todo el proceso es seguro para subprocesos, ¿hay algún elemento de concurrencia que me falte?
# First, submit the function twice using joblib delay delayed_funcs = [joblib.delayed(train)(sub_task) for sub_task in [123, 456]] parallel_pool = joblib.Parallel(n_jobs=2) parallel_pool(delayed_funcs) # Second, submit each training process def train(sub_task): global client if client is None: print('connecting') client = Client() data = some_data_to_train # Third, process the training itself with N workers with joblib.parallel_backend('dask'): X = data[columns] y = data[label] niceties = dict(verbose=False) model = KerasClassifier(build_fn=build_layers, loss=tf.keras.losses.MeanSquaredError(), **niceties) model.fit(X, y, epochs=500, verbose = 0)Esto es pura especulación, pero un posible problema de simultaneidad se debe a que if client is None: parte, donde dos procesos podrían correr para crear un Client .
Si esto se resuelve (por ejemplo, mediante la creación explícita de un cliente por adelantado), entonces el programador de dask se basará en el tiempo de envío para priorizar la tarea (a menos que priority esté claramente asignada) y también en la estructura del gráfico (DAG), hay más detalles disponibles en docs .
La pregunta, tal como se da, podría marcarse fácilmente como "poco clara" para SO. Un par de notas:
global client : hace que el objeto del cliente esté disponible fuera de la función. Pero la función se ejecuta desde otro proceso, no afecta el otro proceso al hacer que el clienteif client is None : este es un error de nombre, su código en realidad no se ejecuta como está escritoclient = Client() : crea un nuevo clúster en cada subproceso, cada uno asumiendo el total de recursos disponibles, sobresuscribiendo esos recursos.Debe preguntarse: ¿por qué está creando procesos para los dos ajustes? ¿Por qué no dejar que Dask descubra su paralelismo, que es para lo que está destinado?
--
-EDITAR-
para responder a la forma de la pregunta formulada en un comentario.
Mi pregunta es si usar la misma variable de cliente en estos dos procesos paralelos crea un problema.
No, las dos variables de client no están relacionadas entre sí. Es posible que vea un mensaje de advertencia sobre la imposibilidad de vincularse a un puerto predeterminado, que puede ignorar con seguridad. Sin embargo, no lo hagas global ya que esto es innecesario y hace que lo que estás haciendo sea menos claro.
--
Creo que debo responder la pregunta tal como está redactada en su comentario, que aconsejo agregar a la pregunta principal
Necesito lanzar dos procesos, usando el mismo cliente dask, donde cada uno entrenará sus respectivos modelos con N trabajadores.
Tienes las siguientes opciones:
Client() y obtenga su dirección (por ejemplo, client._scheduler_identity['address'] ) y conéctese a eseclient.write_scheduler_file y utilíceloTe conectarás en la función con
client = Client(address)o
client = Client(scheduler_file=the_file_you_wrote)