Post

Getting Started with Apache Spark on Databricks: A Practical Guide

In this article, let us get started with Apache Spark on Databricks by building a small data pipeline in a notebook. We will create compute, load data, clean and aggregate it, and finally write the result as a Delta table. This approach is useful when you want to learn Spark without first setting up machines, installing Java, or managing a Spark cluster yourself.

The example is intentionally small, but it covers the same basic flow that most Spark pipelines follow: read, transform, validate, and write.

What are Spark and Databricks?

Apache Spark is a distributed processing engine. Instead of processing a large dataset on one machine, Spark can divide the work across multiple machines. Databricks provides a managed environment around Spark with notebooks, compute, job scheduling, Delta Lake, and other features in one place.

For a beginner, the main benefit is that we can focus on DataFrame code while Databricks handles most of the Spark installation and cluster configuration. It still helps to understand what is running underneath, especially when a pipeline becomes slow or expensive.

OptionBest suited forWhat to keep in mind
Spark DataFrame APIMost transformation pipelinesEasy to test and compose in Python or Scala
Spark SQLAnalysts and SQL-heavy transformationsReadable for joins and aggregations
RDD APILow-level or unusual processingMore control, but usually more code and fewer optimisations

For our use case, we will use PySpark DataFrames and a small amount of Spark SQL.

1. Create Databricks compute

From the Databricks workspace, open Compute and create an all-purpose compute resource. Choose a recent Databricks Runtime LTS version. The runtime already includes Spark and commonly used Python libraries.

For a simple demo, a single-node compute configuration is enough if that option is available in your workspace. There is no benefit in starting a large cluster for a file with a few rows. Also configure automatic termination, perhaps after 15 or 30 minutes, so the compute does not continue running after we close the notebook.

In a production pipeline, I would use job compute rather than leave an interactive cluster running. Job compute is created for a run and terminated afterwards, which gives better isolation and usually makes the cost easier to understand.

2. Create a notebook and inspect Spark

Create a Python notebook and attach it to the compute. Databricks creates a SparkSession named spark, so we do not need to create one manually. Run the below code to confirm it is available.

1
spark.version

We can also check a couple of settings.

1
2
print(spark.conf.get("spark.databricks.clusterUsageTags.sparkVersion"))
print(spark.conf.get("spark.sql.session.timeZone"))

The session time zone is worth noticing. Date and timestamp problems are common when the source system, Spark session, and destination use different time zones. I prefer to make the expected time zone explicit in a real pipeline.

3. Create some source data

Let us use a small sales dataset so the example can be run without downloading anything.

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
from pyspark.sql import functions as F
from pyspark.sql.types import (
    StructType, StructField, StringType,
    IntegerType, DoubleType
)

schema = StructType([
    StructField("order_id", StringType(), False),
    StructField("order_date", StringType(), True),
    StructField("country", StringType(), True),
    StructField("quantity", IntegerType(), True),
    StructField("unit_price", DoubleType(), True)
])

rows = [
    ("O-1001", "2026-09-20", "AU", 2, 35.50),
    ("O-1002", "2026-09-20", "NZ", 1, 20.00),
    ("O-1003", "2026-09-21", "AU", 3, 12.00),
    ("O-1004", "bad-date", "AU", 1, 99.00)
]

raw_df = spark.createDataFrame(rows, schema)
display(raw_df)

When reading a real CSV, define the schema instead of depending on inferSchema. Inference is convenient for exploration, but it adds an extra read and can produce different types when the input changes. A typical read would look like this:

1
2
3
4
raw_df = (spark.read
    .option("header", True)
    .schema(schema)
    .csv("/Volumes/main/landing/sales/orders/*.csv"))

Volumes are the preferred way to access files governed by Unity Catalog. Older examples often use DBFS paths, so check which storage pattern your workspace supports before copying them.

4. Clean and transform the data

Spark transformations are lazy. The following code describes the work, but Spark does not execute it until we call an action such as show, count, or write.

