What you will be able to do
- Save a DataFrame as a persistent managed or external table with saveAsTable and choose the right save mode
- Use partitionBy to write a Hive-style directory layout and read a single partition back
- Combine bucketBy and sortBy on saveAsTable, and say what sortBy actually sorts
- Express the same layout in SQL (PARTITIONED BY, CLUSTERED BY ... SORTED BY) and with DataFrameWriterV2
- Decide from table size whether partitioning will help or hurt retrieval on Databricks
Key concept
Write-time physical layout — When you save a table, the writer methods you chain decide how the files are physically arranged in storage: in directories per column value, in buckets, and sorted within buckets. That arrangement is what later lets a query read less data, so you make the retrieval decision when you write, not when you query.
1.Persistent tables with saveAsTable
A DataFrame lives only as long as your session. To keep its contents so that other queries and sessions can use them, you write it to a table that is registered in the catalog. The DataFrameWriter that you get from df.write does this with saveAsTable(name, format=None, mode=None, partitionBy=None, **options). The simplest call needs only a name:
spark.createDataFrame([
(100, "Alice"), (120, "Bob"), (140, "Tom")],
schema=["age", "name"]
).write.saveAsTable("tblA")That call creates a managed table: the platform decides where its files are stored. An external table points at a storage location that you choose. In SQL you add the EXTERNAL keyword, and the CREATE TABLE reference says that "When creating an external table you must also provide a LOCATION clause." The practical difference shows up when you drop the table. Dropping an external table leaves its files where they are. With the v2 writer, setting tableProperty("location", ...) creates an external (unmanaged) table. If you leave out the format, Databricks uses its default: in SQL, "If USING is omitted, the default is DELTA."
The mode parameter (or the chained .mode() method) controls what happens when the target already exists. On a table, each mode has some details that the exam likes to test:
| Mode | Behaviour when the table exists |
|---|---|
| error / errorifexists (default) | Throws an exception |
| append | Adds the DataFrame's rows; the existing table's format and options are used |
| overwrite | Replaces existing data; the DataFrame schema does not need to match the existing table schema |
| ignore | Silently does nothing |
Column matching also matters. insertInto places columns by position. saveAsTable is different: it "uses column names to find the correct column positions." So if a DataFrame has the right columns in a different order, it still lands correctly through saveAsTable in append mode.
Checkpoint 1 of 8· Check yourself
You append a DataFrame to an existing Parquet table with df.write.mode("append").format("json").saveAsTable("events"). Which format do the appended rows use?
When saveAsTable appends to a table that already exists, it uses that table's own format and options.
“When mode is 'append', if a table already exists, its format and options are used.”Source: docs.databricks.com
Checkpoint 2 of 8· Exam question
A data engineer needs to persist a large `transactions_df` DataFrame as a managed table so that later queries filtering on `transaction_date` can skip irrelevant files entirely, with no bucketing or sorting required. Which code correctly persists the table with this optimization?
Correct answer: A — transactions_df.write.mode("overwrite").partitionBy("transaction_date").saveAsTable("transactions")
- A. `partitionBy` writes the table as Hive-style directories keyed by `transaction_date`, so a later query filtering on that column only reads the matching directory. This is the standard way to get file-skipping without introducing bucketing or sort metadata.
- B. `bucketBy` hashes rows into a fixed number of bucket files rather than creating per-value directories, so a reader cannot tell from the file layout which bucket holds a given `transaction_date`. It optimizes joins and aggregations on the bucketing column, not date-based file skipping.
- C. `sortBy` can only be used together with a preceding `bucketBy` call on the same writer. Calling it alone raises an error before any data is written, so this call never reaches the metastore.
- D. `DataFrameWriter` has no `orderBy` method; ordering a DataFrame before a write is done with `DataFrame.orderBy`, not on the writer object. This call fails with an attribute error at runtime.
- E. `repartition` changes how many in-memory Spark partitions are used during the write, which can affect the number of output files, but it does not create the on-disk directory structure that lets a reader prune by `transaction_date` value.
2.partitionBy: one directory per value
A persistent table helps retrieval when its files are arranged so that a query can skip the ones it doesn't need. partitionBy(*cols) is the first tool for this. It splits the output by the values of the columns you name. The reference describes the result: "the output is laid out on the file system similar to Hive's partitioning scheme." Each distinct value becomes a directory named column=value, and the partition column moves out of the data files and into the path.
import tempfile, os
with tempfile.TemporaryDirectory(prefix="partitionBy") as d:
spark.createDataFrame(
[{"age": 100, "name": "Alice"}, {"age": 120, "name": "Ruifeng Zheng"}]
).write.partitionBy("name").mode("overwrite").format("parquet").save(d)Only age. The read targets the name=Alice directory directly, and the partition value is held in the path, not in the files. The documented output shows a single age column with the value 100.
# Read one partition as a DataFrame.
spark.read.parquet(f"{d}{os.path.sep}name=Alice").show()partitionBy works with path-based writes (save, parquet, orc) and also with saveAsTable. You can chain it as a method, or pass it as the partitionBy argument of saveAsTable. It accepts several columns, as in partitionBy("year", "month"), which gives a directory hierarchy. Choose columns that queries commonly filter on, because a filter on a partition column is what lets whole directories be skipped.
Checkpoint 3 of 8· Fill the gap
This sample saves a managed table laid out in one directory per year. Which method fills the blank?
df.write \
.mode("overwrite") \
.format("parquet") \
. ? ("year") \
.saveAsTable("partitioned_table")On the DataFrameWriter (df.write), partitionBy creates the Hive-style directories. partitionedBy is the DataFrameWriterV2 method, and bucketBy/sortBy do not create per-value directories.
Source: docs.databricks.comCheckpoint 4 of 8· Exam question
Which statement correctly describes the relationship between `sortBy` and `bucketBy` on `DataFrameWriter` when persisting a table?
Correct answer: A — `sortBy` orders rows within each bucket file that `bucketBy` produces; a write that calls `sortBy` without `bucketBy` first raises an error.
- A. Bucketing groups rows into a fixed number of bucket files by hashing a column, and a chained sort orders the rows inside each of those files. Because there is nothing to sort within until buckets exist, the writer raises an error if a sort is requested without bucketing first.
- B. A sort requested this way is scoped to buckets, not the whole dataset, and it cannot be requested on its own; it must follow a bucketing call. There is no single fully sorted output file produced by this mechanism.
- C. Bucket count and column choice are set independently of any sort; the hash function that assigns rows to buckets does not consult sort order at all. The dependency runs the other way: sorting needs buckets to sort within, not the reverse.
- D. These are two distinct writer methods with different jobs: one groups rows into a fixed number of files by hash, the other orders rows inside files that already exist. They are not aliases and using only one of them does not reproduce the other's effect.
- E. Directory-level partitioning is a separate writer method entirely, unrelated to either of these two calls. Neither of these methods governs partition directories, and the in-memory-only claim understates what a bucket-and-sort write actually persists to storage.
Sources4
3.bucketBy and sortBy: sorting inside buckets
Partitioning breaks down when a column has many distinct values, because you get one directory per value. Bucketing takes a different approach. bucketBy(numBuckets, col, *cols) hashes the column values into a fixed number of buckets. The bucketBy reference says the layout is similar to Hive's, "with a different bucket hash function and is not compatible with Hive's bucketing." It also limits where bucketing can be used: it is "Applicable for file-based data sources in combination with DataFrameWriter.saveAsTable." You can't bucket a plain save(path). A bucketed layout exists only as a persistent table.
spark.createDataFrame([
(100, "Alice"), (120, "Alice"), (140, "Bob")],
schema=["age", "name"]
).write.bucketBy(1, "name").sortBy("age").mode(
"overwrite").saveAsTable("sorted_bucketed_table")In this example, sortBy is paired with bucketBy and the result is written with saveAsTable. That matches what the method does: "each bucket" only makes sense once there are buckets. When the example reads the table back, it still calls .sort("age"). The on-disk sort is a storage optimisation, not a promise about the order of query results. The writer's own method table puts the four layout methods side by side:
| Method | What it does |
|---|---|
| partitionBy(*cols) | Partitions the output by the given columns on the file system |
| bucketBy(numBuckets, col, *cols) | Buckets the output by the given columns |
| sortBy(col, *cols) | Sorts the output in each bucket by the given columns on the file system |
| clusterBy(*cols) | Clusters the data by the given columns to optimize query performance |
Checkpoint 5 of 8· Check yourself
A colleague wants a bucketed, sorted copy of a DataFrame and writes df.write.bucketBy(8, "id").sortBy("ts").parquet("/data/out"). What is wrong?
Bucketing (and so sorting within buckets) is documented for file-based sources used together with saveAsTable, so the bucketed layout has to be a persistent table.
“Applicable for file-based data sources in combination with DataFrameWriter.saveAsTable.”Source: docs.databricks.com
Checkpoint 6 of 8· Exam question
A developer writes: ``` df.write.bucketBy(20, "customer_id").sortBy("customer_id").save("/mnt/data/customers_parquet") ``` What happens when this code runs?
Correct answer: A — The write fails immediately with an `AnalysisException`, since `bucketBy` and `sortBy` only work with `saveAsTable` or `insertInto`, not `save`.
- A. Bucketing and its associated sort are table-level metadata that only a catalog entry can store, so they are only accepted on `saveAsTable` and `insertInto`. The file-based `save` path has no catalog entry to record that metadata in, and Spark raises an exception rather than writing partial results.
- B. This describes what a bucketed `saveAsTable` write would produce, but `save` does not accept bucketing at all. The call fails before any files are written, so no such output is produced.
- C. Spark does not silently drop an unsupported writer configuration and continue; a `bucketBy`-plus-`sortBy` call on `save` raises an exception rather than falling back to a default file layout.
- D. `save` never creates a metastore entry regardless of which writer options are chained onto it; it only writes files to the given path. No table registration happens here, and the call fails before reaching that point anyway.
- E. The number of buckets is independent of the cluster's core count and is not capped by available parallelism. The failure here comes from combining bucketing with an unsupported writer method, not from a shuffle-stage parallelism limit.
4.The same layout in SQL and DataFrameWriterV2
Each writer method has a counterpart in CREATE TABLE [USING]. Among the table clauses are PARTITIONED BY, CLUSTER BY, LOCATION and a bucketing clause. In that clause, sorting appears inside the bucket definition, just as sortBy belongs with bucketBy:
clustered_by_clause
{ CLUSTERED BY ( cluster_column [, ...] )
[ SORTED BY ( { sort_column [ ASC | DESC ] } [, ...] ) ]
INTO num_buckets BUCKETS }Note the vocabulary clash. SQL CLUSTERED BY ... INTO n BUCKETS is bucketing, which corresponds to bucketBy. SQL CLUSTER BY and the clusterBy method are a separate clustering feature.
The third route is DataFrameWriterV2, which you reach with df.writeTo(table). Its partition method is spelled partitionedBy, and it can "Partition the output table created by create, createOrReplace, or replace using the given columns or transforms." Transforms let you partition on a value derived from a column without adding that column yourself:
from pyspark.sql.functions import years, months, days
df.writeTo("my_table") \
.partitionedBy(years("date"), months("date")) \
.create()V2 also gives you finer control when rewriting partitioned data. Instead of the all-or-nothing mode("overwrite"), overwritePartitions() replaces only the partitions the DataFrame touches, and overwrite(condition) replaces only the rows that match a filter. The reference also notes that for Databricks tables and Delta Lake, V2 offers "More fine-grained control over partitioning" and "Conditional overwrite capabilities".
Checkpoint 7 of 8· Match them up
Match each DataFrameWriterV2 method to what it does
Tap a term, then the definition that fits it.
These are the descriptions in the DataFrameWriterV2 method table. overwritePartitions works per partition; overwrite works by a row predicate.
“Overwrite all partition for which the data frame contains at least one row”Source: docs.databricks.com
Sources3
5.When partitioning actually helps retrieval
Every earlier section treated partitioning as a way to skip data. On Databricks that only works when the table is large enough. The guidance says that partitioning below the minimum sizes "is likely to negatively affect query performance rather than improve it." The thresholds are based on size:
| Table size | Recommendation |
|---|---|
| Less than 1 TB | Don't partition |
| More than 1 TB to 100 TB | Use liquid clustering instead of partitioning |
| 100 TB or more | Partitioning might help; try liquid clustering first and verify |
| Any partitioned table | Each partition should hold at least 1 GB of data |
The reason is file count. "Tables with fewer, larger partitions tend to outperform tables with many smaller partitions." A high-cardinality partition column creates many small directories, and that costs more than pruning saves. A small Delta table loses little by skipping partitioning: "By using Delta Lake, unpartitioned tables automatically use ingestion time clustering", which gives benefits similar to partitioning on a datetime field.
One more caution concerns the directory layout from the partitionBy section. Spark writes Hive-style partition directories for Parquet, but "Hive-style partitioning is not part of the Delta Lake protocol." Workloads should not reach into a Delta table's column=value folders the way the Parquet example read name=Alice. Use the supported table APIs instead.
Checkpoint 8 of 8· Check yourself
A team partitions a 5 TB Delta table by user_id, which has millions of distinct values. Which statement matches Databricks guidance?
At 1–100 TB Databricks recommends liquid clustering over partitioning. A high-cardinality key would also break the 1 GB-per-partition guideline.
“With more than 1 TB to 100 TB of data, use liquid clustering instead of partitioning.”Source: docs.databricks.com
Sources7
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.If you don't set a mode, saveAsTable appends to (or overwrites) an existing table.Why is that wrong?
The default mode is error/errorifexists, so writing to an existing table without a mode throws an exception.
Covered in Persistent tables with saveAsTable
2.saveAsTable, like insertInto, matches DataFrame columns to table columns by position.Why is that wrong?
saveAsTable matches columns by name. insertInto is the one that matches by position.
Covered in Persistent tables with saveAsTable
3.sortBy on the writer sorts the whole DataFrame, so queries on the table return rows in that order.Why is that wrong?
sortBy orders the data within each bucket on disk. It is a storage layout, and query results still need an explicit sort.
Covered in bucketBy and sortBy: sorting inside buckets
4.Partitioning always speeds up queries, so even small tables should be partitioned on their filter column.Why is that wrong?
Databricks says not to partition tables under 1 TB and to keep each partition at 1 GB or more. Below those sizes, partitioning tends to hurt performance.
Covered in When partitioning actually helps retrieval
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.https://docs.databricks.com/aws/en/pyspark/reference/classes/dataframewriter/saveAsTableOfficial docs
“When mode is 'overwrite', the schema of the DataFrame does not need to match the existing table schema.”
↩︎ Persistent tables with saveAsTable“If the table already exists, the behavior depends on the mode parameter (default is to throw an exception).”
↩︎ Exam trap 1“Unlike DataFrameWriter.insertInto, DataFrameWriter.saveAsTable uses column names to find the correct column positions.”
↩︎ Exam trap 2“If the table already exists, the behavior depends on the mode parameter (default is to throw an exception).”
↩︎ Prediction“When mode is 'append', if a table already exists, its format and options are used.”
↩︎ Checkpoint - 2.https://docs.databricks.com/aws/en/sql/language-manual/sql-ref-syntax-ddl-create-table-usingOfficial docs
“When an external table is dropped the files at the LOCATION will not be dropped.”
↩︎ Persistent tables with saveAsTable“When creating an external table you must also provide a LOCATION clause.”
↩︎ Persistent tables with saveAsTable“If USING is omitted, the default is DELTA.”
↩︎ Persistent tables with saveAsTable - 3.
“to create an EXTERNAL (unmanaged) table.”
↩︎ Persistent tables with saveAsTable“Partition the output table created by create, createOrReplace, or replace using the given columns or transforms.”
↩︎ The same layout in SQL and DataFrameWriterV2“More fine-grained control over partitioning”
↩︎ The same layout in SQL and DataFrameWriterV2“Overwrite all partition for which the data frame contains at least one row”
↩︎ Checkpoint - 4.https://docs.databricks.com/aws/en/pyspark/reference/classes/dataframewriter/partitionByOfficial docs
“If specified, the output is laid out on the file system similar to Hive's partitioning scheme.”
↩︎ partitionBy: one directory per value“Partitions the output by the given columns on the file system.”
↩︎ Key concept - 5.
“with a different bucket hash function and is not compatible with Hive's bucketing.”
↩︎ bucketBy and sortBy: sorting inside buckets“Applicable for file-based data sources in combination with DataFrameWriter.saveAsTable.”
↩︎ Checkpoint - 6.
“Sorts the output in each bucket by the given columns on the file system.”
↩︎ bucketBy and sortBy: sorting inside buckets“Sorts the output in each bucket by the given columns on the file system.”
↩︎ Exam trap 3 - 7.https://docs.databricks.com/aws/en/tables/partitionsOfficial docs
“Tables with fewer, larger partitions tend to outperform tables with many smaller partitions.”
↩︎ When partitioning actually helps retrieval“Hive-style partitioning is not part of the Delta Lake protocol”
↩︎ When partitioning actually helps retrieval“By using Delta Lake, unpartitioned tables automatically use ingestion time clustering.”
↩︎ When partitioning actually helps retrieval“With less than 1 TB of data, don't partition.”
↩︎ Exam trap 4“With less than 1 TB of data, don't partition.”
↩︎ Prediction“With more than 1 TB to 100 TB of data, use liquid clustering instead of partitioning.”
↩︎ Checkpoint