What is a good pattern for aggregating the results from Kubeflow Pipleine kfp.ParallelFor?
What is a good pattern for aggregating the results from Kubeflow Pipleine kfp.ParallelFor?
Not exactly what you asked for, but our workaround was to write the results of the parallelfor tasks into S3 and simply collect them afterwards in a postprocessing task.
with dsl.ParallelFor(preprocessing_task.output) as plant_item:
predict_plant='{}'.format(plant_item)
forecasting_task = forecasting_op(predict_plant, ....).after(preprocessing_task)
postprocessing_task = postprocessing_op(...).after(forecasting_task)
At the moment this might not be supported: