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

解决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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 04:43:10