如何用AWS Glue/PySpark循环合并同名S3表并处理Schema差异
AWS Glue/PySpark 同名CSV表合并方案(兼容Schema差异)
核心思路
- 遍历两个源S3桶,提取所有表名(假设每个表对应独立前缀,如
s3://source-bucket-1/table_a/下的所有CSV属于表table_a) - 筛选出两个桶中同名的表,仅处理这类表
- 对每个同名表,分别读取两个源的CSV数据,自动对齐Schema(补全缺失字段、统一字段类型)
- 合并数据后写入目标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
相关产品推荐
相关产品推荐

