【问题标题】:Parallelize MLflow Project runs with Pandas UDF on Azure Databricks Spark并行化 MLflow 项目在 Azure Databricks Spark 上使用 Pandas UDF 运行
【发布时间】:2023-01-23 12:08:31
【问题描述】:

我试着使用 Azure Databricks 上的 Spark 并行训练多个时间序列.
除了培训,我还想使用 MLflow 记录指标和模型.

代码结构很简单(基本改编自this example)。

  1. Databricks 笔记本触发 MLflow 项目
    mlflow.run(
        uri="/dbfs/mlflow-project",
        parameters={"data_path": "dbfs:/data/", "experiment_name": "test"}, 
        experiment_id=575501044793272,
        use_conda=False,
        backend="databricks",
        backend_config={
            "new_cluster": {
                "spark_version": "9.1.x-cpu-ml-scala2.12",
                "num_workers": 8,
                "node_type_id": "Standard_DS4_v2",
            },
            "libraries": [{"pypi": {"package": "pyarrow"}}]
        },
        synchronous=False
    )
    
    1. 主函数被调用.它基本上执行三个步骤:

      1. 读取由数据路径假如
      2. 定义一个触发“火车入口”MLflow项目
      3. 将此函数作为 Pandas UDF 应用于 Spark DataFrame

      这里的代码:

      sc = sparkContext('local')
      spark = SparkSession(sc)
      
      @click.argument("data_path")
      @click.argument("experiment_name")
      def run(data_path: str, experiment_name: str):
                  
          df = spark.read.format("delta").load(f"{data_path}")
          result_schema = StructType([StructField("key", StringType())])
      
          def forecast(data: pd.DataFrame) -> pd.DataFrame:
              child_run = client.create_run(
                  experiment_id=experiment,
                  tags={MLFLOW_PARENT_RUN_ID: parent_run_id},
              )
              p = mlflow.projects.run(
                  run_id=child_run.info.run_id, 
                  uri=".",
                  entry_points="train",
                  parameters={"data": data.to_json(), "run_id": child_run.info.run_id}, 
                  experiment_id=experiment,
                  backend="local",
                  usa_conda=False,
                  synchronous=False,
              )
      
              # Just a placeholder to use pandas UDF
              out = pd.DataFrame(data={"key": ["1"]})
              return out
      
          client = MLflowClient()
          experiment_path = f"/mlflow/experiments/{experiment_name}"
          experiment = client.create_experiment(experiment_path)
      
          parent_run = client.create_run(experiment_id=experiment)
          parent_run_id = parent_run.run_id
      
          # Apply pandas UDF (count() used just to avoid lazy evaluation)
          df.groupBy("key").applyInPandas(forecast, result_schema).count()
      
      1. 在每个键上调用训练函数.
        这基本上为每个时间序列(即每个键)训练一个 Prophet 模型,同时记录参数和模型。

      从集群 stderr 和 stdout 我可以看到 pandas UDF 被正确应用,因为它根据“关键”列正确划分了整个数据,即一次处理一个时间序列。

      问题是监控集群使用情况 仅使用一个节点,驱动程序节点:工作未分配给可用的工作人员,尽管 pandas UDF 似乎已正确应用。

      这里可能是什么问题? 我可以提供更多细节吗?

      非常感谢您, 马特奥

【问题讨论】:

    标签: apache-spark pyspark azure-databricks mlflow pandas-udf


    【解决方案1】:

    看来您需要重新分区输入数据框。否则 spark 将看到单个分区数据帧并进行相应处理。

    【讨论】:

      猜你喜欢
      • 2021-11-28
      • 2018-07-02
      • 2022-01-23
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-01-13
      • 2022-11-11
      • 2020-09-10
      相关资源
      最近更新 更多