Forum Discussion
Parquet files reading into spark data frame is throwing data type error
- 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 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() - 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.
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.
- Sureshmannem1 year agoFrequent 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.
- Sureshmannem1 year agoFrequent 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.