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

从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每日调度配置

  1. 进入Synapse Studio的集成页面,新建管道
  2. 在管道中添加Notebook活动,选择需要调度的目标Notebook
  3. 点击管道顶部的触发 → 新建/编辑触发器,选择日程触发
  4. 设置触发周期为「每天」,指定执行时间,保存并发布管道
  5. 手动触发测试,确认调度逻辑正常运行

四、常见问题排查

  • pyodbc连接失败:检查ODBC驱动版本(推荐v18)、Azure MI防火墙是否允许Synapse工作区IP访问、账号密码正确性
  • JDBC托管身份认证失败:确认Synapse托管身份已被添加为Azure MI的数据库用户,并分配了对应权限
  • 写入性能问题:调整Spark并行度(spark.sql.shuffle.partitions),或采用批量写入方式减少连接开销

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 18:57:51