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

如何用AWS Glue/PySpark循环合并同名S3表并处理Schema差异

AWS Glue/PySpark 同名CSV表合并方案(兼容Schema差异)

核心思路

  1. 遍历两个源S3桶,提取所有表名(假设每个表对应独立前缀,如s3://source-bucket-1/table_a/下的所有CSV属于表table_a)
  2. 筛选出两个桶中同名的表,仅处理这类表
  3. 对每个同名表,分别读取两个源的CSV数据,自动对齐Schema(补全缺失字段、统一字段类型)
  4. 合并数据后写入目标S3桶的对应前缀

完整代码实现

import boto3
from pyspark.sql import SparkSession
from pyspark.sql.functions import lit
from pyspark.sql.types import StringType

# 初始化Spark会话(Glue环境中可直接使用GlueContext)
spark = SparkSession.builder.appName("S3TableMerge").getOrCreate()

# 配置参数
SOURCE_BUCKET_1 = "your-source-bucket-1"
SOURCE_BUCKET_2 = "your-source-bucket-2"
TARGET_BUCKET = "your-target-bucket"
S3_PREFIX = ""  # 若表在桶的子目录下,填写前缀,如"tables/"

# 1. 获取两个桶中的所有表名(提取S3前缀作为表名)
def get_table_names(bucket, prefix):
    s3 = boto3.client("s3")
    paginator = s3.get_paginator("list_objects_v2")
    response_iterator = paginator.paginate(Bucket=bucket, Prefix=prefix, Delimiter="/")
    
    table_names = set()
    for resp in response_iterator:
        for prefix_obj in resp.get("CommonPrefixes", []):
            # 提取前缀最后一个斜杠前的部分作为表名
            table_name = prefix_obj["Prefix"].rstrip("/").split("/")[-1]
            if table_name:
                table_names.add(table_name)
    return table_names

# 获取两个桶的表名集合
bucket1_tables = get_table_names(SOURCE_BUCKET_1, S3_PREFIX)
bucket2_tables = get_table_names(SOURCE_BUCKET_2, S3_PREFIX)
# 取交集,只处理同名表
common_tables = bucket1_tables.intersection(bucket2_tables)

# 2. 循环处理每个同名表
for table_name in common_tables:
    try:
        # 构建源路径
        path1 = f"s3://{SOURCE_BUCKET_1}/{S3_PREFIX}{table_name}/"
        path2 = f"s3://{SOURCE_BUCKET_2}/{S3_PREFIX}{table_name}/"
        
        # 读取CSV,inferSchema自动推断类型,header=True表示首行是字段名
        df1 = spark.read.csv(path1, header=True, inferSchema=True, multiLine=True)
        df2 = spark.read.csv(path2, header=True, inferSchema=True, multiLine=True)
        
        # 3. 对齐Schema:获取两个DataFrame的所有字段,补全缺失字段为Null
        all_columns = set(df1.columns).union(set(df2.columns))
        
        # 给df1补全缺失字段
        for col in all_columns - set(df1.columns):
            df1 = df1.withColumn(col, lit(None).cast(StringType()))  # 默认转字符串,可根据需求调整类型
        
        # 给df2补全缺失字段
        for col in all_columns - set(df2.columns):
            df2 = df2.withColumn(col, lit(None).cast(StringType()))
        
        # 统一字段顺序(按all_columns排序,避免union时字段顺序不一致)
        sorted_columns = sorted(all_columns)
        df1 = df1.select(sorted_columns)
        df2 = df2.select(sorted_columns)
        
        # 4. 合并数据(unionByName自动匹配字段名,即使顺序不同)
        merged_df = df1.unionByName(df2)
        
        # 5. 写入目标桶,按表名存储,覆盖模式可根据需求调整(append/overwrite)
        target_path = f"s3://{TARGET_BUCKET}/{S3_PREFIX}{table_name}/"
        merged_df.write.csv(
            target_path,
            header=True,
            mode="overwrite",
            compression="snappy"  # 可选,启用压缩减少存储成本
        )
        
        print(f"Successfully merged table: {table_name}")
    
    except Exception as e:
        print(f"Failed to merge table {table_name}: {str(e)}")
        # 可添加日志记录到CloudWatch或S3,方便排查问题
        continue

# 关闭Spark会话
spark.stop()

关键细节说明

  • Schema对齐逻辑:默认将缺失字段转为StringType,如果需要更精准的类型匹配,可以先提取两个表的Schema,对比后选择更兼容的类型(比如一个是Int,一个是Long,统一用Long)
  • 异常处理:单个表处理失败不会中断整个任务,只会跳过并打印错误信息,适合批量处理数千张表的场景
  • 性能优化:
    • 启用multiLine=True处理包含换行符的CSV字段
    • 使用compression="snappy"压缩输出,降低S3存储成本和读写时间
    • Glue环境中可使用GlueContext替代SparkSession,利用Glue的Catalog和优化器提升性能
  • 存储结构适配:代码假设每个表对应S3中的一个独立前缀(文件夹),如果你的CSV文件是直接以表名命名的文件(如table_a.csv),需要修改get_table_names函数,提取文件名(去掉.csv后缀)作为表名

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 23:30:57