如何通过AMLS Python SDK从Pipeline的Pandas DataFrame写入Azure SQL DB?
Absolutely, you can write a Pandas DataFrame back to Azure SQL Database from an Azure Machine Learning Pipeline using the Python SDK. There are two practical approaches depending on whether you want to handle the upsert directly in Python or use a temporary table followed by SQL operations. Let’s walk through both:
Approach 1: Direct Upsert via PythonScriptStep
You can use a PythonScriptStep to connect to Azure SQL (using pyodbc or SQLAlchemy) and perform upserts directly from your DataFrame. The key here is to store your SQL connection string securely in Azure ML’s Key Vault instead of hardcoding it.
Step 1: Store SQL Connection String as a Secret
First, register your connection string in your workspace’s Key Vault:
from azureml.core import Workspace ws = Workspace.from_config() ws.set_default_keyvault() ws.keyvault.set_secret(name="sql-connection-string", value="your_sql_connection_string")
Step 2: Write the Python Script for Upsert
Create a script (e.g., write_to_sql.py) that retrieves the secret, loads your DataFrame, and runs a MERGE statement to update existing rows and insert new ones:
import pandas as pd import pyodbc from azureml.core import Run def main(): run = Run.get_context() ws = run.experiment.workspace # Fetch connection string from Key Vault conn_str = ws.keyvault.get_secret("sql-connection-string") # Replace this with your actual DataFrame logic df = pd.DataFrame({ "id": [1, 2, 3], "customer_name": ["UpdatedName", "NewCustomer1", "NewCustomer2"] }) # Connect to SQL and execute upsert with pyodbc.connect(conn_str) as conn: cursor = conn.cursor() # Use MERGE for upsert (adjust schema/columns to match your table) for _, row in df.iterrows(): cursor.execute(""" MERGE INTO dbo.Customer AS target USING (VALUES (?, ?)) AS source (id, customer_name) ON target.id = source.id WHEN MATCHED THEN UPDATE SET target.customer_name = source.customer_name WHEN NOT MATCHED THEN INSERT (id, customer_name) VALUES (source.id, source.customer_name); """, row["id"], row["customer_name"]) conn.commit() if __name__ == "__main__": main()
Step 3: Define the Pipeline Step
Set up the environment and step to run your script:
from azureml.pipeline.steps import PythonScriptStep from azureml.core import Environment from azureml.core.runconfig import RunConfiguration # Create environment with required packages sql_env = Environment(name="sql-write-env") sql_env.python.conda_dependencies.add_pip_package("pandas") sql_env.python.conda_dependencies.add_pip_package("pyodbc") run_config = RunConfiguration() run_config.environment = sql_env # Define the pipeline step upsert_step = PythonScriptStep( script_name="write_to_sql.py", source_directory="./scripts", runconfig=run_config, compute_target=your_aml_compute_target, allow_reuse=False # Disable reuse if your DataFrame changes each run )
Approach 2: Temp Table + SQL Merge (Better for Large Datasets)
If you prefer handling the upsert via SQL (more efficient for large data), first write your DataFrame to a temporary SQL table using OutputTabularDatasetConfig, then run a SQL script to merge the temp table into your target table.
Step 1: Define Output to Temp SQL Table
Use your registered SQL datastore to set up the temporary table output:
from azureml.data import OutputTabularDatasetConfig sql_datastore = ws.datastores.get("your_registered_sql_datastore") temp_table_output = OutputTabularDatasetConfig( name="temp_customer_table", destination=(sql_datastore, "dbo.TempCustomer") ).register_on_complete(name="temp_customer_dataset")
Step 2: Write DataFrame to Temp Table
Create a script (e.g., write_temp_table.py) to write your DataFrame to the temp table:
import pandas as pd from azureml.core import Run def main(): run = Run.get_context() # Your DataFrame logic here df = pd.DataFrame({...}) # Write DataFrame to the temp table output output_dataset = run.output_datasets["temp_customer_table"] output_dataset.write_dataframe(df, overwrite=True) if __name__ == "__main__": main()
Step 3: Create the Pipeline Step for Temp Table
write_temp_step = PythonScriptStep( script_name="write_temp_table.py", source_directory="./scripts", arguments=["--output", temp_table_output], outputs=[temp_table_output], runconfig=run_config, compute_target=your_aml_compute_target )
Step 4: SQL Merge Step
Create a SQL script (e.g., merge_temp_to_target.sql) to upsert from temp to target:
MERGE INTO dbo.Customer AS target USING dbo.TempCustomer AS source ON target.id = source.id WHEN MATCHED THEN UPDATE SET target.customer_name = source.customer_name WHEN NOT MATCHED THEN INSERT (id, customer_name) VALUES (source.id, source.customer_name); -- Optional: Drop the temp table after merge DROP TABLE dbo.TempCustomer;
Define the SQL pipeline step:
from azureml.pipeline.steps import SqlScriptStep merge_step = SqlScriptStep( script_name="merge_temp_to_target.sql", source_directory="./scripts", sql_datastore=sql_datastore, compute_target=your_aml_compute_target, depends_on=[write_temp_step] # Ensure temp table is written first )
Key Notes
- Ensure your AML compute target has network access to Azure SQL (either via VNet integration or by whitelisting the compute’s public IP in SQL firewall).
- Always use Key Vault for sensitive credentials—never hardcode connection strings or passwords.
- For large datasets, the temp table + SQL merge approach is more performant than row-by-row upserts in Python.
内容的提问来源于stack exchange,提问作者bab689

