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

使用AWS Glue动态转换S3混合类型JSON字段并加载至AWS Redshift

解决S3异构JSON units字段加载到Redshift的类型转换方案

以下三种方案均适配数十亿级数据量场景,可根据现有技术栈选择:

方案1:AWS Glue ETL动态帧处理(适配已有Glue作业链路的场景)

DynamicFrame原生支持读取同名字段不同数据类型的异构JSON,不会出现类型推断报错,转换逻辑如下:

  • 读取S3 JSON数据为DynamicFrame,无需指定固定Schema
  • 自定义行级转换函数,直接将units字段统一强转为整数
  • 转换完成后直接写入Redshift目标表
    示例代码:
from awsglue.context import GlueContext
from awsglue.transforms import Map
from pyspark.context import SparkContext

sc = SparkContext()
glueContext = GlueContext(sc)

# 读取源JSON数据
source_dyf = glueContext.create_dynamic_frame.from_options(
    connection_type="s3",
    connection_options={"paths": ["s3://<源数据桶路径>/"]},
    format="json"
)

# 定义转换逻辑
def units_transform(record):
    record["units"] = int(record["units"])
    return record

# 应用转换
transformed_dyf = Map.apply(frame=source_dyf, f=units_transform)

# 写入Redshift
glueContext.write_dynamic_frame.from_jdbc_conf(
    frame=transformed_dyf,
    catalog_connection="<Redshift连接名称>",
    connection_options={"dbtable": "<目标表名>", "database": "<Redshift库名>"},
    redshift_tmp_dir="s3://<临时文件桶路径>/"
)

方案2:Redshift COPY命令直接转换(无需额外ETL组件,操作成本最低)

利用Redshift临时表做中间转换,全程在Redshift侧完成操作:

  1. 创建临时加载表,将units字段定义为字符串类型,避免COPY时类型不匹配报错:
CREATE TEMP TABLE load_temp (
    other_stuff VARCHAR,
    units VARCHAR
);
  1. 执行COPY命令将S3 JSON数据加载到临时表:
COPY load_temp FROM 's3://<源数据桶路径>/'
IAM_ROLE '<你的Redshift读写S3的IAM角色ARN>'
FORMAT JSON 'auto';
  1. 将临时表数据转换类型后插入正式目标表:
INSERT INTO <正式目标表>
SELECT other_stuff, units::INTEGER FROM load_temp;
  1. 清理临时表即可。

方案3:Spark ETL处理(适配已有Spark作业链路的场景)

关闭Schema推断统一按字符串读取,再显式转换units字段为整数:

from pyspark.sql import SparkSession
from pyspark.sql.functions import col

spark = SparkSession.builder.appName("units_transform").getOrCreate()

# 读取JSON时关闭类型推断,所有字段按字符串读取
df = spark.read.option("inferSchema", "false").json("s3://<源数据桶路径>/")

# 转换units为整数
transformed_df = df.withColumn("units", col("units").cast("integer"))

# 写入Redshift
transformed_df.write \
    .format("io.github.spark_redshift_community.spark.redshift") \
    .option("url", "jdbc:redshift://<Redshift端点>:5439/<库名>") \
    .option("dbtable", "<目标表名>") \
    .option("user", "<用户名>") \
    .option("password", "<密码>") \
    .option("tempdir", "s3://<临时文件桶路径>/") \
    .save()

注意:以上方案均基于units字段所有值为有效数字、非空的前提,若后续需要兼容异常值,可将转换逻辑替换为TRY_CAST避免作业中断。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 02:45:08