ML Engineer MasterClass (October) | 4 seats left

Spark · Write modes
AmazonAmazon Analytics
Amazon · Read and Write Data

Write modes

Understand append, overwrite, and ignore using disposable paths only.

Step 1 of 6 · Learn

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.
}
example.pyPySpark
1. Initial contentsCreate a new temporary destination.
2. Second writeThe requested mode determines whether data is added, replaced, or ignored.
3. Read back and clean upCheck the resulting records, then remove only this example’s directory.
Source customers6 rows
customer_idcustomer_namecity
101AriSeattle
102BoAustin
103CamChicago
104DeeBoston
105EliDenver
106FranPortland
Append another batch
Result4 rows
customer_idcustomer_name
101Ari
102Bo
103Cam
104Dee
The example produces four customers after one appended row. Append itself does not check that keys are new.

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.