如何在AWS Athena、Python中实现列转行并集成到AWS Glue Job?
实现列转行(拆分斜杠分隔值):Athena、Python(PySpark)及Glue Job集成方案
一、AWS Athena SQL实现
直接通过SQL拆分并展开多值列,确保拆分后字段一一对应:
SELECT t.id AS ID, t.name AS Name, CAST(t.sector AS BIGINT) AS Sector -- 按需转换数据类型 FROM your_source_table CROSS JOIN UNNEST( split(ID, '/'), split(Name, '/'), split(Sector, '/') ) AS t(id, name, sector);
逻辑说明:
split函数将斜杠分隔的字符串拆分为数组;UNNEST同时展开三个数组,要求三个数组长度一致,以此保证每行ID、Name、Sector的对应关系;- 若源表有多行数据,该逻辑可批量处理所有记录。
二、Python(PySpark)实现(适配Glue Job)
用PySpark API处理拆分逻辑,代码示例:
from pyspark.sql.functions import split, arrays_zip, explode, col # 读取源数据(Glue中可直接对接Data Catalog或S3) df = spark.read.csv("s3://your-source-bucket/path/", header=True, inferSchema=True) # 1. 将各列拆分为数组 df_with_arrays = df.withColumn("id_array", split(col("ID"), "/")) \ .withColumn("name_array", split(col("Name"), "/")) \ .withColumn("sector_array", split(col("Sector"), "/")) # 2. 打包数组为结构体,再展开为单行单值 df_exploded = df_with_arrays.withColumn("zipped", arrays_zip("id_array", "name_array", "sector_array")) \ .select(explode("zipped").alias("zipped_cols")) \ .select( col("zipped_cols.id_array").alias("ID"), col("zipped_cols.name_array").alias("Name"), col("zipped_cols.sector_array").cast("bigint").alias("Sector") ) # 写入目标存储(示例为Parquet格式) df_exploded.write.parquet("s3://your-target-bucket/path/", mode="overwrite")
三、集成到AWS Glue Job
- 在Glue控制台创建Spark类型的ETL Job,语言选择Python;
- 替换默认脚本为适配Glue的逻辑,关键代码示例:
from awsglue.context import GlueContext from awsglue.job import Job from pyspark.context import SparkContext from awsglue.dynamicframe import DynamicFrame from pyspark.sql.functions import split, arrays_zip, explode, col # 初始化Glue上下文 sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) job.init(args['JOB_NAME'], args) # 从Glue Data Catalog读取源表 source_dyf = glueContext.create_dynamic_frame.from_catalog( database="your_database", table_name="your_source_table" ) df = source_dyf.toDF() # 执行拆分逻辑(同上述PySpark代码) df_with_arrays = df.withColumn("id_array", split(col("ID"), "/")) \ .withColumn("name_array", split(col("Name"), "/")) \ .withColumn("sector_array", split(col("Sector"), "/")) df_exploded = df_with_arrays.withColumn("zipped", arrays_zip("id_array", "name_array", "sector_array")) \ .select(explode("zipped").alias("zipped_cols")) \ .select( col("zipped_cols.id_array").alias("ID"), col("zipped_cols.name_array").alias("Name"), col("zipped_cols.sector_array").cast("bigint").alias("Sector") ) # 转换为DynamicFrame并写入目标表 target_dyf = DynamicFrame.fromDF(df_exploded, glueContext, "target_dyf") glueContext.write_dynamic_frame.from_catalog( frame=target_dyf, database="your_database", table_name="your_target_table", transformation_ctx="write_ctx" ) job.commit()
- 配置Job的IAM角色,确保拥有S3访问、Glue Data Catalog读写权限;
- 启动Job完成ETL处理。
内容的提问来源于stack exchange,提问作者Yoga
相关产品推荐
相关产品推荐

