Append another batch
append adds data to an existing path; it does not deduplicate keys or provide an upsert. This example writes customers 101–103, then appends a non-overlapping batch. Repeating append would add duplicates. Every run creates its own temporary directory, so testing this lesson does not write to a shared dataset.
Lesson reference: PySpark and Scala
Append another batch
append adds data to an existing path; it does not deduplicate keys or provide an upsert. This example writes customers 101–103, then appends a non-overlapping batch. Repeating append would add duplicates. Every run creates its own temporary directory, so testing this lesson does not write to a shared dataset.
PySpark example
from tempfile import TemporaryDirectory
from pyspark.sql.functions import col, count, lit
with TemporaryDirectory(prefix="spark-storage-") as directory:
path = directory + "/data"
compact_path = directory + "/compact"
customers.filter(col("customer_id") <= 103).write.parquet(path)
customers.filter((col("customer_id") > 103) & (col("customer_id") <= 104)).write.mode("append").parquet(path)
loaded = spark.read.parquet(path).select("customer_id", "customer_name").orderBy("customer_id")
# Detach only this tiny teaching result before removing its 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-storage-")
val root = new Path(directory.toUri)
val fs = root.getFileSystem(spark.sparkContext.hadoopConfiguration)
val path = directory.resolve("data").toString
val compact_path = directory.resolve("compact").toString
val result = try {
customers.filter(col("customer_id") <= 103).write.parquet(path)
customers.filter((col("customer_id") > 103) && (col("customer_id") <= 104)).write.mode("append").parquet(path)
val loaded = spark.read.parquet(path).select("customer_id", "customer_name").orderBy("customer_id")
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 this example’s unique temporary directory.
}
Replace existing contents
overwrite replaces data at the target path according to the writer and storage behavior. Here it replaces the whole unpartitioned temporary dataset, not just matching customer IDs. It is not an UPDATE statement, and generic file writes are not a promise of transactional replacement on every storage system. Never point this example at a shared path.
PySpark example
from tempfile import TemporaryDirectory
from pyspark.sql.functions import col, count, lit
with TemporaryDirectory(prefix="spark-storage-") as directory:
path = directory + "/data"
compact_path = directory + "/compact"
customers.write.parquet(path)
customers.filter(col("customer_id") == 101).write.mode("overwrite").parquet(path)
loaded = spark.read.parquet(path).select("customer_id", "customer_name").orderBy("customer_id")
# Detach only this tiny teaching result before removing its 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-storage-")
val root = new Path(directory.toUri)
val fs = root.getFileSystem(spark.sparkContext.hadoopConfiguration)
val path = directory.resolve("data").toString
val compact_path = directory.resolve("compact").toString
val result = try {
customers.write.parquet(path)
customers.filter(col("customer_id") === 101).write.mode("overwrite").parquet(path)
val loaded = spark.read.parquet(path).select("customer_id", "customer_name").orderBy("customer_id")
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 this example’s unique temporary directory.
}
Distinguish ignore from overwrite
ignore leaves an existing path unchanged; it does not merge new records. The default error-if-exists mode would reject a write to that existing path. This example starts with three customers and then attempts to write all six using ignore. Compare the resulting data, not just whether the write raised an error.
PySpark example
from tempfile import TemporaryDirectory
from pyspark.sql.functions import col, count, lit
with TemporaryDirectory(prefix="spark-storage-") as directory:
path = directory + "/data"
compact_path = directory + "/compact"
customers.filter(col("customer_id") <= 103).write.parquet(path)
customers.write.mode("ignore").parquet(path)
loaded = spark.read.parquet(path).select("customer_id", "customer_name").orderBy("customer_id")
# Detach only this tiny teaching result before removing its 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-storage-")
val root = new Path(directory.toUri)
val fs = root.getFileSystem(spark.sparkContext.hadoopConfiguration)
val path = directory.resolve("data").toString
val compact_path = directory.resolve("compact").toString
val result = try {
customers.filter(col("customer_id") <= 103).write.parquet(path)
customers.write.mode("ignore").parquet(path)
val loaded = spark.read.parquet(path).select("customer_id", "customer_name").orderBy("customer_id")
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 this example’s unique temporary directory.
}