Forum Discussion
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 message
org.apache.spark.SparkException: Parquet column cannot be converted in file
filepath/SRV0001148_20250819065539974.parquet. Column: [Syncxxx.xxx:ApplicationArea.xxx:CreationDateTime], Expected: string, Found: INT96.
#1) sample script:
#2) sample script:
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 packagesfrom pyspark.sql import SparkSessionfrom pyspark.sql.types import *import pandas as pdfrom functools import reducefrom notebookutils import mssparkutils# Step 2: Define the Lakehouse pathlakehouse_path = "abfss://[email protected]/xxx.Lakehouse/xxx/"# Step 3: List all Parquet files in the folderfile_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 onceschema_df = spark.read.parquet("abfss://[email protected]/xxx.Lakehouse/Files/xxx/xxx.parquet").toPandas()schema_df = schema_df.head(0) # Empty schema frameschema_columns = schema_df.columns.tolist()# Step 5: Define helper functionsdef clean_column_name(col_name):for sep in ['@', ':']:if sep in col_name:col_name = col_name.split(sep)[-1]return col_namedef 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 processfor file_path in parquet_files:print(f"Processing file: {file_path}")# Load the file into a Spark DataFramesource_df = spark.read.parquet(file_path)# Clean and rename columnsold_columns = source_df.columnsnew_columns = [clean_column_name(col) for col in old_columns]source_df = rename_columns(source_df, old_columns, new_columns)# Convert to Pandassource_df = source_df.toPandas().reset_index(drop=True)# Add missing columnsfor col in schema_columns:if col not in source_df.columns:source_df[col] = pd.NA# Reorder columnssource_df = source_df[schema_columns]# Concatenate with empty schema and convert to stringfinal_df = pd.concat([schema_df, source_df], ignore_index=True, sort=False).astype(str)# Convert back to Spark DataFramefinal_spark_df = spark.createDataFrame(final_df)# Show preview (or write to staging)final_spark_df.show()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.
4 Replies
- v-ssriganeshCommunity Support
Hello Sureshmannem,
Thank you for reaching out to the Microsoft Fabric Community Forum.I have reproduced your scenario in a Fabric Notebook, and I got the expected results. Below I’ll share the steps, the code I used and screenshots of the outputs for clarity.
- Created a DataFrame with sample data
from datetime import datetime from pyspark.sql import Row data = [ Row(ID="1", Name="Ganesh", CreationDateTime=datetime.now().strftime("%Y-%m-%d %H:%M:%S.%f")[:-3]), Row(ID="2", Name="Ravi", CreationDateTime=datetime.now().strftime("%Y-%m-%d %H:%M:%S.%f")[:-3]) ] df = spark.createDataFrame(data) df.printSchema() df.show(truncate=False)Output (Screenshot 1 – Schema & Screenshot 2 – Data):
- Saved DataFrame as a Lakehouse table
df.write.mode("overwrite").saveAsTable("DemoTable")- Verified the table in catalog
spark.catalog.listTables("default")Output (Screenshot 3 – Table Catalog):
With this approach, the table DemoTable was successfully created in the Lakehouse with the expected schema, and data was retrieved correctly with CreationDateTime as a string. It worked in my case because I explicitly formatted the creationdatetime column as a string before saving to the Lakehouse table. By default, Spark can sometimes infer a different data type (like timestamp) depending on how the value is created. Converting it to string ensures consistency and prevents schema mismatch issues.
Best Regards,
Ganesh singamshetty.- SureshmannemFrequent Visitor
Hi Ganesh,
Thanks for your kind support and explanation.
My scneario is slightly different, I am sharing the sample script with masking
I have a scneario to read parquet files stored in lakehouse into a data frame to prepare my data, the issue is happening at very first step itself
source_df = spark.read.parquet("abfss://[email protected]/xxxx.Lakehouse/Files/xxxx/SRV0001148_*.parquet")
error: org.apache.spark.SparkException: Parquet column cannot be converted in file xxxxxxx Expected: string, Found: INT96.
I have tried by defining my schema explicitly, spark is still ignoring and considering only from parquet files.
- SureshmannemFrequent Visitor
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.
- SureshmannemFrequent Visitor
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 packagesfrom pyspark.sql import SparkSessionfrom pyspark.sql.types import *import pandas as pdfrom functools import reducefrom notebookutils import mssparkutils# Step 2: Define the Lakehouse pathlakehouse_path = "abfss://[email protected]/xxx.Lakehouse/xxx/"# Step 3: List all Parquet files in the folderfile_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 onceschema_df = spark.read.parquet("abfss://[email protected]/xxx.Lakehouse/Files/xxx/xxx.parquet").toPandas()schema_df = schema_df.head(0) # Empty schema frameschema_columns = schema_df.columns.tolist()# Step 5: Define helper functionsdef clean_column_name(col_name):for sep in ['@', ':']:if sep in col_name:col_name = col_name.split(sep)[-1]return col_namedef 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 processfor file_path in parquet_files:print(f"Processing file: {file_path}")# Load the file into a Spark DataFramesource_df = spark.read.parquet(file_path)# Clean and rename columnsold_columns = source_df.columnsnew_columns = [clean_column_name(col) for col in old_columns]source_df = rename_columns(source_df, old_columns, new_columns)# Convert to Pandassource_df = source_df.toPandas().reset_index(drop=True)# Add missing columnsfor col in schema_columns:if col not in source_df.columns:source_df[col] = pd.NA# Reorder columnssource_df = source_df[schema_columns]# Concatenate with empty schema and convert to stringfinal_df = pd.concat([schema_df, source_df], ignore_index=True, sort=False).astype(str)# Convert back to Spark DataFramefinal_spark_df = spark.createDataFrame(final_df)# Show preview (or write to staging)final_spark_df.show()