Azure ML运行并调度Python脚本后,Blob存储中CSV文件未随Snowflake数据更新的问题求助
Azure ML运行并调度Python脚本后,Blob存储中CSV文件未随Snowflake数据更新的问题求助
大家好,我最近遇到一个Azure ML的奇怪问题,想请教下各位:
我正在尝试在Azure ML中运行Python脚本并设置调度,核心功能是从Snowflake表读取数据,然后将数据保存为CSV文件存储到Azure Blob存储中。第一次运行一切正常,Blob里的CSV文件也正确生成了,但当我更新Snowflake表中的数据后,不管是手动重新执行管道代码,还是等待调度任务触发,脚本看起来都在正常运行(没有报错),但Blob里的CSV文件就是不更新。
更诡异的是,我试了几个不同的用例,发现只有当我修改Python脚本的内容时,脚本才会真正重新执行,Blob里的文件才会跟着更新。如果脚本内容不变,哪怕Snowflake数据已经变了,跑多少次都是老结果。
下面是我的相关代码,麻烦帮我看看哪里出问题了:
主管道代码
from azure.identity import DefaultAzureCredential from azureml.core import Workspace import os import azureml.core from azureml.core import Workspace, Dataset, Datastore, ComputeTarget, Experiment, ScriptRunConfig from azureml.pipeline.steps import PythonScriptStep from azureml.pipeline.core import Pipeline from azure.ai.ml.entities import AmlCompute from azureml.core import Environment ws = Workspace.from_config() subscription_id = "**" resource_group = "**" workspace = "****" ml_client = MLClient(DefaultAzureCredential(), subscription_id, resource_group, workspace) exp = Experiment(workspace=ws, name='date_test') compute_cluster = "datetime" # 这里我还尝试设置了allow_reuse=False,但似乎没起作用 prep_step = PythonScriptStep(name='prep step', script_name="datetimescript.py", # datetimescript.py包含核心业务逻辑 source_directory="dynamic", compute_target=compute_cluster, allow_reuse=False) job_env = Environment.from_conda_specification(name = 'datetime_testenvironment', file_path = './conda_train.yaml') # conda_train.yaml包含所有依赖包 train_src = ScriptRunConfig(source_directory="dynamic", script='datetimescript.py', compute_target=compute_cluster, environment=job_env) train_step = PythonScriptStep(name='train step', source_directory=train_src.source_directory, script_name=train_src.script, runconfig=train_src.run_config) mybatchpipeline = Pipeline(ws, steps=[train_step]) run = exp.submit(mybatchpipeline) run.wait_for_completion(show_output=True) # 我发现就算手动运行这段代码,结果也不会更新
核心脚本datetimescript.py
import numpy as np import pandas as pd import os from snowflake.snowpark.session import Session import json from azure.storage.blob import BlobClient from datetime import datetime, timezone connection_parameters = json.load(open(r"connection.json")) session = Session.builder.configs(connection_parameters).create() session.sql_simplifier_enabled = True session.sql('USE ROLE "MY_ROLE";').collect() df = pd.DataFrame(session.sql('select * from "MY_TABLE";').collect()) accountUrl="https://**.blob.core.windows.net" containerName="" blobName='tableSf.csv' credentials= "******************************" blob_client = BlobClient(account_url=accountUrl ,container_name= containerName, blob_name= blobName,credential=credentials ) output = df.to_csv(index=False, header=True, encoding="utf-8") # 原代码里的SKILL_MASTER应该是笔误,实际应为df blob_client.upload_blob(output, overwrite=True) curr_dt = datetime.now(timezone.utc) print("job successfull!") print(curr_dt)
我已经确认Snowflake的数据确实更新了,Blob上传参数也设置了overwrite=True,而且脚本运行时没有任何错误日志。有没有人遇到过类似的问题?或者知道Azure ML是不是有什么缓存机制导致脚本没有真正重新执行?
备注:内容来源于stack exchange,提问作者Adaline002
相关产品推荐
相关产品推荐

