如何让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()
关键说明
- 线程池并行处理:通过
ThreadPoolExecutor创建线程池,将每张表的读写任务分配到独立线程并行执行,替代原有的串行循环。 - 任务封装:把单表的读取和写入逻辑封装为
process_table函数,线程池自动调度该函数处理每个表。 - 并发数控制:
max_workers参数控制最大并发线程数,建议根据Glue作业的DPU配置和SQL Server的负载能力调整,避免并发过高导致资源耗尽。
注意事项
- 确认SQL Server数据库能够承受并行读取的压力,避免因并发连接数过高导致数据库性能下降。
- 适当调整Glue作业的DPU数量(例如从默认2个提升至4-8个),确保有足够资源支撑并行任务。
- Glue 4.0的GlueContext具备线程安全性,无需在每个线程内重新初始化上下文。
内容的提问来源于stack exchange,提问作者gbeaven
相关产品推荐
相关产品推荐