1
2
3
4
5
6
7
8
9
clean_df = (raw_df
    .withColumn("order_date", F.to_date("order_date"))
    .withColumn("amount", F.col("quantity") * F.col("unit_price"))
    .filter(F.col("order_date").isNotNull())
    .filter((F.col("quantity") > 0) & (F.col("unit_price") >= 0))
    .select("order_id", "order_date", "country", "quantity", "unit_price", "amount")
)

clean_df.show()

The row containing bad-date becomes null after to_date and is removed. Silently dropping a record is acceptable for this demo, but not ideal for production. I would normally write rejected rows to a quarantine table with a reason, then monitor the number of rejected records. Otherwise, a source format change could remove a large part of the data without anybody noticing.

Let us aggregate the valid records by date and country.

1
2
3
4
5
6
7
8
9
daily_sales_df = (clean_df
    .groupBy("order_date", "country")
    .agg(
        F.countDistinct("order_id").alias("order_count"),
        F.round(F.sum("amount"), 2).alias("sales_amount")
    )
)

display(daily_sales_df.orderBy("order_date", "country"))

We could also create a temporary view and query the same DataFrame using SQL.

1
daily_sales_df.createOrReplaceTempView("daily_sales")
1
2
3
SELECT order_date, country, order_count, sales_amount
FROM daily_sales
ORDER BY order_date, country;

A temporary view only exists for the Spark session. It is useful inside a notebook, but it is not a persisted table.

5. Write the result as a Delta table

Delta is the default table format on Databricks and adds a transaction log on top of data files. It supports reliable writes, schema checks, updates, and history.

1
2
3
4
(daily_sales_df.write
    .format("delta")
    .mode("overwrite")
    .saveAsTable("main.demo.daily_sales"))

Replace main.demo with a catalog and schema that exist in your workspace. You need permission to use the catalog and create the table. We can confirm the write using SQL.

1
2
SELECT * FROM main.demo.daily_sales;
DESCRIBE HISTORY main.demo.daily_sales;

Be careful with mode("overwrite"). It is simple for a repeatable demo, but a production pipeline should decide how reruns are handled. Depending on the source, we might append immutable records, overwrite only affected partitions, or use a Delta MERGE keyed by order_id. Running a full overwrite on a large table can be both slow and dangerous.

6. Check the Spark execution plan

Spark decides how to execute our transformations. We can inspect its plan before trying random performance settings.

1
daily_sales_df.explain("formatted")

For this small dataset, performance does not matter. With larger data, look for expensive shuffles, uneven partitions, and large joins. Avoid calling collect() on a large DataFrame because it sends all records to the driver and can crash it. Use display, show, or a limited result while exploring.

Caching is another feature that is easy to overuse. Cache a DataFrame only when it is expensive to calculate and will be reused several times. Cached data takes cluster memory, and a one-time pipeline often gains nothing from it.

What would change for production?

A notebook is a good place to learn, but I would make a few changes before treating this as a production pipeline:

  1. Store the code in Git and move reusable transformations into Python modules.
  2. Use a Databricks job or workflow with job compute, retries, notifications, and a service principal.
  3. Read from controlled Unity Catalog locations and grant only the permissions needed.
  4. Add checks for duplicate order IDs, nulls, invalid dates, and unexpected volume changes.
  5. Keep rejected records instead of only filtering them out.
  6. Use an incremental write strategy and make reruns idempotent.
  7. Pass catalog names, paths, and processing dates as parameters rather than hard-coding them.

One more practical point is cost. Spark can process very large data, but that does not mean every dataset needs a large cluster. Start with modest compute, review the Spark UI and run duration, and scale only when the evidence shows it is needed.

Conclusion

We created a small Spark pipeline on Databricks, starting from a DataFrame and ending with a Delta table. The example also showed the main habits worth carrying into a real project: define schemas, validate bad records, understand lazy execution, inspect the plan, and be deliberate about how data is written. Once this flow is comfortable, the next step is to place the notebook in a scheduled Databricks job and replace the sample rows with a real source.

This post is licensed under CC BY 4.0 by the author.