使用AWS Glue从MySQL加载40亿条记录大表失败求助
解决AWS Glue迁移MySQL超大表(40亿条)至S3的通信链路故障问题
很多开发者都遇到过类似的AWS Glue迁移MySQL超大表时的通信链路故障问题,结合你已尝试的哈希分区、扩容Worker、书签索引等优化手段,下面补充几个关键的排查和优化方向:
1. 优化MySQL JDBC连接参数,减少链路压力
默认的MySQL JDBC驱动配置不适合超大规模数据拉取,容易引发频繁通信超时:
- 增大
fetchsize:默认值通常很小(比如1000),会导致Glue频繁向MySQL发起数据请求,增加链路负载。建议设置为10000甚至更高,根据数据库性能调整 - 开启游标拉取:添加
useCursorFetch=true,避免驱动一次性拉取全量数据到内存,改为分批从MySQL游标读取 - 延长超时时间:设置
socketTimeout=3600000(1小时)、connectTimeout=30000(30秒),避免长时间数据传输时被断开 - 开启自动重连:添加
autoReconnect=true,临时链路波动时自动恢复连接
2. 调整分区策略,避免数据库过载
你当前设置的hashpartitions=1000可能过大,导致MySQL同时处理数百个并发查询,超过数据库连接数上限或CPU负载,间接引发链路故障:
- 降低哈希分区数,比如调整为
200(与Worker数量匹配,10个Worker每个处理20个分区) - 尝试范围分区替代哈希分区:基于主键
Data103_PK做范围拆分,配合pushDownPredicate分批次读取数据,比如每次读取一个主键区间的记录,降低单批次查询压力
3. 排查网络稳定性
- 确认Glue作业是否与MySQL在同一VPC内,或通过VPC peering、专线连接,避免公网传输的抖动和延迟
- 检查MySQL端的
wait_timeout、interactive_timeout配置,如果连接空闲时间超过阈值会被主动断开,建议调大至7200秒以上,或通过JDBC参数autoReconnect=true规避 - 检查安全组、防火墙规则,确保Glue的IP范围被允许访问MySQL端口
4. 优化Glue资源配置
- 升级Worker类型:使用G.2X或G.4X规格的Worker,提供更大的内存和CPU,避免因内存溢出导致任务中断(间接引发链路异常)
- 调整Spark参数:设置
spark.sql.shuffle.partitions与哈希分区数一致,spark.driver.maxResultSize调至10G以上,避免数据 shuffle 时的内存瓶颈
修改后的示例代码
import sys from awsglue.utils import getResolvedOptions from awsglue.context import GlueContext from awsglue.job import Job from pyspark.context import SparkContext args = getResolvedOptions(sys.argv, ["JOB_NAME","S3_BUCKET_NAME"]) sc = SparkContext() sc._jsc.hadoopConfiguration().set("fs.s3.canned.acl", "BucketOwnerFullControl") # 配置Spark参数优化 sc._jsc.hadoopConfiguration().set("spark.sql.shuffle.partitions", "200") sc._jsc.hadoopConfiguration().set("spark.driver.maxResultSize", "10g") glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) job.init(args["JOB_NAME"], args) bucket_name = args["S3_BUCKET_NAME"] s3_path = "s3://" + bucket_name + "/flow_datad0036103" datasource0 = glueContext.create_dynamic_frame.from_catalog( database="centerstagedb", table_name="repo_flow_datad0036103", additional_options={ "jobBookmarkKeys":["Data103_PK","FLOW_FK"], "jobBookmarkKeysSortOrder":"asc", "hashexpression":"FLOW_FK", "hashpartitions":"200", # 降低分区数 "fetchsize": "10000", # 增大批量拉取数 "socketTimeout": "3600000", # 1小时Socket超时 "connectTimeout": "30000", "useCursorFetch": "true", # 开启游标拉取 "autoReconnect": "true" }, transformation_ctx="datasource0", ) datasink4 = glueContext.write_dynamic_frame.from_options( frame=datasource0, connection_type="s3", connection_options = {"path": s3_path, "partitionKeys": ["FLOW_FK"]}, format="parquet", transformation_ctx="datasink4", ) job.commit()
内容的提问来源于stack exchange,提问作者Bogdan Gheoca
相关产品推荐
相关产品推荐

