Anyone who has worked with insurance data knows the feeling. A partner, a broker or a third-party administrator sends a CSV or Excel file, and it's almost what you expected. The columns are in a different order this month. A date is in American format. Someone has typed "1,250.00" into a premium column, or "N/A", or left a total row at the bottom. A policy description contains a line break inside quotes.
Spark's file readers are fast, but their defaults assume clean data. Read a messy file with schema inference and you'll either get a job failure or, worse, quietly lose values. This post covers the techniques I use in Azure Databricks to read awkward files safely, keep every bad row for investigation, and only let clean data into Delta tables. The examples use Databricks Runtime 7.x (Spark 3.0).
1. Read everything as text first
Schema inference looks at the data and guesses types. With third-party files that's a gamble: one odd value changes a numeric column to a string, or a date column comes through as text one month and as a timestamp the next. I always read as strings and cast later, where I control what happens to bad values.
1: raw = (spark.read 2: .option("header", True) 3: .option("inferSchema", False) # every column arrives as STRING 4: .csv("/mnt/landing/premium/2021/02/premium_extract.csv")) 5: 6: raw.printSchema()
2. Handle quotes, line breaks and encoding
Three options fix most "the row count is wrong" problems:
1: raw = (spark.read 2: .option("header", True) 3: .option("inferSchema", False) 4: .option("multiLine", True) # quoted fields that contain line breaks 5: .option("quote", '"') 6: .option("escape", '"') # "" inside a quoted field, Excel-style 7: .option("encoding", "UTF-8") # try "windows-1252" for files saved from older Excel 8: .csv(path))
If names like "Müller" or "Société" come through garbled, the file usually isn't UTF-8. Files saved from Excel on Windows are often windows-1252. Note that multiLine stops Spark splitting a single file across tasks, so only turn it on where you need it.
3. Keep malformed rows instead of losing them
When a row has the wrong number of columns, Spark's behaviour depends on the mode option. The default, PERMISSIVE, fills what it can and sets the rest to null, which hides the problem. Add a corrupt-record column to the schema and Spark puts the whole original line there instead:
1: from pyspark.sql.types import StructType, StructField, StringType 2: 3: cols = ["policy_ref", "insured_name", "inception_date", "expiry_date", 4: "currency", "gross_premium", "commission"] 5: 6: schema = StructType([StructField(c, StringType(), True) for c in cols] + 7: [StructField("_corrupt_record", StringType(), True)]) 8: 9: df = (spark.read 10: .option("header", True) 11: .option("mode", "PERMISSIVE") 12: .option("columnNameOfCorruptRecord", "_corrupt_record") 13: .schema(schema) 14: .csv(path) 15: .cache()) # see the note below 16: 17: bad_lines = df.filter("_corrupt_record IS NOT NULL") 18: good_lines = df.filter("_corrupt_record IS NULL").drop("_corrupt_record")
The cache() isn't optional. Since Spark 2.3, a query that only references the corrupt-record column is disallowed on the raw file, so cache (or save) the DataFrame first.
On Databricks there's also a simpler option, badRecordsPath, which writes unparseable rows to a folder as JSON with the reason, while the good rows load normally:
1: df = (spark.read 2: .option("header", True) 3: .option("badRecordsPath", "/mnt/quarantine/premium/") 4: .schema(schema_without_corrupt_column) 5: .csv(path))
4. Clean up column names
Headers like "Gross Premium (GBP)" or "Inception Date " (with a trailing space) are awkward in Spark SQL and break Delta writes, because Delta doesn't allow spaces or certain characters in column names. Normalise them straight after reading:
1: import re 2: 3: def clean_name(name: str) -> str: 4: name = name.strip().lower() 5: name = re.sub(r"[^a-z0-9]+", "_", name) # spaces, brackets, slashes -> _ 6: return name.strip("_") 7: 8: df = raw.toDF(*[clean_name(c) for c in raw.columns]) 9: # "Gross Premium (GBP)" -> gross_premium_gbp
5. Excel files
For Excel, the spark-excel library (installed on the cluster from Maven) lets you point at a sheet and a starting cell, which handles the title rows people put above the real header:
1: xl = (spark.read.format("com.crealytics.spark.excel") 2: .option("dataAddress", "'Premium'!A4") # sheet name and the header cell 3: .option("header", True) 4: .option("inferSchema", False) 5: .load("/mnt/landing/premium/2021/02/premium_extract.xlsx"))
For small workbooks, pandas is often simpler, then convert to a Spark DataFrame. Read through the /dbfs path and force every column to string:
1: import pandas as pd 2: 3: pdf = pd.read_excel("/dbfs/mnt/landing/premium/2021/02/premium_extract.xlsx", 4: sheet_name="Premium", skiprows=3, dtype=str) 5: xl = spark.createDataFrame(pdf.fillna(""))
6. Remove blank and total rows
1: from pyspark.sql import functions as F 2: 3: df = df.dropna(how="all") # completely empty rows 4: df = df.filter(~F.lower(F.trim(F.col("policy_ref"))).isin("total", "grand total", "sub total")) 5: df = df.filter(F.trim(F.col("policy_ref")) != "")
7. Cast numbers and dates carefully
Numbers come with thousands separators, currency symbols, and sometimes accounting-style negatives in brackets. Clean the text before casting:
1: def to_amount(col): 2: c = F.trim(col) 3: is_negative = c.startswith("(") & c.endswith(")") 4: c = F.regexp_replace(c, r"[£$€,()\s]", "") 5: amount = c.cast("decimal(18,2)") 6: return F.when(is_negative, -amount).otherwise(amount) 7: 8: df = (df.withColumn("gross_premium_amt", to_amount(F.col("gross_premium"))) 9: .withColumn("commission_amt", to_amount(F.col("commission"))))
Dates need an explicit format. One Spark 3.0 gotcha: the new date parser can raise a "Fail to parse ... in the new parser" error on bad values instead of returning null. Setting the parser policy to CORRECTED gives the behaviour you want for validation, where unparseable dates become null:
1: spark.conf.set("spark.sql.legacy.timeParserPolicy", "CORRECTED") 2: 3: df = (df.withColumn("inception_dt", 4: F.coalesce(F.to_date("inception_date", "dd/MM/yyyy"), 5: F.to_date("inception_date", "yyyy-MM-dd"))) 6: .withColumn("expiry_dt", 7: F.coalesce(F.to_date("expiry_date", "dd/MM/yyyy"), 8: F.to_date("expiry_date", "yyyy-MM-dd"))))
coalesce over a few known formats copes with files that switch format between months. Be careful with ambiguous formats, though. "03/02/2021" is a valid date as both dd/MM and MM/dd, so only allow formats that can't be confused for the same source.
8. Validate, then split good from bad
Keep the original text next to each typed column until validation is done. Then a row is bad if a required value is missing, or if a value was present but couldn't be converted:
1: def failed(raw_col, typed_col): 2: return F.trim(F.col(raw_col)).isNotNull() & (F.trim(F.col(raw_col)) != "") & F.col(typed_col).isNull() 3: 4: checks = F.array( 5: F.when(F.col("policy_ref").isNull(), F.lit("policy_ref missing")), 6: F.when(failed("gross_premium", "gross_premium_amt"), F.lit("gross_premium not numeric")), 7: F.when(failed("inception_date", "inception_dt"), F.lit("inception_date not a valid date")), 8: F.when(F.col("expiry_dt") < F.col("inception_dt"), F.lit("expiry before inception")), 9: ) 10: 11: validated = df.withColumn("errors", 12: F.array_except(checks, F.array(F.lit(None).cast("string")))) 13: 14: clean = validated.filter(F.size("errors") == 0) 15: rejected = validated.filter(F.size("errors") > 0)
The rejected rows, with their original values and a readable reason, are the most useful output of the whole process. Send them back to whoever produced the file. A clear list ("row 214: gross_premium '1.250,00' is not numeric") gets the source fixed. Quietly cleaning it yourself just means you'll be doing it again next month.
9. Let Delta enforce the schema
1: (clean.select("policy_ref", "insured_name", "inception_dt", "expiry_dt", 2: "currency", "gross_premium_amt", "commission_amt") 3: .withColumn("source_file", F.input_file_name()) 4: .withColumn("loaded_at", F.current_timestamp()) 5: .write.format("delta").mode("append") 6: .saveAsTable("silver.premium")) 7: 8: (rejected.write.format("delta").mode("append") 9: .option("mergeSchema", "true") 10: .saveAsTable("quarantine.premium_rejected"))
Delta rejects writes whose columns don't match the table, which is exactly what you want for the clean table: a surprise new column in the file can't silently change it. I only use mergeSchema on the quarantine table, where keeping everything matters more than a fixed shape.
Summary
- Read as strings. Cast later, on your terms.
- Use
multiLine,escapeand the rightencodingbefore blaming the data. - Capture malformed rows with
_corrupt_recordorbadRecordsPath. Never let them disappear. - Normalise column names straight away.
- Clean numbers and dates explicitly, and watch the Spark 3.0 date parser.
- Keep raw and typed values side by side, split good from bad, and send the bad rows back to the source.
None of this is complicated, but together it turns "the file loaded" into "the file loaded and we know exactly what was wrong with it".