Menu

PySpark course · Lesson 2 of 8

PySpark DataFrames, Columns and Schemas

Build PySpark DataFrames with explicit StructType schemas, transform them with select, withColumn and filter, handle nulls, and cast types safely under Spark 4 ANSI mode.

  • Beginner
  • 24 min read
  • Updated Oct 2026
On this page
  1. Sample data
  2. DataFrame API basics
  3. What a DataFrame is
  4. The parts of the API you use every day
  5. Pitfalls
  6. In interviews
  7. Dataset vs DataFrame
  8. What the difference is
  9. How the “type safety” works in PySpark
  10. Pitfalls and interview angle
  11. Defining a schema with StructType
  12. What it is and why it matters
  13. Two ways to write a schema
  14. Common types
  15. Why not infer the schema?
  16. Nullability is a hint, not validation
  17. In interviews
  18. select, withColumn and filter
  19. What they do
  20. SQL strings are also accepted
  21. Pitfalls
  22. In interviews
  23. Column expressions
  24. What a Column is
  25. Building logic with built-in functions
  26. Comparing with NULLs
  27. Pitfalls
  28. In interviews
  29. Handling nulls with na.fill and na.drop
  30. What they do
  31. NULLs and aggregates
  32. Pitfalls
  33. In interviews
  34. Type casting and schema evolution
  35. Casting
  36. Schema evolution: when sources change shape
  37. Pitfalls
  38. In interviews
  39. Practice questions
  40. Key takeaways

A DataFrame is a distributed table: named, typed columns whose rows are split into partitions across a cluster. Almost every PySpark pipeline is a chain of DataFrame operations, and the schema (column names, types and nullability) is the contract every later step depends on. This lesson covers the core API, how to define schemas on purpose, how to write column logic that behaves correctly with NULLs, and how to cast and reconcile types when data changes shape.

Sample data

Every example on this page uses the small orders DataFrame below. It deliberately contains a missing amount, a missing customer, a missing country and a date string that cannot be parsed, because those are the cases that break real pipelines.

from pyspark.sql import SparkSession, functions as F
from pyspark.sql.types import (StructType, StructField, StringType, IntegerType,
                               DoubleType, ArrayType, MapType)

spark = SparkSession.builder.master("local[2]").appName("dataframes").getOrCreate()
spark.sparkContext.setLogLevel("ERROR")

schema = StructType([
    StructField("order_id", IntegerType(), nullable=False),
    StructField("customer", StringType(), nullable=True),
    StructField("country", StringType(), nullable=True),
    StructField("amount", DoubleType(), nullable=True),
    StructField("order_date", StringType(), nullable=True),
])

orders = spark.createDataFrame(
    [(1, "asha", "IN", 120.0, "2026-01-03"),
     (2, "ben", "UK", None, "2026-01-04"),
     (3, None, "UK", 35.5, "2026-01-04"),
     (4, "chen", None, 80.0, "not a date"),
     (5, "dara", "IN", 210.0, "2026-01-06")],
    schema,
)
orders.printSchema()
orders.show()
root
 |-- order_id: integer (nullable = false)
 |-- customer: string (nullable = true)
 |-- country: string (nullable = true)
 |-- amount: double (nullable = true)
 |-- order_date: string (nullable = true)

+--------+--------+-------+------+----------+
|order_id|customer|country|amount|order_date|
+--------+--------+-------+------+----------+
|       1|    asha|     IN| 120.0|2026-01-03|
|       2|     ben|     UK|  NULL|2026-01-04|
|       3|    NULL|     UK|  35.5|2026-01-04|
|       4|    chen|   NULL|  80.0|not a date|
|       5|    dara|     IN| 210.0|2026-01-06|
+--------+--------+-------+------+----------+

DataFrame API basics

What a DataFrame is

A DataFrame is an immutable, lazily evaluated description of a table. “Immutable” means every method returns a new DataFrame; orders.filter(...) never changes orders. “Lazy” means calling a transformation such as filter or select only adds a step to a logical plan. Nothing is read or computed until you call an action such as show, count, collect or write. Spark then hands the whole plan to the Catalyst optimiser, which can reorder, combine and prune steps before running them as tasks on executors. The transformations vs actions lesson covers laziness in depth.

