如何将外部数据源数据传输至Azure Synapse Analytics并使用Python实现ETL
用Python实现CRM数据到Azure Synapse Analytics的ETL流程完整方案
前置准备
- 先安装依赖包:通用依赖为
azure-synapse-spark、pyodbc、pandas、sqlalchemy;如果是Salesforce CRM额外装simple-salesforce,Dynamics CRM额外装dynamics365-ce-python,其他CRM自行替换对应官方SDK即可 - 提前准备两类凭证:CRM的API访问密钥/账号密码、Azure Synapse对应SQL池的连接字符串、存储账户访问密钥(走中间落数场景需要)
两种常用实现路径,根据数据量选型
路径1:小数据量(单表小于10G)直接写入专用SQL池
适合测试、小批量同步场景,脚本可以本地运行也可以部署到Azure Functions调度:
import pandas as pd from simple_salesforce import Salesforce import pyodbc from sqlalchemy import create_engine # 1. 连接CRM拉取源数据 sf = Salesforce(username="CRM账号", password="CRM密码", security_token="CRM安全令牌") # 按需替换为你的CRM查询语句,其他CRM替换为对应拉数逻辑即可 crm_query = "SELECT Id, Name, Account__c, CreatedDate FROM Opportunity" query_result = sf.query_all(crm_query) df = pd.DataFrame(query_result["records"]).drop("attributes", axis=1) # 2. 数据清洗转换(按需调整) df["CreatedDate"] = pd.to_datetime(df["CreatedDate"]) df = df.fillna(value="") # 空值处理适配Synapse表结构 # 3. 写入Synapse专用SQL池 synapse_conn_str = "DRIVER={ODBC Driver 18 for SQL Server};SERVER=<你的Synapse工作区名>-ondemand.sql.azuresynapse.net;DATABASE=<数据库名>;UID=<用户名>;PWD=<密码>;Encrypt=yes;TrustServerCertificate=no;Connection Timeout=30;" engine = create_engine(f"mssql+pyodbc:///?odbc_connect={synapse_conn_str}") # if_exists参数:replace覆盖原表、append追加数据、fail表存在时报错,按需选择 df.to_sql(name="opportunity_import", schema="dbo", con=engine, if_exists="append", index=False, chunksize=1000)
注意:chunksize根据单条数据大小调整,避免单次写入过大触发超时
路径2:大数据量(单表10G以上)走Synapse Spark池导入
生产环境推荐方案,支持TB级数据同步,性能比直接写入高5~10倍,脚本提交到Synapse Spark池运行即可:
from pyspark.sql import SparkSession from simple_salesforce import Salesforce import pandas as pd spark = SparkSession.builder.appName("CRM2Synapse").getOrCreate() # 1. 拉取CRM数据转Spark DataFrame sf = Salesforce(username="CRM账号", password="CRM密码", security_token="CRM安全令牌") crm_query = "SELECT Id, Name, Account__c, CreatedDate FROM Opportunity" query_result = sf.query_all(crm_query) pd_df = pd.DataFrame(query_result["records"]).drop("attributes", axis=1) spark_df = spark.createDataFrame(pd_df) # 2. Spark层做数据转换,性能远高于Pandas spark_df = spark_df.withColumn("CreatedDate", spark_df["CreatedDate"].cast("timestamp")) spark_df = spark_df.na.fill("") # 3. 直接写入Synapse表,默认用列存储索引存储性能最优 spark_df.write \ .synapsesql("dbo.opportunity_import", "append") \ .option("tableOptions", "CLUSTERED COLUMNSTORE INDEX") \ .mode("append") \ .save()
注意:提前在Synapse Spark池的库配置中上传CRM相关的依赖包,运行时会自动加载
常见问题排查
- 写入报类型不匹配:提前对齐DataFrame字段类型和Synapse表的字段类型,比如字符串统一用VARCHAR(n),日期用DATETIME2
- 写入超时:大数据量不要用直接写入专用SQL池的方案,切换到Spark池批量导入
- CRM拉数触发限流:给拉数逻辑加分页和重试机制,主流CRM的官方SDK基本都自带分页能力,特殊场景可以加间隔重试逻辑
内容的提问来源于stack exchange,提问作者Oğuz Kırçiçek
相关产品推荐
相关产品推荐

