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.
}