result = orders.filter(F.col("amount") > 50).select("order_id", "amount")
print(type(result).__name__)     # still a DataFrame: nothing has run yet
print(result.count())            # an action: now Spark executes the plan
print(orders.columns)
print(orders.dtypes)
print(orders.rdd.getNumPartitions())
result.explain()
DataFrame
3
['order_id', 'customer', 'country', 'amount', 'order_date']
[('order_id', 'int'), ('customer', 'string'), ('country', 'string'), ('amount', 'double'), ('order_date', 'string')]
2
== Physical Plan ==
*(1) Project [order_id#0, amount#3]
+- *(1) Filter (isnotnull(amount#3) AND (amount#3 > 50.0))
   +- *(1) Scan ExistingRDD[order_id#0,customer#1,country#2,amount#3,order_date#4]

Notice two things in the plan. The optimiser added isnotnull(amount) on its own, because NULL > 50 can never be true. And count() returned 3, not 4: the row with a NULL amount is not “greater than 50”, so it is filtered out.

The parts of the API you use every day

Category Methods Returns
Column work select, selectExpr, withColumn, withColumns, withColumnRenamed, drop DataFrame
Row filtering filter / where, distinct, dropDuplicates, limit DataFrame
Combining join, union, unionByName DataFrame
Aggregating groupBy(...).agg(...), pivot DataFrame
Ordering orderBy / sort, sortWithinPartitions DataFrame
Actions show, count, collect, take, first, toPandas, write... Python values or side effects
Inspection printSchema, schema, columns, dtypes, explain Metadata (no job, except where a schema must be inferred)

Rows come back as Row objects, which behave like named tuples:

r = orders.first()
print(r)
print(r.customer, r["amount"], r.asDict()["country"])
Row(order_id=1, customer='asha', country='IN', amount=120.0, order_date='2026-01-03')
asha 120.0 IN

Pitfalls

  • collect() and toPandas() pull everything to the driver. On a large DataFrame they exhaust driver memory. Use them on aggregated or limited results only.
  • Each action re-runs the plan from the source unless you cache. Calling count() and then write reads the input twice.
  • Reassigning in a loop builds a deep plan. Hundreds of chained withColumn calls make planning slow; use one select or withColumns instead.

In interviews

Expect “what is a DataFrame and how is it different from a pandas DataFrame?”. A strong answer says it is distributed across executors, immutable, lazily evaluated and optimised as a whole plan by Catalyst, while pandas is eager and in-memory on one machine.

Dataset vs DataFrame

What the difference is

In Scala and Java, Spark has a typed API called a Dataset: Dataset[Order] holds objects of a case class, and the compiler checks field names and types. A DataFrame is simply Dataset[Row]: rows of generic, untyped records whose column names are only checked when the query is analysed at run time.

Python has no Dataset API. Python is dynamically typed, so there is nothing for a compiler to check, and in PySpark everything is a DataFrame. When people ask about Datasets in a PySpark interview, they want to see that you know the distinction exists and why it does not apply to Python.

DataFrame Dataset (Scala/Java only)
Element type Generic Row A typed object such as a case class
When column errors appear Analysis time (when the plan is resolved, before running) Compile time for typed lambdas
Optimisation Full Catalyst optimisation Same for column expressions; typed lambdas (map(o => ...)) are opaque to the optimiser and need objects to be deserialised
Available in Python Yes No
Serialisation Tungsten binary rows Encoders convert between objects and Tungsten rows

How the “type safety” works in PySpark

PySpark still catches column mistakes before running a job: when you reference a column that does not exist, the analyser fails the moment you build the plan, not after an hour of processing. What you do not get is checking of the Python values inside UDFs or rdd.map lambdas. Type hints on pandas UDFs help readers but are not enforced by a compiler.

Pitfalls and interview angle

  • Saying “DataFrames are faster than Datasets” is only half right. Column expressions perform the same in both; the Dataset slowdown comes from typed lambdas, which force object deserialisation and hide logic from Catalyst. The same is true of Python UDFs in PySpark.
  • Since Spark 3.4, PySpark can also run as a thin client over Spark Connect, where the DataFrame API is the only API; there is no RDD access in Connect mode. That makes the DataFrame API the long-term interface for Python.

Defining a schema with StructType

What it is and why it matters

A schema is a StructType: an ordered list of StructField(name, dataType, nullable) entries. You can supply it when creating a DataFrame or reading files. Supplying it yourself means Spark does not have to scan the data to guess types, and it means a bad file cannot silently change the shape of your table.

Two ways to write a schema

ddl = "order_id INT NOT NULL, customer STRING, country STRING, amount DOUBLE, order_date STRING"
print(StructType.fromDDL(ddl) == schema)
print(schema.simpleString())
print(schema["amount"])
print(schema.fieldNames())

nested = StructType([
    StructField("id", IntegerType()),
    StructField("address", StructType([StructField("city", StringType()),
                                       StructField("zip", StringType())])),
    StructField("tags", ArrayType(StringType())),
    StructField("attrs", MapType(StringType(), StringType())),
])
print(nested.simpleString())
True
struct<order_id:int,customer:string,country:string,amount:double,order_date:string>
StructField('amount', DoubleType(), True)
['order_id', 'customer', 'country', 'amount', 'order_date']
struct<id:int,address:struct<city:string,zip:string>,tags:array<string>,attrs:map<string,string>>

The DDL string and the StructType are equivalent. DDL is shorter and easy to keep in config; StructType is easier to build programmatically. schema.json() serialises a schema so you can store it alongside a pipeline and reload it with StructType.fromJson. Nested types (StructType, ArrayType, MapType) are covered in the nested data and explode lesson.

Common types

Spark SQL type PySpark class Notes
INT, BIGINT IntegerType, LongType Python int values are inferred as BIGINT
DOUBLE DoubleType Floating point: do not use for money
DECIMAL(p,s) DecimalType(p, s) Exact; use for currency
STRING StringType
DATE, TIMESTAMP DateType, TimestampType TIMESTAMP is stored in UTC and shown in the session time zone; TIMESTAMP_NTZ has no time zone
BOOLEAN BooleanType

Why not infer the schema?

inferSchema makes Spark read the data an extra time and guess. Guesses change when the data changes: one malformed value and a numeric column becomes a string.

import tempfile, os

path = os.path.join(tempfile.mkdtemp(), "orders.csv")
with open(path, "w") as f:
    f.write("order_id,customer,amount\n1,asha,10\nx,ben,abc\n")

print(spark.read.option("header", True).option("inferSchema", True).csv(path).dtypes)
spark.read.schema("order_id INT, customer STRING, amount DOUBLE").option("header", True).csv(path).show()
[('order_id', 'string'), ('customer', 'string'), ('amount', 'string')]
+--------+--------+------+
|order_id|customer|amount|
+--------+--------+------+
|       1|    asha|  10.0|
|    NULL|     ben|  NULL|
+--------+--------+------+

With an explicit schema the types are stable, and the default PERMISSIVE parse mode turns the unparseable values into NULL. That is better than a changed schema, but still hides bad data: the reading and writing data lesson shows how to capture and count corrupt records instead.

Nullability is a hint, not validation

createDataFrame does check nullable=False for local Python data:

try:
    spark.createDataFrame([(None, "x", "IN", 1.0, "2026-01-01")], schema)
except Exception as e:
    print(type(e).__name__, str(e)[:95])
PySparkValueError [FIELD_NOT_NULLABLE_WITH_NAME] field order_id: This field is not nullable, but got None.

File readers do not: when you read CSV, JSON or Parquet with a user schema, Spark treats every column as nullable. Enforce required fields with explicit checks (count rows where the key is NULL and fail the job if it is not zero), or with a table format such as Delta Lake that supports NOT NULL and CHECK constraints.

In interviews

“How do you handle schema changes when reading files?” Mention explicit schemas (StructType or DDL), why inference is fragile, that nullability is not enforced on read, and how you surface malformed rows.

select, withColumn and filter

What they do

These three methods cover most day-to-day transformation code:

  • select chooses, renames and computes columns. The output has exactly the columns you list.
  • withColumn(name, expr) adds a column or replaces one with the same name, keeping all others. withColumns({...}) does several at once.
  • filter(condition) (alias where) keeps rows where the condition is true. Rows where it is false or NULL are dropped.
enriched = (
    orders
    .select("order_id", "customer", "country", "amount")
    .withColumn("amount_with_tax", F.round(F.col("amount") * 1.2, 2))
    .withColumn("customer", F.initcap("customer"))
    .withColumnRenamed("country", "country_code")
    .filter(F.col("amount").isNotNull() & (F.col("country_code") == "IN"))
)
enriched.show()
+--------+--------+------------+------+---------------+
|order_id|customer|country_code|amount|amount_with_tax|
+--------+--------+------------+------+---------------+
|       1|    Asha|          IN| 120.0|          144.0|
|       5|    Dara|          IN| 210.0|          252.0|
+--------+--------+------------+------+---------------+

SQL strings are also accepted

filter/where accept a SQL predicate string, and selectExpr accepts SQL expressions. They compile to the same plan as the Column API.

orders.where("amount > 30 AND country = 'UK'").show()
orders.selectExpr("order_id", "amount * 2 AS doubled").show(2)
orders.withColumns({
    "is_big": F.col("amount") > 100,
    "year": F.substring("order_date", 1, 4),
}).drop("order_date").show()
+--------+--------+-------+------+----------+
|order_id|customer|country|amount|order_date|
+--------+--------+-------+------+----------+
|       3|    NULL|     UK|  35.5|2026-01-04|
+--------+--------+-------+------+----------+

+--------+-------+
|order_id|doubled|
+--------+-------+
|       1|  240.0|
|       2|   NULL|
+--------+-------+
only showing top 2 rows
+--------+--------+-------+------+------+----+
|order_id|customer|country|amount|is_big|year|
+--------+--------+-------+------+------+----+
|       1|    asha|     IN| 120.0|  true|2026|
|       2|     ben|     UK|  NULL|  NULL|2026|
|       3|    NULL|     UK|  35.5| false|2026|
|       4|    chen|   NULL|  80.0| false|not |
|       5|    dara|     IN| 210.0|  true|2026|
+--------+--------+-------+------+------+----+

Order 2 (UK, amount NULL) is not returned by the first query: NULL > 30 is NULL, and filter drops it. Order 4 shows why string slicing is not date parsing: substring happily returned "not ".

Pitfalls

  • Python operators need parentheses. F.col("a") > 1 & F.col("b") < 2 is parsed by Python as F.col("a") > (1 & F.col("b")) < 2. Always wrap each comparison: (F.col("a") > 1) & (F.col("b") < 2). Use &, |, ~, not and, or, not.
  • withColumn in a loop is slow to plan. Each call adds a projection to the plan. For many columns, build a list and call select once, or use withColumns.
  • withColumn with an existing name silently replaces it. That is useful for cleaning, but a typo creates a new column instead of replacing the old one.
  • drop of a column that does not exist is a no-op, not an error, so typos go unnoticed.

In interviews

You may be asked to translate a SQL query to the DataFrame API or the reverse. Show that select, withColumn and filter map to SELECT, computed columns and WHERE, and mention that rows with a NULL condition are dropped.

Column expressions

What a Column is

F.col("amount") * 1.2 > 100 does not compute anything. It builds a Column expression: a small tree that Spark evaluates for each row when the plan runs. Because it is a tree the optimiser can read, it can push the condition down to the file scan, prune unused columns and generate efficient code. That is why built-in functions are preferred over Python UDFs.

print(F.col("amount") * 1.2 > 100)
print(F.expr("amount * 1.2 > 100"))
Column<'>(*(amount, 1.2), 100)'>
Column<'amount * 1.2 > 100'>

Building logic with built-in functions

orders.select(
    F.col("order_id"),
    F.when(F.col("amount") >= 100, "large")
     .when(F.col("amount").isNotNull(), "small")
     .otherwise("unknown").alias("size"),
    F.concat_ws("-", "country", "customer").alias("label"),
    F.coalesce("customer", F.lit("guest")).alias("customer"),
    F.col("country").isin("IN", "UK").alias("known_country"),
    F.upper(F.col("customer")).startswith("A").alias("starts_a"),
).show()
+--------+-------+-------+--------+-------------+--------+
|order_id|   size|  label|customer|known_country|starts_a|
+--------+-------+-------+--------+-------------+--------+
|       1|  large|IN-asha|    asha|         true|    true|
|       2|unknown| UK-ben|     ben|         true|   false|
|       3|  small|     UK|   guest|         true|    NULL|
|       4|  small|   chen|    chen|         NULL|   false|
|       5|  large|IN-dara|    dara|         true|   false|
+--------+-------+-------+--------+-------------+--------+

Read the NULL behaviour row by row:

  • when checks conditions in order and falls through to otherwise. Without otherwise, unmatched rows get NULL.
  • concat_ws skips NULL inputs (order 3 gives UK, order 4 gives chen). Plain concat returns NULL if any input is NULL.
  • isin on a NULL value returns NULL, not false (order 4).
  • Any function applied to NULL usually returns NULL (order 3, starts_a).

F.lit(value) wraps a Python constant as a column. alias names the result; without it Spark generates names like (amount * 1.2).

Comparing with NULLs

Ordinary comparisons follow SQL three-valued logic: NULL != 'asha' is NULL, so the row is dropped. Use eqNullSafe (SQL <=>) when NULL should be treated as a value:

orders.filter(F.col("customer") != "asha").select("order_id", "customer").show()
orders.filter(~F.col("customer").eqNullSafe("asha")).select("order_id", "customer").show()
+--------+--------+
|order_id|customer|
+--------+--------+
|       2|     ben|
|       4|    chen|
|       5|    dara|
+--------+--------+

+--------+--------+
|order_id|customer|
+--------+--------+
|       2|     ben|
|       3|    NULL|
|       4|    chen|
|       5|    dara|
+--------+--------+

Sorting also has explicit NULL placement: desc_nulls_last(), asc_nulls_first() and so on. By default Spark puts NULLs first in ascending order and last in descending order.

orders.orderBy(F.col("amount").desc_nulls_last()).select("order_id", "amount").show()
+--------+------+
|order_id|amount|
+--------+------+
|       5| 210.0|
|       1| 120.0|
|       4|  80.0|
|       3|  35.5|
|       2|  NULL|
+--------+------+

Pitfalls

  • Referring to columns as df.amount or df["amount"] binds them to that specific DataFrame. After a join or self-join this can point at the wrong side. F.col("amount") is resolved against whatever DataFrame it is used on, which is usually what you want.
  • A Python if cannot inspect a column value: if F.col("amount") > 100: raises an error. Use F.when.
  • Column names with spaces or dots need backticks in string expressions: F.col("`order.id`").

In interviews

“Why prefer built-in functions to UDFs?” is the classic follow-up. The answer is that a Column expression is visible to Catalyst and runs in the JVM, while a Python UDF is a black box that needs data shipped to a Python process. See UDFs and safer alternatives.

Handling nulls with na.fill and na.drop

What they do

df.na (also reachable as dropna, fillna and replace on the DataFrame) gives you bulk NULL handling:

  • na.drop(how="any" | "all", thresh=None, subset=None) removes rows. how="any" (default) drops a row if any of the considered columns is NULL; how="all" only if all are; thresh=n keeps rows with at least n non-NULL values and overrides how.
  • na.fill(value, subset=None) replaces NULLs. A single value only fills columns of a matching type; a dict sets a value per column.
  • na.replace(to_replace, value, subset) swaps specific non-NULL values.
print(orders.na.drop().count())                  # any NULL in any column
print(orders.na.drop(subset=["amount"]).count())  # only look at amount
print(orders.na.drop(how="all").count())
print(orders.na.drop(thresh=5).count())           # at least 5 non-NULL values
orders.na.fill({"customer": "guest", "country": "XX", "amount": 0.0}).show()
orders.na.fill(0).show()
orders.na.replace(["UK"], ["GB"], "country").show()
2
4
5
2
+--------+--------+-------+------+----------+
|order_id|customer|country|amount|order_date|
+--------+--------+-------+------+----------+
|       1|    asha|     IN| 120.0|2026-01-03|
|       2|     ben|     UK|   0.0|2026-01-04|
|       3|   guest|     UK|  35.5|2026-01-04|
|       4|    chen|     XX|  80.0|not a date|
|       5|    dara|     IN| 210.0|2026-01-06|
+--------+--------+-------+------+----------+

+--------+--------+-------+------+----------+
|order_id|customer|country|amount|order_date|
+--------+--------+-------+------+----------+
|       1|    asha|     IN| 120.0|2026-01-03|
|       2|     ben|     UK|   0.0|2026-01-04|
|       3|    NULL|     UK|  35.5|2026-01-04|
|       4|    chen|   NULL|  80.0|not a date|
|       5|    dara|     IN| 210.0|2026-01-06|
+--------+--------+-------+------+----------+

+--------+--------+-------+------+----------+
|order_id|customer|country|amount|order_date|
+--------+--------+-------+------+----------+
|       1|    asha|     IN| 120.0|2026-01-03|
|       2|     ben|     GB|  NULL|2026-01-04|
|       3|    NULL|     GB|  35.5|2026-01-04|
|       4|    chen|   NULL|  80.0|not a date|
|       5|    dara|     IN| 210.0|2026-01-06|
+--------+--------+-------+------+----------+

na.fill(0) only touched the numeric amount column; the string columns kept their NULLs because 0 does not match their type.

NULLs and aggregates

Aggregate functions ignore NULLs, which changes results quietly:

orders.select(
    F.count("*").alias("rows"),
    F.count("amount").alias("non_null_amounts"),
    F.avg("amount").alias("avg_amount"),
    F.sum(F.col("amount").isNull().cast("int")).alias("null_amounts"),
).show()
+----+----------------+----------+------------+
|rows|non_null_amounts|avg_amount|null_amounts|
+----+----------------+----------+------------+
|   5|               4|   111.375|           1|
+----+----------------+----------+------------+

The average is 445.5 / 4, not / 5. If you had filled the missing amount with 0 first, the average would drop to 89.1. Filling is a business decision, not a formatting step.

Pitfalls

  • Filling with a sentinel hides problems. Replacing a missing amount with 0 makes it indistinguishable from a real zero. Prefer keeping NULL, or add a flag column such as amount_missing.
  • na.drop() with no subset drops rows for NULLs in columns you do not care about.
  • Joins and NULL keys: NULL keys never match in an equi-join. Decide whether to drop or fill them before joining, not after.
  • Empty strings are not NULL. CSV readers may produce "" or NULL depending on nullValue and emptyValue options; normalise with F.when(F.trim(c) == "", None).otherwise(c).

In interviews

“How do you handle missing values in Spark?” A strong answer covers profiling first (count NULLs per column), then choosing per column between dropping, filling with a meaningful default, or keeping NULL with a flag, and points out that aggregates ignore NULLs.

Type casting and schema evolution

Casting

cast converts a column to another type; you can pass a type object or a SQL type name such as "int" or "decimal(10,2)". In Spark 4, ANSI mode is on by default (spark.sql.ansi.enabled = true), so an invalid cast raises an error instead of silently returning NULL as Spark 3.x did. When you expect bad values, use try_cast (Column method added in Spark 4.0) or a try_ function such as try_to_date, which return NULL on failure.

typed = orders.select(
    "order_id",
    F.col("order_id").cast("string").alias("id_str"),
    F.col("amount").cast("int").alias("amount_int"),
    F.try_to_date("order_date", "yyyy-MM-dd").alias("order_date"),
)
typed.printSchema()
typed.show()
root
 |-- order_id: integer (nullable = false)
 |-- id_str: string (nullable = false)
 |-- amount_int: integer (nullable = true)
 |-- order_date: date (nullable = true)

+--------+------+----------+----------+
|order_id|id_str|amount_int|order_date|
+--------+------+----------+----------+
|       1|     1|       120|2026-01-03|
|       2|     2|      NULL|2026-01-04|
|       3|     3|        35|2026-01-04|
|       4|     4|        80|      NULL|
|       5|     5|       210|2026-01-06|
+--------+------+----------+----------+

Casting 35.5 to int truncated it to 35; it did not round. Now the failure cases:

try:
    orders.select(F.col("order_date").cast("date")).collect()
except Exception as e:
    print(type(e).__name__, str(e).split(" because")[0])
orders.select(F.col("order_date").try_cast("date").alias("d")).show()

try:
    spark.range(1).select(F.lit("12a").cast("int")).collect()
except Exception as e:
    print(type(e).__name__, str(e).split(" because")[0])
spark.range(1).select(F.lit("12a").try_cast("int").alias("x"),
                      F.lit(" 42 ").cast("int").alias("y")).show()
DateTimeException [CAST_INVALID_INPUT] The value 'not a date' of the type "STRING" cannot be cast to "DATE"
+----------+
|         d|
+----------+
|2026-01-03|
|2026-01-04|
|2026-01-04|
|      NULL|
|2026-01-06|
+----------+

NumberFormatException [CAST_INVALID_INPUT] The value '12a' of the type "STRING" cannot be cast to "INT"
+----+---+
|   x|  y|
+----+---+
|NULL| 42|
+----+---+

Strings with surrounding whitespace such as " 42 " are trimmed and cast successfully.

Schema evolution: when sources change shape

Real sources add columns, reorder them and change types. Three tools handle most of it.

1. Union by name, not by position. union matches columns by position. When the second DataFrame has a different column order, Spark tries to line up mismatched columns and, under ANSI mode, fails the cast; with ANSI off it would silently put values in the wrong column. unionByName(allowMissingColumns=True) matches by name and fills missing columns with NULL.

day1 = spark.createDataFrame([(1, "asha", 120.0)], "order_id INT, customer STRING, amount DOUBLE")
day2 = spark.createDataFrame([(2, 35.5, "promo10")], "order_id INT, amount DOUBLE, coupon STRING")
day1.unionByName(day2, allowMissingColumns=True).show()
try:
    day1.union(day2).show()
except Exception as e:
    print(type(e).__name__, str(e).split(" because")[0])
+--------+--------+------+-------+
|order_id|customer|amount| coupon|
+--------+--------+------+-------+
|       1|    asha| 120.0|   NULL|
|       2|    NULL|  35.5|promo10|
+--------+--------+------+-------+

NumberFormatException [CAST_INVALID_INPUT] The value 'asha' of the type "STRING" cannot be cast to "DOUBLE"

2. Merge Parquet schemas on read. Parquet files written at different times can have different columns. By default Spark takes the schema from a sample of files (spark.sql.parquet.mergeSchema is false), so new columns can be missing from the result. mergeSchema reads every file footer and unions the columns. It costs extra work on large folders.

base = tempfile.mkdtemp()
day1.write.parquet(os.path.join(base, "orders", "load=1"))
day2.write.parquet(os.path.join(base, "orders", "load=2"))
spark.read.parquet(os.path.join(base, "orders")).printSchema()
merged = spark.read.option("mergeSchema", True).parquet(os.path.join(base, "orders"))
merged.printSchema()
merged.orderBy("order_id").show()
root
 |-- order_id: integer (nullable = true)
 |-- customer: string (nullable = true)
 |-- amount: double (nullable = true)
 |-- load: integer (nullable = true)

root
 |-- order_id: integer (nullable = true)
 |-- customer: string (nullable = true)
 |-- amount: double (nullable = true)
 |-- coupon: string (nullable = true)
 |-- load: integer (nullable = true)

+--------+--------+------+-------+----+
|order_id|customer|amount| coupon|load|
+--------+--------+------+-------+----+
|       1|    asha| 120.0|   NULL|   1|
|       2|    NULL|  35.5|promo10|   2|
+--------+--------+------+-------+----+

Without mergeSchema, the coupon column was simply missing. (load comes from the folder name; see partition discovery in the reading and writing lesson.)

3. Conform to a target schema. The most robust pattern is to define the schema you want downstream and project every input onto it: cast existing columns, add missing ones as typed NULLs and drop the rest.

target = StructType.fromDDL("order_id BIGINT, customer STRING, amount DECIMAL(10,2), coupon STRING")

def conform(df, target):
    cols = []
    for field in target.fields:
        if field.name in df.columns:
            cols.append(F.col(field.name).cast(field.dataType).alias(field.name))
        else:
            cols.append(F.lit(None).cast(field.dataType).alias(field.name))
    return df.select(*cols)

conform(day1, target).unionByName(conform(day2, target)).show()
conform(day1, target).printSchema()
+--------+--------+------+-------+
|order_id|customer|amount| coupon|
+--------+--------+------+-------+
|       1|    asha|120.00|   NULL|
|       2|    NULL| 35.50|promo10|
+--------+--------+------+-------+

root
 |-- order_id: long (nullable = true)
 |-- customer: string (nullable = true)
 |-- amount: decimal(10,2) (nullable = true)
 |-- coupon: string (nullable = true)

Table formats take this further: Delta Lake rejects writes that do not match the table schema unless you opt in to evolution. See Delta Lake transactions and schemas and schema evolution patterns.

Pitfalls

  • Casting double to int truncates; use F.round first if you mean rounding.
  • DOUBLE for money produces rounding errors in sums; use DECIMAL.
  • Widening is safe, narrowing is not. INT to BIGINT keeps every value; BIGINT to INT can overflow, which now fails under ANSI.
  • mergeSchema cannot reconcile incompatible types. If one file has amount as string and another as double, the read fails; fix it upstream or read the folders separately and conform.

In interviews

Typical prompts: “a new column appeared in the source, what happens to your job?” and “how do you handle type changes?”. Strong answers mention unionByName, mergeSchema and its cost, conforming to a target schema, ANSI casting with try_cast, and table formats that enforce schemas.

Practice questions

Why does df.filter(F.col(“amount”) > 50) not return rows where amount is NULL?

Comparisons with NULL return NULL, not false, and filter keeps only rows where the condition is true. To include them, add | F.col("amount").isNull() explicitly.

Is there a Dataset API in PySpark? How would you explain the difference to a Scala developer?

No. In Scala, a Dataset is a typed collection of objects checked at compile time, and a DataFrame is Dataset[Row]. Python is dynamically typed, so PySpark exposes only DataFrames. Column expressions get the same Catalyst optimisation in both; typed lambdas in Scala, like Python UDFs, are opaque to the optimiser.

What is wrong with df.filter(F.col(“a”) > 1 & F.col(“b”) < 5)?

Python’s & binds more tightly than > and <, so the expression is parsed as F.col("a") > (1 & F.col("b")) < 5, which fails or does the wrong thing. Wrap each comparison in parentheses: (F.col("a") > 1) & (F.col("b") < 5).

You fill missing amounts with 0 and the monthly average drops. Why, and what would you do instead?

avg ignores NULLs, so before filling it divided by the number of known amounts. After filling, the zeros are counted as real values and pull the average down. Keep the NULL (perhaps with a flag column) unless the business rule really says a missing amount means zero.

A job that worked on Spark 3.5 fails on Spark 4 with CAST_INVALID_INPUT. What changed and how do you fix it?

Spark 4 enables ANSI mode by default, so invalid casts raise errors instead of returning NULL. Find the cast that sees bad data and either clean the data or use try_cast / try_to_date / try_to_number where NULL on failure is acceptable, and count the resulting NULLs so the bad data stays visible.

Why is union dangerous when combining daily extracts, and what should you use?

union matches columns by position. If a source reorders or adds columns, values land in the wrong column (or, under ANSI, the cast fails). unionByName(allowMissingColumns=True) matches by name and fills missing columns with NULL; better still, conform every input to a defined target schema first.

Key takeaways

  • A DataFrame is an immutable, lazy plan over distributed data; actions trigger execution and each action re-runs the plan.
  • PySpark has only DataFrames; the typed Dataset API exists in Scala and Java.
  • Define schemas explicitly with StructType or DDL strings; inference is slow and unstable, and nullability is not enforced on file reads.
  • Column expressions are optimisable trees: build logic with when, coalesce and other built-ins, and remember three-valued NULL logic.
  • Handle NULLs per column on purpose; aggregates ignore them and fills change results.
  • Spark 4 casts are strict by default: use try_cast where failure is expected, and reconcile drifting sources with unionByName, mergeSchema or a conform-to-target step.

By Data Career Hub Editorial · Last reviewed Oct 2026 · All examples run on PySpark 4.2.0 in local mode (local[2]) with Java 21. Error messages are from Spark 4.2; Spark 3.x had ANSI mode off by default, so casts there return NULL instead of failing.

Progress is saved in this browser only. No account needed.

Search
Filter by type