使用DynamicFrame.fromDF后Glue写入SQL Server性能骤降问题排查
AWS Glue作业SQL Server读写性能优化问题
我有一个AWS Glue作业,需从SQL Server读取2张表,执行关联/转换操作后写回SQL Server的新表或已截断表,待写入数据量约15GB。尝试了两种实现方案,性能差异极大,目标是将作业时长控制在10分钟内。
方案1 - 总耗时约17分钟
流程:从SQL Server读取数据→转换→写入S3临时存储→从S3读取→写回SQL Server
- 从SQL Server读取数据到Spark DataFrame(约3-5秒)
- 在Spark DataFrame上执行转换操作(约5秒)
- 将数据写入S3临时存储(约8分钟)
- 使用
glueContext.create_dynamic_frame.from_options()从S3读取数据为DynamicFrame - 使用
glueContext.write_from_options()写入SQL Server表(约9分钟)
方案1代码模板
import sys from awsglue.transforms import * from awsglue.utils import getResolvedOptions from pyspark.context import SparkContext from awsglue.context import GlueContext from awsglue.dynamicframe import DynamicFrame from awsglue.job import Job from pyspark.sql.types import * from pyspark.sql.functions import * import time from py4j.java_gateway import java_import ## @params: [JOB_NAME] args = getResolvedOptions(sys.argv, ['JOB_NAME']) sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) job.init(args['JOB_NAME'], args) logger = glueContext.get_logger() # 步骤1 -- 从表中读取数据到DataFrame # ----------------------------------------------- # 步骤2 -- 执行转换(如有),并写入数据湖S3 #---------------------------------------------------------------------- df.write.mode("overwrite").csv("s3://<<bucket-name>>/temp_result") # 步骤3 -- 从S3读取数据到新的DataFrame #------------------------------------------------ newdf = glueContext.create_dynamic_frame.from_options( connection_type='s3', connection_options={"paths": ["s3://<<bucket-name>>/temp_result"]}, format='csv' ) # 步骤4 -- 截断目标表(因为目标表始终是全量刷新) #--------------------------------------------------------------------------------- cstmt = conn.prepareCall("TRUNCATE TABLE mytable_in_db"); results = cstmt.execute(); # 步骤5 -- 从DataFrame写入目标表 # ---------------------------------------------- glueContext.write_from_options( frame_or_dfc=newdf, connection_type="sqlserver", connection_options=connection_sqlserver_options ) job.commit()
方案2 - 总耗时约50分钟
流程:从SQL Server读取数据→转换→直接写回SQL Server
- 从SQL Server读取数据到Spark DataFrame(约3-5秒)
- 在Spark DataFrame上执行转换操作(约5秒)
- 使用
DynamicFrame.fromDF()将Spark DataFrame转换为DynamicFrame - 使用
glueContext.write_from_options()写入SQL Server表(约43分钟)
方案2代码模板
import sys from awsglue.transforms import * from awsglue.utils import getResolvedOptions from pyspark.context import SparkContext from awsglue.context import GlueContext from awsglue.dynamicframe import DynamicFrame from awsglue.job import Job from pyspark.sql.types import * from pyspark.sql.functions import * import time from py4j.java_gateway import java_import ## @params: [JOB_NAME] args = getResolvedOptions(sys.argv, ['JOB_NAME']) sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) job.init(args['JOB_NAME'], args) logger = glueContext.get_logger() # 步骤1 -- 从表中读取数据到DataFrame # ----------------------------------------------- # 步骤2 -- 执行转换(如有)并存储到df #---------------------------------------------------------------------- # df contains transformed data # 步骤3 -- 将Spark DataFrame转换为DynamicFrame #-------------------------------------------------------- newdf2 = DynamicFrame.fromDF(df, glueContext , "newdf2") # 步骤4 -- 截断目标表(因为目标表始终是全量刷新) #--------------------------------------------------------------------------------- cstmt = conn.prepareCall("TRUNCATE TABLE mytable_in_db"); results = cstmt.execute(); # 步骤5 -- 从DataFrame写入目标表 # ---------------------------------------------- glueContext.write_from_options( frame_or_dfc=newdf2, connection_type="sqlserver", connection_options=connection_sqlserver_options ) job.commit()
方案2耗时更长的原因
- 分区并行度不足:方案2中,转换后的DynamicFrame继承了原DataFrame的分区数,通常默认分区数较少,导致写入SQL Server时只有少量并发连接,无法利用SQL Server的并行写入能力。而方案1写入S3再读取时,Spark会根据S3文件数量自动重分区,提升并行度。
- 序列化效率差异:直接从DataFrame转DynamicFrame写入时,序列化逻辑未经过优化;而通过S3中转(尤其是用列式存储格式)时,数据序列化更高效,降低写入时的CPU/IO开销。
- Glue写入机制的差异:Glue的
write_from_options针对从S3读取的DynamicFrame有更优的并行调度,而对转换自DataFrame的DynamicFrame未充分利用集群资源。
优化建议
方案2针对性优化
- 手动调整分区数:转换为DynamicFrame前,对DataFrame重分区,根据集群核心数设置合理值(比如每个核心对应2-4个分区):
df = df.repartition(100) # 示例值,根据集群规模调整 newdf2 = DynamicFrame.fromDF(df, glueContext , "newdf2") - 改用Spark原生JDBC写入:放弃DynamicFrame,直接用Spark DataFrame的JDBC写入,配置并行参数:
df.write \ .mode("overwrite") \ .jdbc( url="jdbc:sqlserver://<server>:<port>;databaseName=<db>", table="mytable_in_db", properties={ "user": "<user>", "password": "<password>", "driver": "com.microsoft.sqlserver.jdbc.SQLServerDriver", "numPartitions": 100, # 并行写入的分区数 "batchsize": 10000, # 每批次写入行数 "rewriteBatchedStatements": "true" # 开启批量语句优化 } ) - 优化JDBC参数:
batchsize建议设置为10000-50000,减少网络交互;开启rewriteBatchedStatements可大幅提升SQL Server批量写入效率。
方案1进一步优化(压缩至10分钟内)
- 替换CSV为Parquet格式:Parquet是列式存储,读写速度更快且支持压缩,减少S3传输和存储时间:
# 写入S3 df.write.mode("overwrite").parquet("s3://<<bucket-name>>/temp_result") # 读取S3 newdf = glueContext.create_dynamic_frame.from_options( connection_type='s3', connection_options={"paths": ["s3://<<bucket-name>>/temp_result"]}, format='parquet' ) - 提前重分区:写入S3前对DataFrame重分区,避免生成过多小文件,同时保证并行度:
df.repartition(80).write.mode("overwrite").parquet("s3://<<bucket-name>>/temp_result") - 合并截断与写入:使用Spark JDBC的
mode("overwrite")可自动截断表,省去手动执行TRUNCATE的步骤。
通用优化
- 升级Glue集群配置:使用更大的Worker类型(如G.2X)和更多Worker数量,15GB数据建议至少配置10个G.2X Worker,提升计算和IO能力。
- 检查SQL Server端瓶颈:确保SQL Server的最大连接数足够(匹配Glue的并行分区数),磁盘IO和网络带宽充足,避免成为写入瓶颈。
- 减少无意义转换:如果不需要Glue的Schema Registry等特性,直接使用Spark DataFrame操作,避免DataFrame与DynamicFrame之间的转换开销。
内容的提问来源于stack exchange,提问作者Yuva
相关产品推荐
相关产品推荐

