如何使用Glue Job将S3中的JSON文件转换为CSV并存回同桶
Glue Job ETL 实现示例
以下是基于Spark的AWS Glue Job代码示例,可用于完成数据抽取、转换、加载的典型任务:
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 # 获取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) # 1. 读取数据源(示例:从S3读取CSV) source_dyf = glueContext.create_dynamic_frame.from_options( connection_type="s3", connection_options={"paths": ["s3://your-source-bucket/path/to/data/"]}, format="csv", format_options={"withHeader": True, "separator": ","} ) # 2. 数据转换(示例:字段映射、过滤) transformed_dyf = ApplyMapping.apply( frame=source_dyf, mappings=[ ("old_col1", "string", "new_col1", "string"), ("old_col2", "int", "new_col2", "int"), ("old_col3", "date", "new_col3", "date") ] ) # 过滤掉new_col2为空的记录 filtered_dyf = Filter.apply( frame=transformed_dyf, f=lambda x: x["new_col2"] is not None ) # 3. 写入目标存储(示例:写入S3 Parquet格式) glueContext.write_dynamic_frame.from_options( frame=filtered_dyf, connection_type="s3", connection_options={"path": "s3://your-target-bucket/path/to/output/"}, format="parquet", format_options={"compression": "snappy"} ) # 提交Job job.commit()
关键说明
- 使用
getResolvedOptions获取Glue Job内置参数(如JOB_NAME),也可自定义业务参数传入 DynamicFrame是Glue特有的数据结构,兼容Spark DataFrame,同时支持对接Glue Data Catalog元数据- 可按需替换数据源/目标存储类型(如JDBC、Redshift、Glue Catalog表等)
- 转换步骤可灵活扩展,比如添加
ResolveChoice处理字段类型冲突、DropFields移除无用字段等
内容的提问来源于stack exchange,提问作者Leor s
相关产品推荐
相关产品推荐

