从Azure Synapse Analytics Spark池连接到Azure MI的问题
Synapse Notebook连接Azure MI实现数据迁移及调度方案
一、Azure MI连接方案(替代TokenLibrary/mssparkutils)
由于Synapse链接服务未内置Azure MI支持,无法直接通过TokenLibrary或mssparkutils获取连接串,可采用以下两种直接连接方式:
1. pyodbc连接(SQL账号认证)
- 先安装依赖库:
!pip install pyodbc - 连接及操作示例代码:
import pyodbc # 配置Azure MI参数 server = "你的MI服务器名.database.windows.net" database = "目标数据库名" username = "SQL认证账号" password = "SQL认证密码" driver = "{ODBC Driver 18 for SQL Server}" # 建立连接 conn = pyodbc.connect( f"DRIVER={driver};SERVER={server};DATABASE={database};UID={username};PWD={password};" "Encrypt=yes;TrustServerCertificate=no;Connection Timeout=30;" ) cursor = conn.cursor() # 测试查询 cursor.execute("SELECT TOP 5 * FROM 测试表") for row in cursor.fetchall(): print(row) # 写入数据示例(以ADLS Gen2的Delta表数据为例) delta_df = spark.read.format("delta").load("abfss://容器名@ADLS账号.dfs.core.windows.net/Delta表路径") for row in delta_df.collect(): cursor.execute( "INSERT INTO 目标MI表(col1, col2) VALUES (?, ?)", row["col1"], row["col2"] ) conn.commit() # 关闭资源 cursor.close() conn.close()
2. JDBC连接(托管身份认证)
适合用Synapse工作区托管身份做无密码认证,需先给托管身份分配Azure MI的数据库权限(如db_datawriter、db_datareader):
from pyspark.sql import SparkSession spark = SparkSession.builder.appName("MI_JDBC_Connect").getOrCreate() # JDBC连接配置 jdbc_url = "jdbc:sqlserver://你的MI服务器名.database.windows.net:1433;" \ "databaseName=目标数据库名;encrypt=true;trustServerCertificate=false;" \ "hostNameInCertificate=*.database.windows.net;loginTimeout=30;" conn_props = { "driver": "com.microsoft.sqlserver.jdbc.SQLServerDriver", "authentication": "ActiveDirectoryMSI" } # 读取ADLS Gen2的Delta表并处理(关联/子查询) delta_df = spark.read.format("delta").load("abfss://容器名@ADLS账号.dfs.core.windows.net/Delta表路径") filtered_df = delta_df.filter("create_time >= DATEADD(day, -1, CURRENT_DATE())") \ .join(spark.table("关联表"), on="关联键", how="inner") # 写入Azure MI filtered_df.write.jdbc( url=jdbc_url, table="目标MI表", mode="append", properties=conn_props )
二、Delta表数据筛选处理
直接用Spark DataFrame API或SQL完成关联/子查询逻辑,示例如下:
# Spark SQL方式处理 spark.sql(""" SELECT d.id, d.content, o.tag FROM delta.`abfss://容器名@ADLS账号.dfs.core.windows.net/Delta表路径` d LEFT JOIN 其他Delta表 o ON d.id = o.ref_id WHERE d.update_time >= CURRENT_TIMESTAMP() - INTERVAL 1 DAY """).createOrReplaceTempView("filtered_data") processed_df = spark.table("filtered_data")
三、Notebook每日调度配置
- 进入Synapse Studio的集成页面,新建管道
- 在管道中添加Notebook活动,选择需要调度的目标Notebook
- 点击管道顶部的触发 → 新建/编辑触发器,选择日程触发
- 设置触发周期为「每天」,指定执行时间,保存并发布管道
- 手动触发测试,确认调度逻辑正常运行
四、常见问题排查
- pyodbc连接失败:检查ODBC驱动版本(推荐v18)、Azure MI防火墙是否允许Synapse工作区IP访问、账号密码正确性
- JDBC托管身份认证失败:确认Synapse托管身份已被添加为Azure MI的数据库用户,并分配了对应权限
- 写入性能问题:调整Spark并行度(
spark.sql.shuffle.partitions),或采用批量写入方式减少连接开销
内容的提问来源于stack exchange,提问作者jemi
相关产品推荐
相关产品推荐

