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

如何通过配置文件用Glue(Python/PySpark)将多表从数据源同步到S3

Glue多表从RDBMS同步到S3的PySpark实现方案

1 调整多表同步配置文件

原配置为单表结构,你可以修改为多表数组格式方便批量读取,参考配置如下:

{
  "sync_tables": [
    {
      "source_type": "rdbms",
      "source_schema": "DATABASE",
      "source_table": "DATABASE.Table_1",
      "s3_output_path": "s3://你的存储桶路径/table_1/",
      "write_mode": "overwrite"
    },
    {
      "source_type": "rdbms",
      "source_schema": "DATABASE",
      "source_table": "DATABASE.Table_2",
      "s3_output_path": "s3://你的存储桶路径/table_2/",
      "write_mode": "append"
    }
  ]
}

配置文件上传到S3后即可在Glue脚本中读取调用。

2 完整Glue PySpark脚本

import sys
import json
from awsglue.context import GlueContext
from awsglue.job import Job
from pyspark.context import SparkContext

# 初始化Glue运行环境
sc = SparkContext()
glueContext = GlueContext(sc)
spark = glueContext.spark_session
job = Job(glueContext)
job.init(sys.argv[1], sys.argv)

# 自定义配置项,替换为你的实际参数
# 配置文件在S3上的存储路径
CONFIG_PATH = "s3://你的存储桶路径/sync_config.json"
# 数据库连接参数,敏感信息建议存放在AWS Secrets Manager中,不要硬编码
JDBC_URL = "jdbc:mysql://你的数据库地址:3306/"
DB_USER = "数据库用户名"
DB_PASSWORD = "数据库密码"
JDBC_DRIVER = "com.mysql.cj.jdbc.Driver"

# 读取并解析同步配置
config_str = spark.read.text(CONFIG_PATH).collect()[0][0]
sync_config = json.loads(config_str)
table_list = sync_config.get("sync_tables", [])

# 循环同步所有配置表
for table_conf in table_list:
    source_table = table_conf.get("source_table")
    output_path = table_conf.get("s3_output_path")
    write_mode = table_conf.get("write_mode", "overwrite")

    # 读取关系型数据库表数据
    source_df = spark.read.format("jdbc") \
        .option("url", JDBC_URL) \
        .option("dbtable", source_table) \
        .option("user", DB_USER) \
        .option("password", DB_PASSWORD) \
        .option("driver", JDBC_DRIVER) \
        .load()

    # 写入S3,默认存储格式为parquet,可根据需求修改为csv、json等格式
    source_df.write.mode(write_mode).parquet(output_path)

job.commit()

3 注意事项

  • 数据库账号密码等敏感信息禁止硬编码在脚本中,建议通过Glue作业参数、AWS Secrets Manager传递
  • 不同类型的关系型数据库需要替换对应的JDBC驱动和URL格式,脚本中默认以MySQL为例
  • 大表同步可以新增分批读取、分区写入参数优化性能,避免作业OOM
  • 需要增量同步的场景可以在配置中新增增量字段、同步时间范围等参数,读取数据源时加where条件过滤增量数据即可

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 06:06:03