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

如何在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

  1. 在Glue控制台创建Spark类型的ETL Job,语言选择Python;
  2. 替换默认脚本为适配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()
  1. 配置Job的IAM角色,确保拥有S3访问、Glue Data Catalog读写权限;
  2. 启动Job完成ETL处理。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 17:57:14