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

如何让AWS Glue 4.0 PySpark作业并行抽取数据库表数据

AWS Glue 4.0 PySpark作业并行化优化

问题背景

现有AWS Glue 4.0 PySpark作业通过串行循环读取SQL Server多张表并写入S3 Parquet格式,代码可正常运行但串行执行效率低下,需要调整为并行执行以提升速度。

原作业代码:

import sys
from awsglue.transforms import *
from awsglue.utils import getResolvedOptions
from pyspark.context import SparkContext
from awsglue.context import GlueContext
from awsglue.job import Job

args = getResolvedOptions(sys.argv, ["JOB_NAME"])
sc = SparkContext()
glueContext = GlueContext(sc)
spark = glueContext.spark_session
job = Job(glueContext)
job.init(args["JOB_NAME"], args)

tables = ['table1','table2','table3','table4','table5']

for table in tables:
    dyf = glueContext.create_dynamic_frame.from_options(
        connection_type="sqlserver",
        connection_options={
            "useConnectionProperties": "true",
            "dbtable": table,
            "connectionName": "my-glue-jdbc-connection",
        }
    )
    glueContext.write_dynamic_frame.from_options(
        frame=dyf,
        connection_type="s3",
        format="parquet",
        connection_options={"path": f"s3://my-s3-bucket/{table}", "partitionKeys": []},
        format_options={"compression": "gzip"}
    )

解决方案

利用Python的concurrent.futures.ThreadPoolExecutor实现多线程并行处理每张表的读写操作,具体修改如下:

修改后的代码

import sys
from concurrent.futures import ThreadPoolExecutor
from awsglue.transforms import *
from awsglue.utils import getResolvedOptions
from pyspark.context import SparkContext
from awsglue.context import GlueContext
from awsglue.job import Job

args = getResolvedOptions(sys.argv, ["JOB_NAME"])
sc = SparkContext()
glueContext = GlueContext(sc)
spark = glueContext.spark_session
job = Job(glueContext)
job.init(args["JOB_NAME"], args)

tables = ['table1','table2','table3','table4','table5']

def process_table(table):
    # 读取SQL Server表为Dynamic Frame
    dyf = glueContext.create_dynamic_frame.from_options(
        connection_type="sqlserver",
        connection_options={
            "useConnectionProperties": "true",
            "dbtable": table,
            "connectionName": "my-glue-jdbc-connection",
        }
    )
    # 写入S3 Parquet
    glueContext.write_dynamic_frame.from_options(
        frame=dyf,
        connection_type="s3",
        format="parquet",
        connection_options={"path": f"s3://my-s3-bucket/{table}", "partitionKeys": []},
        format_options={"compression": "gzip"}
    )

# 初始化线程池,max_workers可根据资源调整(建议匹配表数量或合理并发数)
with ThreadPoolExecutor(max_workers=len(tables)) as executor:
    executor.map(process_table, tables)

job.commit()

关键说明

  1. 线程池并行处理:通过ThreadPoolExecutor创建线程池,将每张表的读写任务分配到独立线程并行执行,替代原有的串行循环。
  2. 任务封装:把单表的读取和写入逻辑封装为process_table函数,线程池自动调度该函数处理每个表。
  3. 并发数控制:max_workers参数控制最大并发线程数,建议根据Glue作业的DPU配置和SQL Server的负载能力调整,避免并发过高导致资源耗尽。

注意事项

  • 确认SQL Server数据库能够承受并行读取的压力,避免因并发连接数过高导致数据库性能下降。
  • 适当调整Glue作业的DPU数量(例如从默认2个提升至4-8个),确保有足够资源支撑并行任务。
  • Glue 4.0的GlueContext具备线程安全性,无需在每个线程内重新初始化上下文。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 10:25:26