Forum Discussion

Sureshmannem's avatar
Sureshmannem
Frequent Visitor
1 year ago
Solved

Parquet files reading into spark data frame is throwing data type error

Dear All,   I have a requirement to read parquet files form the directory into a data frame to prepare the data form Bronze Lakehouse to Silver Lakehouse. while reading files, it is throwing error ...
  • Sureshmannem's avatar
    1 year ago

    Dear Community,

    Thank you for your continued support.

    I’m happy to share that I’ve resolved the issue I was facing, and I’d like to outline the approach I followed in case it helps others encountering similar challenges.

     

    Initial Observation

    The issue occurred when I attempted to load over 50 Parquet files into a single PySpark DataFrame using a wildcard path. PySpark inferred the schema from the data in each file, but inconsistencies arose—some files interpreted a particular attribute as an integer, while others treated the same attribute as a string.

    This led to data type mismatch errors during the read operation.

     

    Testing

    To investigate further, I loaded each file individually into a DataFrame. This worked as expected, confirming that the wildcard-based bulk load was failing due to schema inference conflicts across files.

     

    Solution

    I modified my script to iterate through each file individually, applying the full processing logic per file. This approach bypasses the schema inference conflict and successfully loads and processes all files.

     

    # Step 1: Import required packages
    from pyspark.sql import SparkSession
    from pyspark.sql.types import *
    import pandas as pd
    from functools import reduce
    from notebookutils import mssparkutils

     

    # Step 2: Define the Lakehouse path
    lakehouse_path = "abfss://[email protected]/xxx.Lakehouse/xxx/"

     

    # Step 3: List all Parquet files in the folder
    file_list = mssparkutils.fs.ls(lakehouse_path)
    parquet_files = [f.path for f in file_list if f.path.endswith(".parquet")]

     

    # Step 4: Read schema reference file once
    schema_df = spark.read.parquet("abfss://[email protected]/xxx.Lakehouse/Files/xxx/xxx.parquet").toPandas()
    schema_df = schema_df.head(0)  # Empty schema frame
    schema_columns = schema_df.columns.tolist()

     

    # Step 5: Define helper functions
    def clean_column_name(col_name):
        for sep in ['@', ':']:
            if sep in col_name:
                col_name = col_name.split(sep)[-1]
        return col_name

     

    def rename_columns(df, old_names, new_names):
        return reduce(
            lambda data, idx: data.withColumnRenamed(old_names[idx], new_names[idx]),
            range(len(old_names)),
            df
        )

     

    # Step 6: Loop through each file and process
    for file_path in parquet_files:
        print(f"Processing file: {file_path}")
       
        # Load the file into a Spark DataFrame
        source_df = spark.read.parquet(file_path)
       
        # Clean and rename columns
        old_columns = source_df.columns
        new_columns = [clean_column_name(col) for col in old_columns]
        source_df = rename_columns(source_df, old_columns, new_columns)
       
        # Convert to Pandas
        source_df = source_df.toPandas().reset_index(drop=True)
       
        # Add missing columns
        for col in schema_columns:
            if col not in source_df.columns:
                source_df[col] = pd.NA
       
        # Reorder columns
        source_df = source_df[schema_columns]
       
        # Concatenate with empty schema and convert to string
        final_df = pd.concat([schema_df, source_df], ignore_index=True, sort=False).astype(str)
       
        # Convert back to Spark DataFrame
        final_spark_df = spark.createDataFrame(final_df)
       
        # Show preview (or write to staging)
        final_spark_df.show()
       
  • Sureshmannem's avatar
    Sureshmannem
    1 year ago

    Dear Community,

    Thank you for your continued support.

    I’m happy to share that I’ve resolved the issue I was facing, and I’d like to outline the approach I followed in case it helps others encountering similar challenges.

    Initial Observation

    The issue occurred when I attempted to load over 50 Parquet files into a single PySpark DataFrame using a wildcard path. PySpark inferred the schema from the data in each file, but inconsistencies arose—some files interpreted a particular attribute as an integer, while others treated the same attribute as a string.

    This led to data type mismatch errors during the read operation.

    Testing

    To investigate further, I loaded each file individually into a DataFrame. This worked as expected, confirming that the wildcard-based bulk load was failing due to schema inference conflicts across files.

    Solution

    I modified my script to iterate through each file individually, applying the full processing logic per file. This approach bypasses the schema inference conflict and successfully loads and processes all files.