You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PromptFlow流水线Serverless连接异常及MLFlow指标日志问题

问题解决方案

1. 修复Serverless连接异常

根据报错信息「Unknown connection... category Serverless, please upgrade your promptflow sdk version」,核心原因是当前使用的PromptFlow 1.18.0版本不支持Serverless类型连接,按以下步骤修复:

  • 升级PromptFlow SDK:在本地环境和流水线运行环境中升级到支持Serverless连接的版本(建议升级到最新稳定版),执行命令:
    pip install --upgrade promptflow promptflow-azureml
    
  • 更新环境依赖:修改requirements.txt文件,指定升级后的PromptFlow版本,例如:
    promptflow>=1.28.0
    promptflow-azureml>=1.28.0
    
  • 重新注册Flow组件:升级SDK后,重新执行ml_client.components.create_or_update(flow_component),确保流水线使用的是新版本的组件定义
  • 验证连接配置:确认flow.dag.yaml中provider: Serverless配置正确,连接名称Llama-4-Scout-17B-16E-Instruct-j为工作区级别且在PromptFlow UI中可正常访问

2. MLFlow追踪实验指标及结果聚合

假设连接问题已修复,按以下步骤实现结果汇总、聚合并记录MLFlow指标:

步骤1:创建聚合组件

新建aggregate_results.py文件,用于读取PromptFlow节点输出并执行聚合逻辑:

import mlflow
import json
import os

def aggregate(input_dir):
    # 读取子节点输出的JSONL文件
    results = []
    for filename in os.listdir(input_dir):
        if filename.endswith(".jsonl"):
            with open(os.path.join(input_dir, filename), "r", encoding="utf-8") as f:
                for line in f:
                    results.append(json.loads(line))
    
    # 示例聚合逻辑:统计输出内容的平均长度
    total_length = sum(len(item["joke"]) for item in results)
    avg_length = total_length / len(results) if results else 0
    
    # 记录MLFlow指标
    mlflow.log_metric("average_joke_length", avg_length)
    # 可添加更多自定义指标,如关键词出现次数等
    
    # 保存聚合结果(可选)
    with open("aggregated_results.json", "w", encoding="utf-8") as f:
        json.dump({"average_length": avg_length}, f)
    
    return {"aggregated_result": "aggregated_results.json"}

if __name__ == "__main__":
    import argparse
    parser = argparse.ArgumentParser()
    parser.add_argument("--input_dir", type=str, required=True)
    parser.add_argument("--output_dir", type=str, required=True)
    args = parser.parse_args()
    
    # 流水线运行时自动继承工作区MLFlow配置,本地测试需手动设置tracking URI
    mlflow.set_tracking_uri(os.environ.get("MLFLOW_TRACKING_URI"))
    
    with mlflow.start_run():
        result = aggregate(args.input_dir)
        os.makedirs(args.output_dir, exist_ok=True)
        with open(os.path.join(args.output_dir, "result.json"), "w") as f:
            json.dump(result, f)

步骤2:定义聚合组件的YAML配置(aggregate_component.yaml)

$schema: https://azuremlschemas.azureedge.net/latest/component.schema.json
name: aggregate_promptflow_results
version: 1
type: command
inputs:
  input_dir:
    type: uri_folder
outputs:
  output_dir:
    type: uri_folder
code: ./
environment:
  python_requirements_txt: requirements.txt
command: >-
  python aggregate_results.py
  --input_dir ${{inputs.input_dir}}
  --output_dir ${{outputs.output_dir}}

对应的requirements.txt需包含:

mlflow>=2.0.0
azureml-mlflow>=1.50.0

步骤3:修改流水线定义(orchestrator.py)

加载聚合组件并添加到流水线中,关联PromptFlow节点的输出:

# 新增:加载聚合组件
aggregate_component = load_component(source="aggregate_component.yaml")
ml_client.components.create_or_update(aggregate_component, version="1")

@pipeline()
def eval_pipeline():
    eval_node = flow_component(
        data=eval_data,
        topic="${data.topic}"
    )
    eval_node.compute                      = cluster_name
    eval_node.max_concurrency_per_instance = 1
    eval_node.mini_batch_error_threshold   = 5

    # 新增:添加聚合节点,关联PromptFlow输出
    aggregate_node = aggregate_component(
        input_dir=eval_node.outputs.joke
    )
    aggregate_node.compute = cluster_name

pipeline_job = eval_pipeline()
# 其余配置保持不变

步骤4:查看实验指标

提交流水线后,在Azure ML工作区的「实验」页面找到对应MLFlow实验,即可查看记录的average_joke_length等自定义指标。

内容的提问来源于stack exchange,提问作者Lopez

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.12 22:14:50