ML Engineer MasterClass (October) | 4 seats left

Spark · Files and schemas
AmazonAmazon Analytics
Amazon · Read and Write Data

Files and schemas

Read JSON and CSV with deliberate schemas while keeping temporary-file examples self-contained.

Step 1 of 6 · Learn

Read a JSON dataset

Spark writes JSON as a directory of part files, then reads the directory as a dataset. This example infers a schema and explicitly selects output columns; file and field discovery order should not define your result. The temporary directory is unique to this run. We collect only these six teaching rows to detach the result before cleanup; collecting large datasets this way is not suitable for production.

Lesson reference: PySpark and Scala

Read a JSON dataset

Spark writes JSON as a directory of part files, then reads the directory as a dataset. This example infers a schema and explicitly selects output columns; file and field discovery order should not define your result. The temporary directory is unique to this run. We collect only these six teaching rows to detach the result before cleanup; collecting large datasets this way is not suitable for production.

PySpark example

from tempfile import TemporaryDirectory
from pyspark.sql.functions import col

# This directory belongs only to this example.
with TemporaryDirectory(prefix="spark-lesson-") as directory:
    path = directory + "/data"
    products.write.json(path)
    loaded = spark.read.json(path).select("product_id", "product_name")
    # Detach these few rows before deleting their source files.
    result = spark.createDataFrame(loaded.collect(), loaded.schema).orderBy("product_id")
    result.show(truncate=False)

Scala example

import org.apache.spark.sql.functions._
import java.nio.file.Files
import org.apache.hadoop.fs.Path

val directory = Files.createTempDirectory("spark-lesson-")
val root = new Path(directory.toUri)
val fs = root.getFileSystem(spark.sparkContext.hadoopConfiguration)
val path = directory.resolve("data").toString
val result = try {
    products.write.json(path)
    val loaded = spark.read.json(path).select("product_id", "product_name")
    // Detach these few rows before deleting their source files.
    val output = spark.createDataFrame(
        spark.sparkContext.parallelize(loaded.collect().toSeq), loaded.schema
    ).orderBy("product_id")
    output.show(false)
    output
} finally {
    fs.delete(root, true) // Only the unique directory created above.
}

Supply an explicit schema

An explicit schema controls field types rather than relying on inference. Here quantity is INT and line_total is DECIMAL(10,2), preserving the monetary type. Schema choice is not a complete data-quality policy: missing or malformed fields still need deliberate handling. This fixture matches its schema, and line_total describes a whole order line, not a unit price.

PySpark example

from tempfile import TemporaryDirectory
from pyspark.sql.functions import col

# This directory belongs only to this example.
with TemporaryDirectory(prefix="spark-lesson-") as directory:
    path = directory + "/data"
    order_items.write.json(path)
    loaded = spark.read.schema("order_item_id LONG, order_id LONG, product_id LONG, quantity INT, line_total DECIMAL(10,2)").json(path).select("order_item_id", "line_total")
    # Detach these few rows before deleting their source files.
    result = spark.createDataFrame(loaded.collect(), loaded.schema).orderBy("order_item_id")
    result.show(truncate=False)

Scala example

import org.apache.spark.sql.functions._
import java.nio.file.Files
import org.apache.hadoop.fs.Path

val directory = Files.createTempDirectory("spark-lesson-")
val root = new Path(directory.toUri)
val fs = root.getFileSystem(spark.sparkContext.hadoopConfiguration)
val path = directory.resolve("data").toString
val result = try {
    order_items.write.json(path)
    val loaded = spark.read.schema("order_item_id LONG, order_id LONG, product_id LONG, quantity INT, line_total DECIMAL(10,2)").json(path).select("order_item_id", "line_total")
    // Detach these few rows before deleting their source files.
    val output = spark.createDataFrame(
        spark.sparkContext.parallelize(loaded.collect().toSeq), loaded.schema
    ).orderBy("order_item_id")
    output.show(false)
    output
} finally {
    fs.delete(root, true) // Only the unique directory created above.
}

Read CSV headers deliberately

CSV needs explicit choices about headers, separators, and types. This example writes a header and reads it with header=true and a matching schema. The header is not a customer row. Schema fields follow the written column order; real files with reordered or mismatched headers require validation, not blind assumptions.

PySpark example

from tempfile import TemporaryDirectory
from pyspark.sql.functions import col

# This directory belongs only to this example.
with TemporaryDirectory(prefix="spark-lesson-") as directory:
    path = directory + "/data"
    customers.write.option("header", "true").csv(path)
    loaded = spark.read.schema("customer_id LONG, customer_name STRING, city STRING").option("header", "true").csv(path).select("customer_id", "customer_name")
    # Detach these few rows before deleting their source files.
    result = spark.createDataFrame(loaded.collect(), loaded.schema).orderBy("customer_id")
    result.show(truncate=False)

Scala example

import org.apache.spark.sql.functions._
import java.nio.file.Files
import org.apache.hadoop.fs.Path

val directory = Files.createTempDirectory("spark-lesson-")
val root = new Path(directory.toUri)
val fs = root.getFileSystem(spark.sparkContext.hadoopConfiguration)
val path = directory.resolve("data").toString
val result = try {
    customers.write.option("header", "true").csv(path)
    val loaded = spark.read.schema("customer_id LONG, customer_name STRING, city STRING").option("header", "true").csv(path).select("customer_id", "customer_name")
    // Detach these few rows before deleting their source files.
    val output = spark.createDataFrame(
        spark.sparkContext.parallelize(loaded.collect().toSeq), loaded.schema
    ).orderBy("customer_id")
    output.show(false)
    output
} finally {
    fs.delete(root, true) // Only the unique directory created above.
}
example.pyPySpark
1. Write a unique temporary datasetOnly the example’s own directory is created.
2. Read and projectChoose reader options, schema, and output columns.
3. Detach and clean upMaterialize these few rows before removing the temporary source.
Source products6 rows
product_idproduct_namecategory
201Wireless keyboardElectronics
202Laptop standOffice
203Desk lampHome
204USB-C hubElectronics
205Travel backpackTravel
206Notebook setOffice
Read a JSON dataset
Result6 rows
product_idproduct_name
201Wireless keyboard
202Laptop stand
203Desk lamp
204USB-C hub
205Travel backpack
206Notebook set
The reader loads all part files in the temporary data directory. The returned rows remain usable after that directory is removed.

The example is loaded in the editor. Run it as written, then try a small change.

Runs on the Spark backend. First startup may take a moment.

Run your code to see the result.