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

使用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耗时更长的原因

  1. 分区并行度不足:方案2中,转换后的DynamicFrame继承了原DataFrame的分区数,通常默认分区数较少,导致写入SQL Server时只有少量并发连接,无法利用SQL Server的并行写入能力。而方案1写入S3再读取时,Spark会根据S3文件数量自动重分区,提升并行度。
  2. 序列化效率差异:直接从DataFrame转DynamicFrame写入时,序列化逻辑未经过优化;而通过S3中转(尤其是用列式存储格式)时,数据序列化更高效,降低写入时的CPU/IO开销。
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 18:40:34