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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 14:05:24