如何在AWS Glue中从单源表处理事实表与维度表
AWS Glue 拆分JSON数据源到维度表&事实表解决方案
1. 先通过Glue爬虫构建数据源元数据
- 创建Glue爬虫,指定S3上JSON文件的存储路径,配置爬虫识别JSON格式(如果是每行一个JSON对象,选择
JSON格式;嵌套JSON按需调整) - 运行爬虫后,Glue Data Catalog会生成对应的源表,后续ETL作业可直接引用该表,无需手动定义Schema
2. 编写PySpark风格的Glue ETL作业
通过单份脚本一次性完成维度表生成、事实表关联写入,避免多任务拆分的繁琐:
初始化Glue上下文
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)
读取S3 JSON源数据
# 从Data Catalog读取爬取好的JSON源表 source_dyf = glueContext.create_dynamic_frame.from_catalog( database="你的Glue数据库名称", table_name="你的JSON源表名称" ) # 转换为Spark DataFrame方便操作 source_df = source_dyf.toDF()
处理维度表(去重+生成主键)
# 提取维度字段:替换成你要导入维度表的字段(排除amount、balance) dim_fields = ["account_id", "account_name", "user_region", "account_type"] dim_df = source_df.select(*dim_fields).distinct() # 生成维度表主键(优先用业务唯一键;无业务键时用Spark自增ID) dim_df = dim_df.withColumn("dim_account_key", dim_df["account_id"]) # 无业务键时替换为:dim_df = dim_df.withColumn("dim_account_key", monotonically_increasing_id())
写入维度表到目标存储
# 转换回Glue DynamicFrame并写入(以S3 Parquet为例,也可写入Redshift等数仓) dim_dyf = DynamicFrame.fromDF(dim_df, glueContext, "dim_account_dyf") glueContext.write_dynamic_frame.from_options( frame=dim_dyf, connection_type="s3", connection_options={"path": "s3://你的维度表存储桶路径/dim_account/"}, format="parquet", format_options={"compression": "snappy"} )
处理事实表(关联维度键+提取度量字段)
# 关联维度表获取维度引用键 fact_df = source_df.join( dim_df, on=dim_fields, # 用维度字段做关联条件 how="left" ) # 提取事实表字段:维度键 + 金额、余额度量字段 fact_df = fact_df.select("dim_account_key", "amount", "balance")
写入事实表到目标存储
fact_dyf = DynamicFrame.fromDF(fact_df, glueContext, "fact_account_balance_dyf") glueContext.write_dynamic_frame.from_options( frame=fact_dyf, connection_type="s3", connection_options={"path": "s3://你的事实表存储桶路径/fact_account_balance/"}, format="parquet", format_options={"compression": "snappy"} )
提交作业
job.commit()
3. 验证与调度
- 运行作业后,分别查看维度表和事实表的存储路径,确认数据字段、关联关系正确
- 可在Glue控制台设置作业定时触发器,实现增量或全量同步
内容的提问来源于stack exchange,提问作者saravanan subramanian
相关产品推荐
相关产品推荐

