Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing

PySpark for Distributed Processing

Progress0/21
x

PySpark Architecture

  • PySpark architecture: the one explanation35m
  • When do you need Spark?12m
  • What is a cluster?12m
  • Reading and writing data12m
  • SQL inside Spark12m

Mental model & DataFrames

  • PySpark architecture & mental model10m
  • DataFrame basics10m
  • Column operations & built-in functions10m

Aggregations, windows & nested data

  • Aggregations & groupings10m
  • PySpark window functions12m
  • Joins & optimization strategies10m
  • Handling complex & nested data12m
  • Partitioning, repartition & coalesce10m
  • Medallion pipeline project14m

Performance at Scale

  • Caching & persistence12m
  • Data skew detection and salting14m
  • Reading a Catalyst physical plan12m

Production Spark

  • UDFs in depth, and why to avoid them12m
  • Structured Streaming & watermarks14m
  • Delta Lake: MERGE, time travel, ACID12m
  • Capstone part 2: incremental MERGE16m
Back to track
  1. Learn
  2. PySpark for Distributed Processing
  3. PySpark Architecture
  4. Reading and writing data

Lesson 4 of 21 · Theory first, then run it

Reading and writing data

pysparkbeginner12 min

Overview

On a cluster you spark.read.parquet and write Parquet. In this tab the catalog is spark.table.

On this page7 sections›
  1. 1What you will do
  2. 2Why this skill
  3. 3How the code works
  4. 4Worked examples
  5. 5Common beginner questions
  6. 6What comes next
  7. 7Practice

What you will do

Spark jobs read files from storage and write files back. The usual production format is Parquet: columnar, compressed, with types stored in the file. CSV is text with separators. JSON is text with nested objects. Each format has a reader: spark.read.parquet, spark.read.csv, spark.read.json.

When tables are registered in a catalog, you do not pass a path every time. spark.table("orders") means "open the orders table the catalog knows." In this tab, that is how every exercise loads data. The files are already registered. There is no cluster disk for you to write to.

You still need to recognize the path-based API for production. A job on a company lake looks like spark.read.parquet("s3://bucket/orders/").write.mode("overwrite").parquet("s3://bucket/silver/orders/"). Here you select columns from spark.table("orders") instead.

Why this skill

Choosing CSV for a 200-column event log means every job reparses text and guesses types. Parquet lets Spark read only the columns you select. That is why lakes standardize on Parquet (often inside Delta or Iceberg). If you only ever call spark.table, you still need to know what the catalog is pointing at.

Writes matter for the next person. A thousand tiny files from a high partition count makes the next read slow. A CSV dump loses types. This lesson is the map of formats. Later lessons cover partition counts and Delta logs.

How the code works

A catalog is a front desk: you ask for "orders" by name. A path reader is going to the warehouse aisle yourself with a file address.

Two ways to get a DataFrame
Storage filesCatalog orspark.readDataFrame planAction: show orwrite

table() uses the catalog. read.parquet uses a path. Both are lazy plans on a real cluster.

On a cluster you pick a reader. In this tab you pick a table name.

FormatTypes in the file?Typical useWatch-out
ParquetYesLakes, production tablesDefault choice for pipelines
CSVNo (you guess or supply a schema)Exports, messy vendor dropsSlow parse; always set schema in production
JSONNoAPI payloads, nested eventsSame schema warning; one object per line is common
Columns you will select from orders
order_idorder_totalorder_status104284.50paid104319.00cancelled1044122.40paid

Project only what you need. That habit maps to Parquet column pruning on a cluster.

Worked examples

Run the example below in this tab. Read the input, follow the code, then check the output matches what you expect.

PythonCluster read/write API as comments; this tab uses spark.table
# On a cluster you would write:
# df = spark.read.parquet("s3://lake/orders/")
# df.write.mode("overwrite").parquet("s3://lake/silver/orders/")
#
# CSV / JSON (also cluster API, not run here as a file path):
# spark.read.option("header", True).csv("s3://landing/orders.csv")
# spark.read.json("s3://landing/events.json")

df = spark.table("orders")
result = df.select("order_id", "order_total", "order_status").limit(8)
result.show()

The commented lines are what you will type against S3, GCS, or ADLS. spark.table is the catalog spelling. select plus limit is the exercise: a narrow projection from a registered table.

PythonInspect schema, then project three columns
orders = spark.table("orders")
orders.printSchema()
result = (
    orders
    .select("order_id", "order_total", "order_status")
    .limit(8)
)
result.show()

printSchema() is how you confirm types after a read. In production, CSV and JSON should get an explicit schema so Spark does not guess wrong. Parquet already carries types.

Do not invent a file path in this tab

spark.read.parquet("/some/path") will not open your laptop disk here. Load spark.table("orders") (or ecommerce_events) and select columns. The path API is for a real cluster.

If you know Pandas, here is the translation

pd.read_parquet / to_parquet become spark.read.parquet and df.write.parquet. pd.read_csv becomes spark.read.csv with a schema. pandas writes one file from one process. Spark writes many part files from many executors unless you coalesce.

Copy-paste without reading the output

Run Sample first. If the numbers or row count look wrong, stop and re-read the previous section before changing code.

Common beginner questions

Why is Parquet preferred?

It stores columns separately, compresses well, and keeps types. Spark can skip columns you did not select. CSV is for humans and awkward vendor files.

What does mode overwrite mean?

Replace the output folder's contents. append adds files. error/ignore handle "already exists." You will not write files in this tab.

Is spark.table the same as SQL FROM orders?

Yes in spirit: a named table. The next lesson is Spark SQL, where you can write SELECT ... FROM orders as text.

What comes next

You can load a table. Spark also accepts SQL strings that compile to the same engine as DataFrame calls. That is next.

Practice

Run Sample to project three columns. Then complete Exercise: from spark.table("orders"), select order_id, order_total, and order_status, limit 8, assign to result.

You are practicing catalog load plus column prune. Save path-based read/write for a cluster.

Practicals · load into the editor

After you read the theory, run these in the pane on the right. They execute in this tab, no cluster.

Rate:
Was this useful?
What is a cluster?SQL inside Spark