解决AWS Glue+Athena中HIVE_PARTITION_SCHEMA_MISMATCH分区架构不匹配问题
AWS数据管道Athena分区Schema不匹配问题解决方案
问题背景
在AWS中搭建数据管道,数据存储于名为input-bucket的S3存储桶,桶内包含多份压缩文件。通过Glue Job完成数据解压、CSV格式转换后,存入目标存储桶。
Glue Job代码
import boto3 import tempfile from pyspark.sql import SparkSession import zipfile import os import pyspark.sql.functions as f import datetime from pyspark.sql.functions import col spark = ( SparkSession.builder.appName("Unzip_Zipped_Files") .config("spark.driver.extraJavaOptions", "-Duser.timezone=GMT+5:30") .config("spark.executor.extraJavaOptions", "-Duser.timezone=GMT+5:30") .getOrCreate() ) def list_folders_and_files(read_bucket_name, prefix): s3_client = boto3.client('s3') paginator = s3_client.get_paginator('list_objects_v2') page_iterator = paginator.paginate(Bucket=read_bucket_name, Prefix=prefix, Delimiter='/') target_bucket_name = "target-bucket" for page in page_iterator: for prefix in page.get('CommonPrefixes', []): folder_name = prefix.get('Prefix').rstrip('/') print(f"Found folder: {folder_name}") list_folders_and_files(read_bucket_name, folder_name+ '/') for obj in page.get('Contents', []): key = obj['Key'] if key.endswith('.zip'): file_name = key.split('/')[-1] current_folder_name = key.split("/")[0] print(f"unextracted file name is {key}") unzip_and_upload_file(s3_client, read_bucket_name, target_bucket_name, key, current_folder_name) def unzip_and_upload_file(s3_client, read_bucket_name, target_bucket_name, key, new_target_folder_key): with tempfile.TemporaryDirectory() as tmpdir: download_path = os.path.join(tmpdir, os.path.basename(key)) s3_client.download_file(read_bucket_name, key, download_path) with zipfile.ZipFile(download_path, 'r') as zip_ref: zip_ref.extractall(tmpdir) for file in os.listdir(tmpdir): if file != os.path.basename(key): file_name = file.split(".csv")[0] if file_name.lower() != ("static_pool_data_ftd") and file_name != ("mShakti_Usage_Report_FTD.xls"): new_key = new_target_folder_key + "/"+ file print(f"file to be uploaded under this path {new_key}") s3_client.upload_file(os.path.join(tmpdir, file), target_bucket_name, new_key) df = spark.read.option("header", "true").option("inferSchema", "true").csv(f"s3://{target_bucket_name}/{new_target_folder_key}/{file}", sep='|') result_df= df.withColumn("partition_date",f.lit(new_target_folder_key)) if file_name =="Loan_Closure_and_Foreclosure_Report_MTD": clean_Loan_Closure_and_Foreclosure_Report_MTD_data(result_df, target_bucket_name, file_name, new_target_folder_key) elif file_name=="NPA_Recovery_Report_MTD": clean_NPA_recovery_report_rate(result_df, target_bucket_name, file_name, new_target_folder_key) elif file_name == "customer_wise_disbursement_report_mtd": result_df = result_df.withColumn("sales officer code", col("sales officer code").cast("double")) result_df.coalesce(1).write.mode("overwrite").partitionBy("partition_date").option("header", "true").csv(f"s3://{target_bucket_name}/{file_name}/") elif file_name == "collection_recon_report_mtd": result_df = result_df.withColumn("total collection", col("total collection").cast("string")) result_df.coalesce(1).write.mode("overwrite").partitionBy("partition_date").option("header", "true").csv(f"s3://{target_bucket_name}/{file_name}/") else: result_df.coalesce(1).write.mode("overwrite").partitionBy("partition_date").option("header","true").csv(f"s3://{target_bucket_name}/{file_name}/") print(f"s3://{target_bucket_name}/{new_target_folder_key}/{file} is the file read...") s3_client.delete_object(Bucket=target_bucket_name,Key=f"{new_target_folder_key}/{file}") print(f"Objects deleted {new_target_folder_key}/{file}") def clean_NPA_recovery_report_rate(df, target_bucket_name, file_name, new_target_folder_key): try: df = df.withColumn("Outstanding amount as on NPA date", f.regexp_replace(f.col("Outstanding amount as on NPA date"), ",", "")) df.coalesce(1).write.mode("overwrite").option("header","true").csv(f"s3://{target_bucket_name}/{file_name}/{new_target_folder_key}/") except Exception as e: raise e def clean_Loan_Closure_and_Foreclosure_Report_MTD_data(df, target_bucket_name, file_name, new_target_folder_key): try: for column in ["BRANCH CODE","CUSTOMER NUMBER","ACCOUNT NUMBER"]: df = df.withColumn(column, f.regexp_replace(f.col(column), "'", "")) df.coalesce(1).write.mode("overwrite").option("header","true").csv(f"s3://{target_bucket_name}/{file_name}/{new_target_folder_key}/") except Exception as e: raise e if __name__=="__main__": list_folders_and_files("input-bucket","28092024")
问题现象
执行Glue Job后,通过Glue Crawler生成数据目录供Athena查询,首次查询正常。次日以OVERWRITE模式加载数据后,Athena报错HIVE_PARTITION_SCHEMA_MISMATCH,提示某表新旧分区数据类型(BIGINT与DOUBLE)不匹配。尝试转换数据类型无效,需确保Athena仅展示最新数据。
解决方案
1. 彻底清理旧分区数据
OVERWRITE模式可能未完全覆盖旧分区文件,导致新旧数据混合。在Glue Job写入新数据前,先删除目标分区的所有旧文件:
# 在写入逻辑前添加删除代码 s3_client = boto3.client('s3') # 示例:删除指定表下的旧分区目录 s3_client.delete_objects( Bucket="target-bucket", Delete={ "Objects": [{"Key": f"customer_wise_disbursement_report_mtd/partition_date={new_target_folder_key}/"}] } )
2. 修改Glue Crawler配置
关闭Crawler的分区自动推断,避免重新扫描旧分区导致Schema混乱:
- 进入Glue Crawler配置页,取消勾选分区推断
- 使用自定义CSV分类器,指定分隔符为
|、包含表头 - 配置Crawler仅扫描最新写入的分区路径,而非全表目录
3. 手动指定Schema替代inferSchema
Spark的inferSchema易因数据波动导致类型推断不一致,改为手动定义固定Schema:
from pyspark.sql.types import StructType, StructField, StringType, DoubleType # 以customer_wise_disbursement_report_mtd为例定义Schema custom_schema = StructType([ StructField("sales officer code", DoubleType(), True), StructField("其他字段1", StringType(), True), StructField("partition_date", StringType(), True) # 按实际字段补充完整 ]) # 读取数据时使用自定义Schema df = spark.read.option("header", "true").schema(custom_schema).csv(f"s3://{target_bucket_name}/{new_target_folder_key}/{file}", sep='|')
4. 手动同步Glue表Schema
若已出现分区Schema不匹配,手动统一表结构:
- 进入Glue数据目录,找到对应表,选择编辑Schema
- 将冲突字段统一为目标数据类型(如统一为Double)
- 在Athena中执行
MSCK REPAIR TABLE <表名>刷新分区信息
5. 使用Athena CTAS创建仅含最新数据的表
通过CTAS语句直接生成仅包含最新分区的表,彻底隔离旧数据:
CREATE TABLE latest_data_table WITH (format = 'CSV', external_location = 's3://target-bucket/latest-data/') AS SELECT * FROM original_table WHERE partition_date = '最新日期值';
内容的提问来源于stack exchange,提问作者Khushal Ghathalia
相关产品推荐
相关产品推荐

