使用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侧完成操作:
- 创建临时加载表,将
units字段定义为字符串类型,避免COPY时类型不匹配报错:
CREATE TEMP TABLE load_temp ( other_stuff VARCHAR, units VARCHAR );
- 执行COPY命令将S3 JSON数据加载到临时表:
COPY load_temp FROM 's3://<源数据桶路径>/' IAM_ROLE '<你的Redshift读写S3的IAM角色ARN>' FORMAT JSON 'auto';
- 将临时表数据转换类型后插入正式目标表:
INSERT INTO <正式目标表> SELECT other_stuff, units::INTEGER FROM load_temp;
- 清理临时表即可。
方案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
相关产品推荐
相关产品推荐

