CertSafari

    Free Databricks Certified Associate Developer for Apache Spark Sample Questions

    35 free sample questions from our bank of 352+, covering every exam domain, with answers and detailed explanations. Updated September 2026.

    Domain 1: Apache Spark Architecture and Components

    Subdomain 1.2: Identify the role of core components of Apache Spark™'s Architecture, including cluster, driver node, worker nodes/executors, CPU cores, and memory.

    1.A PySpark application calls three actions in sequence: `df1.count()`, `df2.write.parquet(path)`, and `df3.collect()`. On the Spark UI's Jobs tab, how many jobs and what triggers each one?

    1. A.Three separate jobs, because Spark launches one job per action, and each job is then broken into stages and tasks that run on the executors.
    2. B.One job, because Spark batches every action in a script into a single job that only starts once the driver reaches the end of the file.
    3. C.Three jobs, but only the two write actions actually launch tasks on executors; `count()` is resolved entirely inside the driver's JVM.
    4. D.Two jobs, because `collect()` and `count()` are lazy aggregation calls that Spark merges into the job triggered by the `write.parquet` action.
    5. E.One job per unique DataFrame lineage, so `df1`, `df2`, and `df3` each contribute to the same job since they share the same SparkSession.
    Show answer & explanation

    Correct answer: A — Three separate jobs, because Spark launches one job per action, and each job is then broken into stages and tasks that run on the executors.

    • A. Correct. Each action triggers Spark's lazy-evaluated transformations to actually execute, and the driver submits a distinct job per action, which is then split into stages and tasks.
    • B. Incorrect. Spark does not defer job submission to the end of the script; each action independently triggers its own job as soon as it is called.
    • C. Incorrect. `count()` still requires reading and aggregating the underlying partitions across executors; it is not resolved purely inside the driver without executor work.
    • D. Incorrect. Each action produces its own job in the Spark UI; `collect()` and `count()` are not merged into a job triggered by an unrelated `write.parquet` call on a different DataFrame.
    • E. Incorrect. Job counting is tied to actions, not to shared lineage or a shared SparkSession; three actions on three DataFrames still produce three jobs even under one session.

    Subdomain 1.4: Explain the Apache Spark™ Architecture execution hierarchy.

    2.A platform team submits a long-running Spark application to a YARN cluster using `spark-submit`. They want the driver process itself to run on a node inside the cluster so that the application keeps running even if the machine that issued the submit command disconnects. Which deploy mode should they specify?

    1. A.Client mode
    2. B.Cluster mode
    3. C.Local mode
    4. D.Standalone mode
    5. E.Interactive mode
    Show answer & explanation

    Correct answer: B — Cluster mode

    • A. This is incorrect because in this mode the driver runs on the machine that issued the submit command, so the application would stop if that machine disconnected.
    • B. This is correct because this deploy mode launches the driver process on a node inside the cluster itself, so the application continues running independently of the machine that originally issued the submit command.
    • C. This is incorrect because this mode runs the driver and executors within a single JVM on one machine and is not used for submitting a job to a YARN cluster.
    • D. This is incorrect because this term refers to Spark's built-in cluster manager, not to a deploy mode that controls where the driver process runs.
    • E. This is incorrect because this is not a recognized Spark deploy mode; the deploy mode options control driver placement, not an interactive session type.

    Subdomain 1.3: Describe the architecture of Apache Spark™, including DataFrame and Dataset concepts, SparkSession lifecycle, caching, storage levels, and garbage collection.

    3.A developer is reviewing this notebook cell to understand which lines actually trigger a Spark job: ```python df1 = spark.read.parquet("/data/raw") df2 = df1.withColumn("year", F.year("event_time")) df3 = df2.filter(df2.year == 2026) total = df3.count() df3.write.mode("overwrite").parquet("/data/2026") ``` Choose 2 answers: which two lines trigger the execution of a Spark job?(Select 2)

    1. A.`total = df3.count()`, because `count()` is an action that forces Spark to execute the accumulated transformations and return a value to the driver.
    2. B.`df3.write.mode("overwrite").parquet("/data/2026")`, because writing output is an action that forces Spark to execute the transformation chain and materialize files.
    3. C.`df1 = spark.read.parquet("/data/raw")`, because reading a file eagerly loads all of its records straight into the driver's memory before any transformation can even be defined.
    4. D.`df2 = df1.withColumn("year", F.year("event_time"))`, because adding a derived column requires Spark to immediately materialize the new column's values.
    5. E.`df3 = df2.filter(df2.year == 2026)`, because filtering rows is treated as an action in Spark and therefore executes a job as soon as the line runs.
    Show answer & explanation

    Correct answers: A, B — `total = df3.count()`, because `count()` is an action that forces Spark to execute the accumulated transformations and return a value to the driver.; `df3.write.mode("overwrite").parquet("/data/2026")`, because writing output is an action that forces Spark to execute the transformation chain and materialize files.

    • A. This is correct: `count()` is a classic action that returns a value to the driver, forcing Spark to execute the read, column addition, and filter that precede it in the lazy plan.
    • B. This is also correct: writing to storage is an action, so this line triggers its own job that executes the full transformation chain again to produce and persist the output files.
    • C. `spark.read.parquet()` only builds a logical plan describing the source; it does not eagerly load records into driver memory, so this line alone does not trigger a job.
    • D. `withColumn()` is a transformation, not an action; it lazily extends the logical plan with a new column expression without materializing any values at this point.
    • E. `filter()` is a transformation, not an action; row filtering is recorded in the plan and only executes once a downstream action like `count()` or `write()` runs.

    Subdomain 1.5: Configure Spark partitioning in distributed data processing, including shuffles and partitions

    4.Which two statements about `coalesce()` are correct? (Choose 2 answers)(Select 2)

    1. A.`coalesce()` can only decrease a DataFrame's partition count and can never raise it beyond however many partitions already exist.
    2. B.`coalesce()` avoids a full shuffle by combining existing partitions on the same executor, moving less data than `repartition()`.
    3. C.`coalesce()` always redistributes rows evenly across the resulting partitions, matching the guarantees that `repartition()` provides.
    4. D.`coalesce()` triggers a full shuffle across the cluster identical to calling `repartition()` with the same target partition count.
    5. E.`coalesce()` recalculates the DataFrame's schema, dropping any columns that a subsequent action does not reference anywhere downstream.
    Show answer & explanation

    Correct answers: A, B — `coalesce()` can only decrease a DataFrame's partition count and can never raise it beyond however many partitions already exist.; `coalesce()` avoids a full shuffle by combining existing partitions on the same executor, moving less data than `repartition()`.

    • A. Coalesce is designed only to merge existing partitions together, so it cannot raise the partition count above what the DataFrame already has; requesting more partitions than currently exist has no effect.
    • B. By combining nearby partitions on the same executor instead of redistributing every row across the network, coalesce moves far less data than a full shuffle, which is exactly why it is the cheaper way to reduce partition count.
    • C. Because coalesce merges whatever partitions already exist rather than hashing rows into new ones, the resulting partitions can be uneven in size, unlike the more even distribution a full shuffle in `repartition` typically produces.
    • D. A full shuffle moving every row across the cluster is the behavior of `repartition`, not `coalesce`; coalesce specifically avoids that cost when it is only reducing the partition count.
    • E. Coalesce only changes how a DataFrame's existing partitions are grouped together; it has no effect on the DataFrame's schema and does not drop or recompute any columns.

    Subdomain 1.1: Identify the advantages and challenges of implementing Spark.

    5.A retail company runs separate Hadoop MapReduce jobs for nightly ETL, a dedicated SQL engine for analyst reporting, and a standalone cluster for demand-forecasting models, and wants to consolidate onto one platform without losing performance. Which advantage of Apache Spark most directly supports this consolidation?

    1. A.Spark exposes a unified engine where Spark SQL, DataFrames, Structured Streaming, and MLlib share the same execution model, so ETL, reporting, and forecasting run on one cluster.
    2. B.Spark rewrites every submitted job into a MapReduce-compatible execution plan, so the existing Hadoop tooling keeps running unmodified while it still gains in-memory execution speed.
    3. C.Spark enforces schema-on-write for every file format it reads, which removes the analyst reporting layer because raw files become query-ready without any engine.
    4. D.Spark bundles a built-in charting dashboard that replaces the reporting tool entirely, so analysts no longer need a separate SQL engine to produce reports.
    5. E.Spark's cluster manager provisions a dedicated GPU pool automatically whenever a machine learning job is submitted, replacing the standalone forecasting cluster.
    Show answer & explanation

    Correct answer: A — Spark exposes a unified engine where Spark SQL, DataFrames, Structured Streaming, and MLlib share the same execution model, so ETL, reporting, and forecasting run on one cluster.

    • A. This is correct: Spark's unified engine lets the same cluster and APIs handle batch ETL, SQL reporting, and MLlib model training, which is exactly the consolidation the company wants. This eliminates the need to operate three separate specialized systems.
    • B. Spark does not translate jobs into MapReduce plans; it compiles work into its own DAG of stages and tasks executed by its own scheduler. Framing Spark as a MapReduce-compatibility layer misstates how its execution engine works.
    • C. Spark does not enforce schema-on-write, and reading raw files does not remove the need for a query engine to serve reporting workloads. Schema inference and enforcement are configurable, not automatic replacements for a SQL layer.
    • D. Spark itself has no built-in charting dashboard; visualization is handled by separate tools such as notebooks or BI products that connect to Spark. This option describes a capability Spark does not provide.
    • E. Spark does not auto-provision GPU pools when ML code runs; hardware allocation depends on the cluster manager and the cluster configuration chosen by an administrator. This overstates what happens automatically.

    Subdomain 1.6: Describe the execution patterns of the Apache Spark™ engine, including actions, transformations, and lazy evaluation.

    6.A data engineer runs the following code in a Databricks notebook cell: ```python df = spark.read.parquet("/mnt/sales/transactions") filtered = df.filter(df.amount > 1000) grouped = filtered.groupBy("region").sum("amount") ``` No further code is executed in this cell. What happens when this cell finishes running?

    1. A.No Spark job is submitted to the cluster; Spark only records `filtered` and `grouped` as logical plans built on top of `df`, since every call in the cell is a transformation.
    2. B.Spark submits one job that reads the Parquet files and applies the filter, then submits a second job later to compute the grouped sum once that first job finishes.
    3. C.Spark submits a single job immediately because `groupBy` is a wide transformation that always forces materialization of its parent DataFrame to build shuffle partitions.
    4. D.Spark raises a runtime error when the cell finishes executing because `grouped` was never assigned to an action, leaving the driver unable to resolve an incomplete DataFrame plan.
    5. E.Spark reads the Parquet file's footer statistics to execute the filter eagerly, but defers only the final `groupBy` aggregation until an action is later called.
    Show answer & explanation

    Correct answer: A — No Spark job is submitted to the cluster; Spark only records `filtered` and `grouped` as logical plans built on top of `df`, since every call in the cell is a transformation.

    • A. This is correct: `read`, `filter`, and `groupBy`/`sum` are all transformations, so the cell only builds up a logical plan referencing `df`. No job is submitted to the cluster because no action was called.
    • B. This is incorrect because no action appears anywhere in the cell, so no job is submitted at all, let alone two separate ones. Filtering and grouping stay purely as plan-building steps until something forces execution.
    • C. This is incorrect because being a wide transformation does not make `groupBy` eager. Wide transformations still only add to the logical plan and only cause a shuffle when a later action forces the plan to run.
    • D. This is incorrect because leaving a DataFrame unused by an action is not an error condition in Spark; the cell simply finishes with `grouped` holding an unexecuted plan, and no exception is raised.
    • E. This is incorrect because Spark does not selectively execute some transformations eagerly based on file metadata. Every transformation in the chain, including the filter, remains unexecuted until an action is called.

    Subdomain 1.7: Identify the features of the Apache Spark Modules, including Core, Spark SQL, DataFrames, Pandas API on Spark, Structured Streaming, and MLib.

    7.A developer wants to chain a `StringIndexer` that encodes a categorical `category` column and a `LogisticRegression` estimator into a single reusable workflow that can be fit on training data and later applied consistently to new data: ```python indexer = StringIndexer(inputCol="category", outputCol="categoryIndex") lr = LogisticRegression(featuresCol="features", labelCol="label") ``` Which snippet correctly builds and fits that workflow?

    1. A.`pipeline = Pipeline(stages=[indexer, lr])` `model = pipeline.fit(train_df)`, which chains both stages in order and trains them together with one `fit()` call.
    2. B.`model = lr.fit(indexer.fit(train_df))`, which passes a fitted `StringIndexerModel` object into `lr.fit()` instead of the transformed DataFrame of features that estimator expects.
    3. C.`model = [indexer, lr].fit(train_df)`, which calls `.fit()` on a plain Python list, and a list has no such method, so this raises an `AttributeError`.
    4. D.`model = lr.fit(train_df, transformers=[indexer])`, which passes a `transformers` keyword that `LogisticRegression.fit()` does not accept, so this raises a `TypeError`.
    5. E.`model = Pipeline(indexer, lr).fit(train_df)`, which passes the stages as separate positional arguments instead of the required `stages` keyword list, so this raises a `TypeError`.
    Show answer & explanation

    Correct answer: A — `pipeline = Pipeline(stages=[indexer, lr])` `model = pipeline.fit(train_df)`, which chains both stages in order and trains them together with one `fit()` call.

    • A. `Pipeline` accepts an ordered list of stages, and a single `fit()` call runs each stage's transformation or training in sequence, which is exactly the reusable, chained workflow the developer wants.
    • B. `indexer.fit(train_df)` produces a fitted `StringIndexerModel`, not a transformed DataFrame; `lr.fit()` needs an actual DataFrame of features and labels, so feeding it a model object is a type mismatch.
    • C. A Python list is not a Spark ML construct and has no `.fit()` method, so this line fails immediately rather than chaining the indexer and the classifier.
    • D. `LogisticRegression.fit()` only accepts a DataFrame (and optional params), not a `transformers` argument, so this call does not exist as written and raises a `TypeError`.
    • E. `Pipeline` expects its stages to be passed as a `stages` keyword argument holding a list, not as separate positional arguments, so constructing it this way raises a `TypeError` before fitting begins.

    Domain 2: Using Spark SQL

    Subdomain 2.2: Execute SQL queries directly on files, including ORC Files, JSON Files, CSV Files, Text Files, and Delta Files, and understand the different save modes for outputting data in Spark SQL.

    8.A streaming ingestion job writes a fresh batch of transactions every hour to a Delta table and must add each batch alongside all previously written data, never replacing it. Which write call implements this correctly?

    1. A.`batch_df.write.format("delta").mode("overwrite").save(path)`, which commits each hour's rows as the sole contents of the table, replacing everything written before it
    2. B.`batch_df.write.format("delta").mode("errorifexists").save(path)`, which raises an exception on every run after the first because the target path already contains data
    3. C.`batch_df.write.format("delta").save(path)` with no mode specified, which defaults to merging incoming rows with existing rows by matching primary key values
    4. D.`batch_df.write.format("delta").mode("ignore").save(path)`, which silently discards every hourly batch after the first because the target path already has data in it
    5. E.`batch_df.write.format("delta").mode("append").save(path)`, which commits the new rows as an additional atomic transaction on top of the existing table history
    Show answer & explanation

    Correct answer: E — `batch_df.write.format("delta").mode("append").save(path)`, which commits the new rows as an additional atomic transaction on top of the existing table history

    • A. Overwrite mode replaces the entire table contents with the new batch, so each hourly run would erase every transaction written in earlier runs rather than accumulating them.
    • B. Errorifexists mode is the default guard against accidentally clobbering an existing table and raises an exception once the path already holds data, so it cannot support repeated hourly writes at all.
    • C. Spark's DataFrameWriter has no automatic key-based merge behavior triggered by omitting the mode; without an explicit mode it falls back to `errorifexists`, which would fail after the first successful write.
    • D. Ignore mode is designed to silently skip a write only when the target already contains data, so it would drop every batch after the first hour instead of accumulating the streaming transactions.
    • E. Append mode commits each new batch as an additional atomic transaction without touching existing rows, so every hour's transactions accumulate on top of the prior history exactly as the job requires.

    Subdomain 2.4: Register DataFrames as temporary views in Spark SQL, allowing them to be queried with SQL syntax.

    9.What is the key difference between `DataFrame.createTempView(name)` and `DataFrame.createOrReplaceTempView(name)` in Spark SQL?

    1. A.`createTempView` raises an exception if a temporary view with that name already exists, while `createOrReplaceTempView` silently overwrites it.
    2. B.`createTempView` registers the view in the `global_temp` database, while `createOrReplaceTempView` registers it in the session's default database.
    3. C.`createTempView` persists the view to the metastore across cluster restarts, while `createOrReplaceTempView` keeps the view in memory only.
    4. D.`createTempView` requires the DataFrame to be cached first, while `createOrReplaceTempView` can register an uncached DataFrame directly.
    5. E.`createTempView` accepts a SQL query string as its argument, while `createOrReplaceTempView` accepts only a list of column names.
    Show answer & explanation

    Correct answer: A — `createTempView` raises an exception if a temporary view with that name already exists, while `createOrReplaceTempView` silently overwrites it.

    • A. This is the documented behavioral distinction: the non-replacing method fails with an already-exists error on a duplicate name, while the replacing method overwrites the prior registration without error.
    • B. Both methods register a session-scoped temporary view in the session's default namespace; neither one uses the `global_temp` database, which is reserved for the separate global temp view methods.
    • C. Neither method writes anything to the metastore; both create purely in-memory, session-scoped catalog entries that are lost when the session ends.
    • D. Neither method requires the DataFrame to be cached beforehand; registering a view only binds a name to the DataFrame's logical plan, independent of caching.
    • E. Both methods take only the view name string as their argument; neither accepts a SQL query string or a column list.

    Subdomain 2.1: Utilize common data sources such as JDBC, files, etc., to efficiently read from and write to Spark DataFrames using Spark SQL, including overwriting and partitioning by column.

    10.Two notebooks are attached to the same cluster and run in separate Spark sessions. A developer in the first notebook wants a DataFrame to be queryable by name from the second notebook without writing it to storage. Which two statements are true about achieving this? Choose 2 answers.(Select 2)

    1. A.Calling `df.createOrReplaceTempView("orders")` in the first notebook will not make `orders` visible to a query in the second notebook, because a plain temp view is scoped to the session that created it.
    2. B.Calling `df.createOrReplaceGlobalTempView("orders")` in the first notebook makes the view visible from the second one too, because global temp views are shared across sessions in the same Spark application.
    3. C.The second notebook must query the view as `SELECT * FROM orders`, because the `global_temp` prefix is only required when calling `spark.sql` from Python, not from a SQL cell.
    4. D.Calling `df.createOrReplaceTempView("orders")` in the first notebook automatically registers `orders` as a permanent, shared entry in the metastore, so the second notebook can query it like any managed table.
    5. E.Global temp views persist to disk under a `global_temp` directory, so the second notebook can still read them even after the cluster restarts and the original session ends.
    Show answer & explanation

    Correct answers: A, B — Calling `df.createOrReplaceTempView("orders")` in the first notebook will not make `orders` visible to a query in the second notebook, because a plain temp view is scoped to the session that created it.; Calling `df.createOrReplaceGlobalTempView("orders")` in the first notebook makes the view visible from the second one too, because global temp views are shared across sessions in the same Spark application.

    • A. Correct — a session-scoped temp view created with `createOrReplaceTempView` only exists for the Spark session that created it, so a different session's queries cannot see it by name.
    • B. Correct — a global temp view is tied to the `global_temp` system database, which is shared by every session running in the same Spark application, so other sessions can query it.
    • C. Incorrect — the `global_temp.` prefix is required for every session accessing a global temp view, regardless of whether the query runs from a SQL cell, `spark.sql`, or another language.
    • D. Incorrect — a plain temp view is an in-memory, session-local logical plan; it is never written to the metastore, so no other session or later job can discover it as a table.
    • E. Incorrect — global temp views live only in memory for the lifetime of the Spark application; they are not persisted to disk and disappear once the application (and its cluster) stops.

    Subdomain 2.3: Save data to persistent tables while applying sorting and partitioning to optimize data retrieval.

    11.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?

    1. A.transactions_df.write.mode("overwrite").partitionBy("transaction_date").saveAsTable("transactions")
    2. B.transactions_df.write.mode("overwrite").bucketBy(50, "transaction_date").saveAsTable("transactions")
    3. C.transactions_df.write.mode("overwrite").sortBy("transaction_date").saveAsTable("transactions")
    4. D.transactions_df.write.mode("overwrite").orderBy("transaction_date").saveAsTable("transactions")
    5. E.transactions_df.repartition("transaction_date").write.mode("overwrite").saveAsTable("transactions")
    Show answer & explanation

    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.

    Domain 3: Developing Apache Spark™ DataFrame/DataSet API Applications

    Subdomain 3.1: Manipulate columns, rows, and table structures by adding, dropping, splitting, renaming column names, applying filters, and exploding arrays.

    12.An analyst has a DataFrame `orders` with columns `order_id`, `amount`, and `region`. They need to add a new column `amount_tier` that labels rows as `"high"` when `amount` is greater than 1000 and `"standard"` otherwise, without mutating the original DataFrame: ```python orders_tiered = orders.____( "amount_tier", when(col("amount") > 1000, "high").otherwise("standard") ) ``` Which method correctly completes this task?

    1. A.select
    2. B.withColumnRenamed
    3. C.withColumn
    4. D.filter
    5. E.withColumnsRenamed
    Show answer & explanation

    Correct answer: C — withColumn

    • A. select projects a set of column expressions but does not take a target column name and a value expression as two positional arguments the way this call is written, so this signature does not fit.
    • B. withColumnRenamed only accepts an existing column name and a new name string; it cannot evaluate a when/otherwise expression, so it cannot produce a computed column.
    • C. withColumn takes a target column name and a Column expression, adding the result as a new column while returning a new DataFrame, which matches this call exactly.
    • D. filter restricts which rows are kept based on a boolean condition; it does not add or compute a new column, so it does not fit this signature.
    • E. withColumnsRenamed takes a dictionary mapping existing names to new names for bulk renaming, not a name plus a computed expression, so it does not match this call.

    Subdomain 3.2: Perform data deduplication and validation operations on DataFrames.

    13.A `sensor_df` DataFrame has a `reading` column of type double. Some rows contain Python `None` (missing sensor data) and others contain the result of `0.0/0.0` (an invalid calculation), which Spark stores as `NaN`. You need to identify only the rows with the invalid `NaN` calculation, not the rows with missing sensor data. Which filter should you use?

    1. A.sensor_df.filter(F.isnan("reading"))
    2. B.sensor_df.filter(F.col("reading").isNull())
    3. C.sensor_df.filter(F.col("reading").isNotNull())
    4. D.sensor_df.na.drop(subset=["reading"])
    5. E.sensor_df.filter(F.isnull("reading"))
    Show answer & explanation

    Correct answer: A — sensor_df.filter(F.isnan("reading"))

    • A. isnan() specifically flags Spark's floating-point NaN marker used for invalid numeric results, distinguishing it from a genuinely missing None value.
    • B. isNull() matches rows where the value is Python None/SQL NULL, which is exactly the missing-data case this scenario wants to exclude, not the invalid calculation.
    • C. Filtering with isNotNull() returns both valid numeric readings and NaN rows together, since NaN is a defined floating-point value rather than a null, so it fails to isolate the invalid rows.
    • D. na.drop(subset=["reading"]) removes rows with null values and returns everything else, which keeps the NaN rows mixed in with valid readings rather than isolating them.
    • E. isnull() is functionally the same null check as isNull(), so it identifies missing-data rows rather than the NaN rows produced by the invalid calculation.

    Subdomain 3.3: Perform aggregate operations on DataFrames such as count, approximate count distinct, and mean, summary.

    14.A data engineer wants a single-row report showing the total row count and the average `order_total` for the entire `orders` DataFrame, with no grouping. Which line of code produces this correctly?

    1. A.`orders.select(F.count("*"), F.avg("order_total")).groupBy("order_total")`, which regroups the already-aggregated single row by `order_total`, producing one row per distinct total rather than a single summary row.
    2. B.`orders.groupBy().count().alias("avg_total")`, which returns the row count as a DataFrame and then renames the entire DataFrame object rather than adding an average column.
    3. C.`orders.agg(F.count("*").alias("total_orders"), F.avg("order_total").alias("avg_total"))`, which applies both aggregate expressions across the whole DataFrame and returns one row labeled by the given aliases.
    4. D.`orders.count(), orders.select(F.avg("order_total"))`, which evaluates two independent expressions rather than a single line of code, so the result is a Python tuple, not a DataFrame.
    5. E.`orders.avg("order_total").withColumn("total_orders", F.count("*"))`, which fails because `DataFrame` has no `.avg()` method and `F.count()` cannot be added as a column outside an aggregate context.
    Show answer & explanation

    Correct answer: C — `orders.agg(F.count("*").alias("total_orders"), F.avg("order_total").alias("avg_total"))`, which applies both aggregate expressions across the whole DataFrame and returns one row labeled by the given aliases.

    • A. Once `.select()` has collapsed the data to a single aggregated row, grouping by `order_total` afterward just partitions that one row by its own value, which does not restore a meaningful per-order breakdown or match the requested single summary row.
    • B. `.alias()` on a DataFrame does not rename a column or introduce an average; it is a no-op renaming of the DataFrame reference itself, so no average column is produced.
    • C. This is correct: calling `.agg()` directly on the DataFrame without a preceding `.groupBy()` key aggregates across every row, and the two aliased expressions produce exactly the two labeled columns the report needs in one row.
    • D. Writing two comma-separated expressions on one line creates a tuple of two separate results in Python rather than combining them into a single DataFrame row, so it does not produce the requested report.
    • E. `DataFrame` objects do not expose a bare `.avg()` method, and adding an aggregate function as a new column outside of `.agg()` or `.select()` is not valid, so this raises an error before producing any result.

    Subdomain 3.5: Combine DataFrames with operations such as Inner join, left join, broadcast join, multiple keys, cross join, union, and union all.

    15.```python orders = spark.createDataFrame( [(101, "widget", 3), (102, "gadget", 1), (103, "gizmo", 5)], ["order_id", "item", "qty"] ) promotions = spark.createDataFrame( [(101, "SAVE10"), (103, "SAVE20")], ["order_id", "code"] ) result = orders.join(promotions, on="order_id", how="left") ``` What value appears in the `code` column for order `102` after this left join?

    1. A.An empty string `""`, because Spark fills missing right-side string columns with empty strings rather than nulls.
    2. B.The row for order `102` is removed entirely, because left joins drop rows that lack a match on either side.
    3. C.`null`, because `promotions` has no row with `order_id` = 102, so the left join keeps the order row with a null for the unmatched column.
    4. D.`"SAVE10"`, because Spark falls back to the first available promotion code when no exact match is found.
    5. E.The literal value `"gadget"`, because Spark reuses the closest matching column value when no promotion code exists.
    Show answer & explanation

    Correct answer: C — `null`, because `promotions` has no row with `order_id` = 102, so the left join keeps the order row with a null for the unmatched column.

    • A. Spark does not substitute empty strings for missing values produced by a join; unmatched right-side columns are populated with `null`, not `""`.
    • B. A left join keeps every row from the left DataFrame no matter what, filling right-side columns with null when there is no match, rather than dropping the row.
    • C. This is correct. A left join preserves all rows from `orders`, and since `promotions` has no row for `order_id` 102, the `code` column for that row is populated with `null`.
    • D. Spark does not borrow a value from an unrelated matched row; there is no fallback to reuse `"SAVE10"` for an order that has no corresponding promotion row.
    • E. Spark never copies a value from a different column, such as `item`, into an unrelated output column; a missing match produces a null, not a reused string from elsewhere in the row.

    Subdomain 3.4: Manipulate and utilize Date data type, such as Unix epoch to date string, and extract date component.

    16.A dashboard computes "days since signup" for each user from a `DateType` column `signup_date`. The developer writes: ```python df.withColumn("days_since_signup", datediff(col("signup_date"), current_date())) ``` and the resulting values are negative for every existing user. What should change, and why?

    1. A.Swap the arguments to `datediff(current_date(), col("signup_date"))`; `datediff(end, start)` returns `end` minus `start` in days, so the earlier date first yields a negative count.
    2. B.Replace `datediff()` with `date_add(col("signup_date"), -1)`, because `date_add()` measures elapsed days between two columns while `datediff()` only shifts a single date forward or backward.
    3. C.Cast `signup_date` to `TimestampType` before calling `datediff()`, because `datediff()` returns negative numbers whenever it receives a `DateType` argument instead of a `TimestampType`.
    4. D.Wrap the result in `abs()`, because `datediff()` always returns a negative value whenever its first argument is chronologically earlier than the true current date being compared.
    5. E.Use `months_between(col("signup_date"), current_date())` instead, because `months_between()` returns whole day counts while `datediff()` is limited to whole calendar months.
    Show answer & explanation

    Correct answer: A — Swap the arguments to `datediff(current_date(), col("signup_date"))`; `datediff(end, start)` returns `end` minus `start` in days, so the earlier date first yields a negative count.

    • A. `datediff()` computes its first argument minus its second argument in days, so passing the earlier `signup_date` as the first argument and today as the second produces a negative number; swapping the order fixes the sign.
    • B. `date_add()` shifts a single date by a fixed number of days and does not compute the elapsed days between two columns at all, so it cannot replace `datediff()` for this calculation.
    • C. `datediff()` works correctly with `DateType` arguments and the column's data type is not the source of the negative sign; the cause is the order in which the two dates are passed.
    • D. Wrapping the result in `abs()` hides the sign without fixing the underlying argument order, and the claim that `datediff()` always returns a negative value for an earlier first argument describes exactly the bug rather than a real universal rule to rely on.
    • E. `months_between()` returns a number of months, not days, and `datediff()` returns days, not months, so this option reverses what each function actually measures.

    Subdomain 3.6: Manage input and output operations by writing, overwriting, and reading DataFrames with schemas.

    17.A data engineer reads a daily CSV extract with no header row. The `order_id` column must be read as a long (not the int Spark would infer) to avoid downstream overflow, and schema inference is too slow on the multi-gigabyte file. Which line fills the blank so the file is read with the required types without triggering an inference scan? ```python from pyspark.sql.types import StructType, StructField, LongType, StringType order_schema = StructType([ StructField("order_id", LongType(), True), StructField("customer", StringType(), True) ]) df = spark.read.____.csv("/mnt/raw/orders/") ```

    1. A.option("inferSchema", "true")
    2. B.schema(order_schema)
    3. C.option("header", "true")
    4. D.format("csv").load()
    Show answer & explanation

    Correct answer: B — schema(order_schema)

    • A. Turning on `inferSchema` makes Spark scan the file to guess column types, which is exactly the slow, unreliable path the engineer needs to avoid, and it would not guarantee `order_id` comes back as a long.
    • B. Passing the `StructType` directly forces every row to be parsed according to the declared field names and types, so `order_id` is read as a `LongType` with no inference scan required.
    • C. Declaring a header only tells Spark to treat the first line as column names; it says nothing about column types and does nothing to prevent inference from running.
    • D. Chaining `format` and `load` is just an alternate way to specify the CSV source path and does not attach a schema, so the file would still fall back to default typing behavior.

    Subdomain 3.7: Perform operations on DataFrames such as sorting, iterating, printing schema, and conversion between DataFrame and sequence/list formats.

    18.Given: ```python df = spark.createDataFrame([(1, "a"), (3, "c"), (2, "b")], ["id", "val"]) result = df.orderBy("id", ascending=False).select("id").collect() ``` What are the `id` values in `result`, in order?

    1. A.3, 2, 1
    2. B.1, 2, 3
    3. C.1, 3, 2 (the original insertion order is preserved)
    4. D.A `TypeError` is raised because `ascending` must be a list matching the number of sort columns
    5. E.3, 1, 2
    Show answer & explanation

    Correct answer: A — 3, 2, 1

    • A. `ascending=False` reverses the sort direction for the single `id` column, so rows are ordered from the highest value to the lowest, producing 3, then 2, then 1.
    • B. This would be the result of an ascending sort, but `ascending=False` explicitly requests descending order, so the lowest-to-highest sequence 1, 2, 3 is not what this call produces.
    • C. `orderBy()` always reshuffles rows into the requested sort order; a `DataFrame` has no guaranteed row ordering to preserve in the first place, so the output reflects the sort, not the construction order.
    • D. `ascending` accepts a single boolean when sorting by one column; a list is only required when sorting by multiple columns and specifying a direction per column, so a bare `False` here is valid and does not raise an error.
    • E. This sequence matches neither a full ascending nor a full descending sort on `id`, so it does not correspond to any ordering that `orderBy("id", ascending=False)` would produce.

    Subdomain 3.10: Describe the purpose and implementation of broadcast joins

    19.A Spark job automatically uses a broadcast hash join for a small dimension table with no explicit hint applied anywhere in the code. Which configuration property controls this automatic behavior, and what does its default value represent?

    1. A.`spark.sql.autoBroadcastJoinThreshold` sets the largest estimated table size Spark broadcasts automatically, defaulting to 10 MB unless the session overrides it.
    2. B.Setting `spark.sql.shuffle.partitions` controls how many post-shuffle partitions Spark creates, and it defaults to 200 regardless of any table's estimated size.
    3. C.The property `spark.sql.broadcastTimeout` controls how many seconds Spark waits for a broadcast to complete, defaulting to 300 seconds rather than a size limit.
    4. D.Increasing `spark.driver.maxResultSize` raises the largest table Spark will broadcast automatically, since broadcast tables are collected through the driver's result buffer.
    5. E.Only `spark.sql.adaptive.enabled` determines whether Spark can choose a broadcast join at runtime, with automatic size-based broadcasting disabled until AQE is turned on.
    Show answer & explanation

    Correct answer: A — `spark.sql.autoBroadcastJoinThreshold` sets the largest estimated table size Spark broadcasts automatically, defaulting to 10 MB unless the session overrides it.

    • A. This is correct: `spark.sql.autoBroadcastJoinThreshold` is the size cutoff the planner compares a table's estimated size against, and its default of 10 MB is why small dimension tables get broadcast without any explicit hint. Setting it to -1 disables this automatic behavior entirely.
    • B. `spark.sql.shuffle.partitions` only controls the number of partitions used after a shuffle occurs and has nothing to do with the size threshold that triggers automatic broadcasting. Its default of 200 is unrelated to table size estimation.
    • C. `spark.sql.broadcastTimeout` governs how long a broadcast exchange is allowed to run before timing out, not what size table qualifies for broadcasting in the first place. Its 300-second default is a wait limit, not a size limit.
    • D. `spark.driver.maxResultSize` limits how much data a driver can collect from actions like `collect()`, and it is not the property the planner consults when deciding whether a table is small enough to broadcast. Raising it does not change broadcast eligibility.
    • E. Automatic size-based broadcasting is independent of `spark.sql.adaptive.enabled` and works whether or not adaptive query execution is turned on. AQE can additionally convert a join to broadcast at runtime, but it is not a prerequisite for the size-based default behavior.

    Subdomain 3.9: Describe different types of variables in Spark, including broadcast variables and accumulators.

    20.A data engineer has a small Python dictionary `country_codes` (well under 1 MB) that maps ISO codes to country names, and needs to look it up inside a PySpark UDF applied to a 500 million row DataFrame. To avoid Spark serializing and shipping a fresh copy of `country_codes` with every task, which change should the engineer make to the following code? ```python def lookup_country(code): return country_codes.get(code, "Unknown") lookup_udf = udf(lookup_country, StringType()) df = df.withColumn("country_name", lookup_udf(col("iso_code"))) ```

    1. A.Wrap `country_codes` with `spark.sparkContext.broadcast(country_codes)` and reference `bc.value.get(code, "Unknown")`, so every executor keeps one cached copy.
    2. B.Move the dictionary definition inside `lookup_country` itself, rebuilding `country_codes` on every call so each task always works with its own freshly built local copy of the mapping.
    3. C.Convert `country_codes` into a Spark DataFrame and perform a `join` against the 500 million row DataFrame on the ISO code column before applying any further transformations.
    4. D.Register `country_codes` as a Spark accumulator with `sc.accumulator(country_codes)` so tasks can read and update the shared mapping while the batch job keeps running.
    5. E.Raise `spark.sql.autoBroadcastJoinThreshold` so Spark automatically distributes the dictionary as part of query planning for this specific UDF call in the job.
    Show answer & explanation

    Correct answer: A — Wrap `country_codes` with `spark.sparkContext.broadcast(country_codes)` and reference `bc.value.get(code, "Unknown")`, so every executor keeps one cached copy.

    • A. Broadcasting caches a single read-only copy of the dictionary on each executor, and referencing `.value` inside the UDF avoids re-serializing and re-sending the dictionary with every task.
    • B. Rebuilding the dictionary inside the function makes the problem worse: Spark still ships whatever data built it, and the work of constructing the mapping now repeats on every single row.
    • C. A join is a valid pattern for larger reference tables, but for a small in-memory dictionary it adds an unnecessary shuffle and does not fix the serialization behavior of the existing UDF code.
    • D. Accumulators are write-only aggregators whose running total only the driver can read; they are not a mechanism for distributing an arbitrary read-only mapping for executors to look values up in.
    • E. `spark.sql.autoBroadcastJoinThreshold` only influences whether Catalyst chooses a broadcast join strategy for DataFrame joins, and has no effect on a plain Python closure referenced inside a UDF.

    Subdomain 3.8: Create and invoke user-defined functions with or without stateful operators, including StateStores.

    21.Which `timeoutConf` value should be passed to `applyInPandasWithState` so a group's state expires after a fixed amount of wall-clock time has elapsed since it last received data, independent of any event-time watermark in the incoming records?

    1. A.GroupStateTimeout.NoTimeout()
    2. B.GroupStateTimeout.EventTimeTimeout()
    3. C.GroupStateTimeout.ProcessingTimeTimeout()
    4. D.GroupStateTimeout.WatermarkTimeout()
    Show answer & explanation

    Correct answer: C — GroupStateTimeout.ProcessingTimeTimeout()

    • A. This disables automatic expiration entirely, so state for a group is retained indefinitely until the per-group function itself calls `state.remove()`, which does not satisfy a fixed wall-clock expiration requirement.
    • B. This ties expiration to the watermark computed from an event-time column in the data, so how quickly state expires depends on data arrival patterns rather than a fixed span of processing time.
    • C. This measures elapsed wall-clock time against the duration set with `state.setTimeoutDuration()`, independent of any event-time column or watermark, which is exactly the behavior described.
    • D. `GroupStateTimeout` has no member with this name; the only valid timeout configurations are `NoTimeout()`, `ProcessingTimeTimeout()`, and `EventTimeTimeout()`.

    Subdomain 3.8: Create and invoke user-defined functions with or without stateful operators, including StateStores.

    22.A team reviews an `applyInPandasWithState` implementation that declares `stateStructType = StructType([StructField("count", LongType()), StructField("last_seen", TimestampType())])` and passes `GroupStateTimeout.NoTimeout()` as the `timeoutConf`. Which two statements about this setup are correct? (Choose 2 answers)(Select 2)

    1. A.`state.setTimeoutDuration()` can still be called safely to expire idle keys even though `timeoutConf` is `GroupStateTimeout.NoTimeout()`, since the timeout configuration only affects `EventTimeTimeout()`.
    2. B.State declared with `stateStructType` is scoped to the entire streaming query rather than per grouping key, so two different `device_id` values processed by the same function share one combined state object.
    3. C.Every call to `state.update()` inside the per-group function must pass a two-element tuple whose values line up with `count` and `last_seen`, since the value is validated against `stateStructType`.
    4. D.Because `timeoutConf` is `GroupStateTimeout.NoTimeout()`, `state.hasTimedOut` will always evaluate to `False` inside the function, so no branch of the function will ever run purely due to a timeout.
    Show answer & explanation

    Correct answers: C, D — Every call to `state.update()` inside the per-group function must pass a two-element tuple whose values line up with `count` and `last_seen`, since the value is validated against `stateStructType`.; Because `timeoutConf` is `GroupStateTimeout.NoTimeout()`, `state.hasTimedOut` will always evaluate to `False` inside the function, so no branch of the function will ever run purely due to a timeout.

    • A. Calling `state.setTimeoutDuration()` while `timeoutConf` is `NoTimeout()` raises an error at runtime, since duration-based timeouts are only valid when `ProcessingTimeTimeout()` is the configured timeout type.
    • B. State declared with `stateStructType` is maintained independently per grouping key produced by `groupBy()`; each distinct key gets its own isolated state entry, not one shared object across keys.
    • C. `state.update()` validates its argument against the declared `stateStructType`, so passing a tuple of the wrong arity or incompatible types raises an error; here that means a two-element tuple matching `count` and `last_seen`.
    • D. Timeout callbacks only fire when `timeoutConf` enables them; with `NoTimeout()` selected, Spark never invokes the function purely because time has elapsed, so this property stays `False` on every call.

    Domain 4: Troubleshooting and Tuning Apache Spark DataFrame API Applications

    Subdomain 4.1: Implement performance tuning strategies & optimize cluster utilization, including partitioning, repartitioning, coalescing, identifying data skew, and reducing shuffling

    23.A DataFrame with 1000 partitions is filtered down to a small result set. An engineer wants the final write to produce exactly 8 output files while spending as little execution time as possible on the operation itself: ```python result = spark.read.table("sales.transactions").filter(col("region") == "EMEA") result.write.mode("overwrite").parquet("/mnt/output/emea_sales") ``` Which change achieves this most efficiently?

    1. A.Insert `result = result.coalesce(8)` before the write, since it merges existing partitions on the same executors and avoids a full shuffle across the cluster.
    2. B.Insert `result = result.repartition(8)` before the write, since renumbering partition metadata this way never triggers a network shuffle of the underlying rows.
    3. C.Set `spark.conf.set("spark.sql.shuffle.partitions", 8)` right before the write call, since that setting fixes the output partition count for any write operation.
    4. D.Insert `result = result.repartition(8, "region")` before the write, since hashing on the filter column guarantees exactly 8 evenly sized output files here.
    5. E.Call `result.persist()` immediately before the write, since caching a DataFrame in memory automatically collapses its partition count to match the write target.
    Show answer & explanation

    Correct answer: A — Insert `result = result.coalesce(8)` before the write, since it merges existing partitions on the same executors and avoids a full shuffle across the cluster.

    • A. `coalesce()` reduces the number of partitions by combining existing ones on the same executors, so it can shrink partition count without a full shuffle when the target is fewer partitions than the current count. This makes it the cheapest way to hit exactly 8 output files here.
    • B. `repartition()` always performs a full shuffle to redistribute rows across the new partition count, regardless of direction, so it is more expensive than needed when only reducing the number of partitions.
    • C. `spark.sql.shuffle.partitions` only controls the number of partitions produced by shuffle operations such as joins and aggregations; it has no effect on the partition count of a plain filter-then-write DataFrame operation.
    • D. Hash partitioning on a single low-cardinality column like `region` does not guarantee balanced output; rows with the same region hash to the same partition, so file sizes can still be very uneven, and this still forces a full shuffle.
    • E. Caching a DataFrame stores its existing partitions in memory or on disk; it does not change the partition count at all, so the write would still emit the original number of files.

    Subdomain 4.2: Describe Adaptive Query Execution (AQE) and its benefits.

    24.A data engineer wants Spark to pick the number of shuffle partitions automatically for each job rather than always using a fixed count, letting AQE's coalescing feature size the partitions itself. Which setting achieves this on a Databricks cluster with AQE enabled?

    1. A.`spark.conf.set("spark.sql.shuffle.partitions", "auto")` starts each shuffle with a large partition count and lets AQE coalesce it down toward the advisory partition size.
    2. B.`spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "auto")` switches partition coalescing between on and off automatically depending on the size of the input DataFrame.
    3. C.`spark.conf.set("spark.sql.shuffle.partitions", "0")` tells Spark to compute the shuffle partition count from the number of executor cores available at job submission time.
    4. D.`spark.conf.set("spark.databricks.adaptive.autoOptimizeShuffle.enabled", "false")` forces Spark back onto a manually specified partition count for every shuffle stage in the job.
    5. E.`spark.conf.set("spark.sql.files.maxPartitionBytes", "auto")` lets Spark size the input read partitions dynamically, which also determines the partition count for every downstream shuffle.
    Show answer & explanation

    Correct answer: A — `spark.conf.set("spark.sql.shuffle.partitions", "auto")` starts each shuffle with a large partition count and lets AQE coalesce it down toward the advisory partition size.

    • A. This is correct because setting the shuffle partition count to `"auto"` on Databricks starts with a large initial partition count and delegates the final sizing to AQE's coalescing logic based on actual shuffled data volume.
    • B. This is incorrect because the coalescing enabled flag is a plain boolean; it does not accept an `"auto"` value that conditionally toggles behavior by input size.
    • C. This is incorrect because setting the partition count to zero is not a supported way to derive partitions from executor core count; it does not enable automatic AQE-driven sizing.
    • D. This is incorrect because disabling this Databricks-specific auto-optimization flag does the opposite of what is requested, forcing a fixed manual partition count instead of automatic sizing.
    • E. This is incorrect because this setting controls input file read partitioning, not the number of partitions used for shuffles triggered later by joins or aggregations in the query.

    Subdomain 4.3: Perform logging and monitoring of Spark applications - publish, customize, and analyze Driver logs and Executor logs to diagnose out-of-memory errors, cluster underutilization, etc.

    25.A join between a 2 TB fact table and what the engineer believed was a small dimension table causes repeated executor `OutOfMemoryError` failures: ```python fact = spark.read.format("delta").load("/mnt/gold/fact_sales") dim = spark.read.format("delta").load("/mnt/gold/dim_store") result = fact.join(broadcast(dim), "store_id") ``` Investigation shows `dim_store` has grown to 12 GB over time and `spark.sql.autoBroadcastJoinThreshold` is set to `-1` in the cluster's Spark config. Which two actions would most directly resolve the OutOfMemoryError? (Choose 2 answers.)(Select 2)

    1. A.Remove the explicit `broadcast(dim)` hint so the join falls back to a standard sort-merge join, since forcing a 12 GB table to be broadcast to every executor is what is overwhelming executor memory regardless of the threshold setting.
    2. B.Reduce `dim_store` to only the columns and rows the join actually needs before the join executes, and confirm the filtered DataFrame is genuinely small so the existing `broadcast(dim)` hint can safely execute without overwhelming any executor.
    3. C.Increase `spark.sql.shuffle.partitions` to a much higher value so the sort-merge join that `broadcast(dim)` supposedly triggers under the hood can process the 12 GB dimension table across more, smaller shuffle partitions instead.
    4. D.Persist `dim_store` with `StorageLevel.DISK_ONLY` before the join, since caching a DataFrame to disk exempts it from the broadcast join threshold check and allows the explicit `broadcast()` hint to succeed without exceeding executor memory.
    5. E.Rewrite the join as `fact.crossJoin(dim).filter(fact.store_id == dim.store_id)`, since a cross join followed by a filter processes the two tables independently on each executor and avoids ever broadcasting `dim_store` in full.
    Show answer & explanation

    Correct answers: A, B — Remove the explicit `broadcast(dim)` hint so the join falls back to a standard sort-merge join, since forcing a 12 GB table to be broadcast to every executor is what is overwhelming executor memory regardless of the threshold setting.; Reduce `dim_store` to only the columns and rows the join actually needs before the join executes, and confirm the filtered DataFrame is genuinely small so the existing `broadcast(dim)` hint can safely execute without overwhelming any executor.

    • A. An explicit `broadcast()` hint forces Spark to ship the entire hinted DataFrame to every executor regardless of the `autoBroadcastJoinThreshold` setting, so removing the hint lets the optimizer fall back to a standard sort-merge join that processes `dim_store` in shuffled partitions instead of copying all 12 GB into every executor's memory.
    • B. The underlying problem is that `dim_store` grew past the point where broadcasting it is safe, so trimming it down to only the join-relevant columns and rows — and verifying the result is actually small — restores the assumption the original code was written under and lets the broadcast hint execute without exhausting executor memory.
    • C. `spark.sql.shuffle.partitions` controls partitioning for shuffle-based joins such as sort-merge joins, but an explicit `broadcast()` hint bypasses the shuffle path entirely and copies the whole DataFrame to every executor, so raising this setting has no effect on the broadcast operation that is actually failing.
    • D. Setting a storage level on `dim_store` controls how that DataFrame's data is cached for reuse, but it has no relationship to the broadcast join threshold or to how the explicit `broadcast()` hint decides to ship data to executors, so this change would not prevent the same broadcast operation from failing.
    • E. A cross join followed by a filter still requires materializing the full pairwise combination of both tables' partitions on the executors that process them, which is typically far more expensive than either a sort-merge join or a safely sized broadcast, and does not avoid memory pressure — it usually makes it worse.

    Domain 5: Structured Streaming

    Subdomain 5.1: Explain the Structured Streaming engine in Spark, including its functions, programming model, micro-batch processing, exactly-once semantics, and fault tolerance mechanisms.

    26.A data engineer starts the following streaming query without specifying a `.trigger(...)` clause: ```python query = (streamingDF.writeStream .format("console") .outputMode("append") .start()) ``` Which best describes how Spark schedules micro-batches for this query?

    1. A.Spark uses the default micro-batch trigger, immediately starting the next micro-batch as soon as the previous one finishes processing
    2. B.Spark falls back to continuous processing mode automatically, executing records with millisecond-level latency instead of batching them
    3. C.Spark waits indefinitely for an explicit trigger interval to be configured, and the query remains idle until one is set
    4. D.Spark processes exactly one micro-batch containing all currently available data, then stops the query automatically after that batch
    Show answer & explanation

    Correct answer: A — Spark uses the default micro-batch trigger, immediately starting the next micro-batch as soon as the previous one finishes processing

    • A. This is correct: when no trigger is specified, Spark uses the default micro-batch trigger, which starts a new batch as soon as the prior one finishes, giving the lowest latency the micro-batch engine can offer without explicit tuning.
    • B. This is incorrect because continuous processing is an experimental, opt-in execution mode requiring `.trigger(continuous=...)`; omitting a trigger does not switch the engine into that mode.
    • C. This is incorrect because the absence of a trigger clause does not pause the query; Spark falls back to its default trigger behavior rather than waiting for configuration.
    • D. This is incorrect because processing a single batch and stopping is the behavior of a one-time trigger such as `Trigger.Once()`, not the default trigger that is used when no trigger is specified.

    Subdomain 5.1: Explain the Structured Streaming engine in Spark, including its functions, programming model, micro-batch processing, exactly-once semantics, and fault tolerance mechanisms.

    27.Which statement accurately distinguishes Spark's default micro-batch processing model from its continuous processing mode?

    1. A.Micro-batch mode polls the source and batches accumulated data, while continuous processing runs long tasks that handle records as they arrive for lower latency
    2. B.Micro-batch mode requires a fixed cluster size that cannot ever autoscale, while continuous processing dynamically adds new executors for every single incoming record
    3. C.Micro-batch mode only supports file-based sinks, while continuous processing is required whenever the sink is a Delta table or a JDBC connection
    4. D.Micro-batch mode guarantees at-most-once delivery semantics, while continuous processing is the only Spark mode capable of exactly-once guarantees
    Show answer & explanation

    Correct answer: A — Micro-batch mode polls the source and batches accumulated data, while continuous processing runs long tasks that handle records as they arrive for lower latency

    • A. This is correct: micro-batch execution repeatedly polls the source and processes whatever has accumulated in discrete batches, while continuous processing keeps long-running tasks reading and processing records individually to achieve much lower latency.
    • B. This is incorrect because cluster autoscaling behavior is independent of which streaming execution mode is chosen; neither mode inherently fixes or dynamically resizes the cluster per record.
    • C. This is incorrect because micro-batch mode supports a wide range of sinks including Delta tables and JDBC, not only files; sink support is not the distinguishing factor between the two modes.
    • D. This is incorrect because micro-batch mode is the execution model that provides Structured Streaming's well-supported exactly-once guarantees, while continuous processing offers weaker, at-least-once guarantees for its supported operations.

    Subdomain 5.3: Perform basic operations on Streaming DataFrames and Streaming Datasets, such as selection, projection, window and aggregation.

    28.A team processes streaming user activity events and wants to group consecutive events from the same `userId` into a single session that closes automatically whenever there is at least a 5-minute gap of inactivity, rather than using a fixed calendar window. Which construct should they use?

    1. A.`session_window()` builds one window per group that keeps extending while new events arrive within the gap duration, and finalizes once that inactivity gap has elapsed.
    2. B.`window()` with a 5-minute duration and no slide duration approximates session grouping, since a short tumbling window mirrors typical gaps between a user's events.
    3. C.`groupBy("userId").agg(max("event_time"))` combined with a watermark tracks each user's latest event, which is treated as equivalent to a session boundary in practice.
    4. D.`window()` with a 5-minute slide and a 10-minute duration merges adjacent events through overlapping windows, approximating the session boundaries the team wants.
    5. E.`dropDuplicatesWithinWatermark("userId")` applied before the aggregation removes repeat events per user, which approximates grouping a user's activity into one session.
    Show answer & explanation

    Correct answer: A — `session_window()` builds one window per group that keeps extending while new events arrive within the gap duration, and finalizes once that inactivity gap has elapsed.

    • A. This is correct because `session_window` is purpose-built for dynamic-length windows keyed on a gap duration: a session stays open while new events keep arriving within that gap and is finalized once the gap is exceeded, matching the requirement exactly.
    • B. This is incorrect because a tumbling window has a fixed boundary aligned to the clock regardless of activity, so it would split or merge events based on wall-clock time rather than actual inactivity gaps.
    • C. This is incorrect because tracking only the maximum event time per user produces a single running aggregate value, not distinct session windows that open and close based on gaps between events.
    • D. This is incorrect because sliding windows still have fixed durations and slide intervals tied to the clock, so they do not dynamically expand or contract based on the actual gap between a user's events.
    • E. This is incorrect because `dropDuplicatesWithinWatermark` removes duplicate rows within a watermark bound; it does not group events into variable-length sessions based on inactivity gaps.

    Subdomain 5.2: Create and write Streaming DataFrames and Streaming Datasets, including the basic output modes and output sinks.

    29.A developer joins two streaming DataFrames and writes the result: ```python joined = orders_stream.join(shipments_stream, "order_id") query = joined.writeStream.outputMode(______).format("console").start() ``` Which two output modes are **not** valid for this streaming-to-streaming join and will raise an `AnalysisException` when passed to `.outputMode()`? (Choose 2 answers)(Select 2)

    1. A.`append`
    2. B.`update`
    3. C.`complete`
    4. D.`once`
    5. E.`default`
    Show answer & explanation

    Correct answers: B, C — `update`; `complete`

    • A. Stream-to-stream joins are only supported with the `append` output mode, so this mode is valid and does not raise an exception.
    • B. `update` mode is not yet supported for stream-to-stream joins; passing it to `.outputMode()` for a join query raises an `AnalysisException`.
    • C. `complete` mode is not supported for stream-to-stream joins either, since the join does not produce a bounded result table that can be rewritten in full each trigger.
    • D. `once` is a trigger setting, not an output mode string, so it is not a candidate answer to this question about `.outputMode()` values.
    • E. `default` is not a recognized output mode identifier in Structured Streaming; the valid mode names are `append`, `update`, and `complete`.

    Subdomain 5.4: Perform Streaming Deduplication in Structured Streaming, both with and without watermark usage.

    30.A sensor telemetry stream has columns `device_id`, `reading_time`, and `temperature`. The team writes: ```python cleaned = sensor_df \ .withWatermark("reading_time", "15 minutes") \ .dropDuplicates(["device_id", "reading_time"]) ``` A record for `device_id="A1"` with `reading_time="10:00:00"` arrives, is processed, and the watermark subsequently advances to `10:20:00`. A duplicate record for the same `device_id` and `reading_time="10:00:00"` then arrives. What happens to this late duplicate?

    1. A.The state for that key was already evicted once the watermark passed 10:00:00 by more than the 15-minute delay, so the engine can no longer recognize the duplicate and the record passes through as new.
    2. B.The engine always retains a permanent record of every key it has ever emitted, kept in a separate index outside the bounded state store, so the duplicate is correctly filtered out no matter how late it arrives.
    3. C.The query throws a runtime exception because the incoming event time is now older than the current watermark, which Structured Streaming treats as an unrecoverable error in the input data.
    4. D.The engine automatically extends the watermark delay backward to accommodate the late arrival, silently widening the effective lateness bound to keep the original key's state alive longer.
    5. E.The record is buffered by the engine until the next checkpoint interval, at which point it is compared against a stored snapshot of every key deleted so far before being emitted.
    Show answer & explanation

    Correct answer: A — The state for that key was already evicted once the watermark passed 10:00:00 by more than the 15-minute delay, so the engine can no longer recognize the duplicate and the record passes through as new.

    • A. This is correct: once the watermark advances past a key's event time by more than the declared delay, the engine treats that key's state as safe to remove, so a genuine duplicate arriving after that point is no longer matched and is emitted as a new row. This is the expected tradeoff of bounding state with a watermark.
    • B. There is no secondary unbounded index kept outside the state store; the watermark's entire purpose is to let the engine physically discard state, and once discarded there is no other record of the key to check against.
    • C. Structured Streaming does not throw an exception for late data relative to the watermark; late records that fall outside the bound are simply treated as new input rather than causing the query to fail.
    • D. The watermark delay is a fixed configuration set when the query is defined; the engine does not dynamically extend it backward in response to individual late arrivals, since that would defeat the purpose of bounding state at all.
    • E. There is no buffering-until-checkpoint mechanism that reconciles late records against deleted keys; the state for expired keys is gone, and buffering indefinitely would reintroduce the unbounded growth the watermark is meant to prevent.

    Subdomain 5.4: Perform Streaming Deduplication in Structured Streaming, both with and without watermark usage.

    31.Which statement about setting the watermark delay threshold for streaming deduplication is accurate?

    1. A.The threshold should exceed the maximum expected gap between an event and its latest duplicate, so genuine duplicates are not mistaken for new records once the watermark passes.
    2. B.The threshold has no effect on correctness at all, and only changes how frequently the Spark UI refreshes the numbers shown on its streaming query statistics panel.
    3. C.The threshold must always be set to exactly equal the micro-batch trigger interval, since Spark derives all watermark expiry decisions directly from the configured trigger duration.
    4. D.The threshold is measured in a count of records rather than a time duration, so it should be set to the expected number of duplicate retries observed per key.
    5. E.The threshold only affects the aggregation operators present in a query, and has no bearing on how streaming deduplication state is retained or expired over time.
    Show answer & explanation

    Correct answer: A — The threshold should exceed the maximum expected gap between an event and its latest duplicate, so genuine duplicates are not mistaken for new records once the watermark passes.

    • A. This is correct: the watermark delay directly determines how long state is retained for a given key, so it must cover the realistic lateness of duplicate arrivals or the engine will expire state too early and fail to catch genuine duplicates.
    • B. The watermark threshold directly controls state retention and correctness of duplicate detection, not merely how often UI statistics are refreshed, which is an unrelated concern from query configuration.
    • C. The watermark delay is an independently configured duration passed to `withWatermark` and is not derived from or required to match the micro-batch trigger interval, which controls how often batches are processed.
    • D. The watermark threshold is expressed as a time duration relative to event time, such as minutes or hours, not as a count of expected retries per key.
    • E. The watermark threshold governs state retention for any stateful operator that uses it, including deduplication, not only aggregation operators; deduplication state expiry depends directly on this same threshold.

    Domain 6: Using Spark Connect to deploy applications

    Subdomain 6.1: Describe the features of Spark Connect.

    32.A platform team upgrades the Spark version running on their Spark Connect server to pick up performance improvements, while client applications keep running unmodified PySpark scripts such as: ```python df = spark.read.table("events").groupBy("user_id").count() ``` Which two statements correctly explain why this server upgrade does not require the client applications to be rebuilt or redeployed? (Choose 2 answers)(Select 2)

    1. A.The thin client library only encodes DataFrame operations into unresolved logical plans; it does not depend on the internal execution engine version running on the server.
    2. B.Spark Connect's gRPC protocol maintains compatibility across versions, so an older client library can keep issuing requests to a newer server without a matching client-side upgrade.
    3. C.The client and server exchange a shared JVM classpath at connection time, so any server-side JAR upgrade is automatically mirrored onto the client machine.
    4. D.The server converts the client's compiled Java bytecode into an updated intermediate representation, which removes any dependency on the original Spark version.
    5. E.Client applications never actually run against the upgraded server, because Spark Connect automatically pins each client to the exact server version it was originally built against.
    Show answer & explanation

    Correct answers: A, B — The thin client library only encodes DataFrame operations into unresolved logical plans; it does not depend on the internal execution engine version running on the server.; Spark Connect's gRPC protocol maintains compatibility across versions, so an older client library can keep issuing requests to a newer server without a matching client-side upgrade.

    • A. Correct — because the client only builds and transmits an unresolved logical plan, everything version-specific about execution stays on the server side, so a server upgrade does not change what the client needs to send.
    • B. Correct — Spark Connect is designed so the gRPC-based protocol between client and server stays compatible across versions, letting an existing client keep working against a newly upgraded server without needing its own rebuild.
    • C. Incorrect — there is no shared JVM classpath exchanged between client and server; Spark Connect deliberately avoids that kind of tight coupling, which is why client and server can be upgraded independently.
    • D. Incorrect — clients submit DataFrame operations as plans, not compiled Java bytecode, so there is no bytecode-to-intermediate-representation conversion happening as part of a server upgrade.
    • E. Incorrect — Spark Connect does not pin clients to a fixed server version; the whole point of protocol compatibility is that a client can keep talking to a server that has since been upgraded.

    Subdomain 6.2: Describe the different deployment mode types (Client, Cluster, Local) in the Apache Spark™ environment.

    33.A data scientist installs Databricks Connect in a local IDE and writes PySpark DataFrame code that references a table in Unity Catalog, then runs the script. Which description matches how Databricks Connect executes this code?

    1. A.Databricks Connect acts as a Spark Connect client: DataFrame operations build a logical plan locally and send it over gRPC to the remote cluster for execution.
    2. B.Databricks Connect copies the entire Databricks Runtime onto the local IDE's machine, so all DataFrame operations execute locally before results sync back to Unity Catalog.
    3. C.Databricks Connect requires the script to be uploaded as a notebook first, since local IDEs cannot submit Spark Connect requests directly to a remote cluster.
    4. D.Databricks Connect runs the driver locally and the executors locally as well, only contacting Unity Catalog at the very end to write the final output table to storage.
    5. E.Databricks Connect converts every DataFrame call into a SQL warehouse query behind the scenes, bypassing the Spark execution engine on the remote cluster entirely.
    Show answer & explanation

    Correct answer: A — Databricks Connect acts as a Spark Connect client: DataFrame operations build a logical plan locally and send it over gRPC to the remote cluster for execution.

    • A. Databricks Connect is built as a Spark Connect client, so DataFrame code written locally builds a logical plan that is sent over gRPC to the remote Databricks cluster, which executes it with access to Unity Catalog. This matches the intended workflow.
    • B. Databricks Connect does not install or copy the Databricks Runtime onto the local machine; the whole point of the thin client design is that execution happens remotely, not that the runtime is duplicated locally.
    • C. One of the main benefits of Databricks Connect is that it lets a local IDE submit Spark Connect requests directly to a remote cluster without ever converting the script into a notebook first.
    • D. Running both the driver and executors locally would make Databricks Connect indistinguishable from local mode, but its purpose is remote execution against a real Databricks cluster and its Unity Catalog access, not local computation.
    • E. Databricks Connect submits work through the Spark execution engine via Spark Connect, not by rewriting DataFrame operations into SQL warehouse queries; the underlying engine and plan structure remain Spark's own.

    Domain 7: Using Pandas API on Spark

    Subdomain 7.1: Explain the advantages of using Pandas API on Spark.

    34.Which statement accurately describes the relationship between `pyspark.pandas` and the native `pyspark.sql.DataFrame` API?

    1. A.`pyspark.pandas` is a pandas-compatible layer built on top of Spark DataFrames, so operations written with its API are translated into the same distributed Spark execution plan that native PySpark DataFrames use.
    2. B.`pyspark.pandas` is a completely separate execution engine from Spark SQL, so DataFrames created with it run on a dedicated single-node process that exists outside of the Spark cluster entirely.
    3. C.`pyspark.pandas` replaces `pyspark.sql.DataFrame` starting from Spark 3.2, meaning the native DataFrame API is deprecated and all new PySpark code should use pandas-on-Spark DataFrames exclusively from that version onward.
    4. D.`pyspark.pandas` only works when Spark is running in local mode, so any pandas-on-Spark code submitted to a multi-node cluster automatically falls back to a native PySpark DataFrame instead.
    5. E.`pyspark.pandas` is a client-side charting library that renders pandas-style plots from query results, and it does not provide any DataFrame transformation or aggregation methods of its own.
    Show answer & explanation

    Correct answer: A — `pyspark.pandas` is a pandas-compatible layer built on top of Spark DataFrames, so operations written with its API are translated into the same distributed Spark execution plan that native PySpark DataFrames use.

    • A. This is correct: `pyspark.pandas` provides pandas-compatible syntax that is executed through the same underlying Spark DataFrame execution engine, so it is a layer on top of, not a replacement for, native PySpark DataFrames.
    • B. This is incorrect because `pyspark.pandas` runs within the Spark cluster and relies on the same distributed engine as native DataFrames; it is not a separate single-node execution engine outside the cluster.
    • C. This is incorrect because the native `pyspark.sql.DataFrame` API remains fully supported and is not deprecated; `pyspark.pandas` is an additional, complementary interface rather than a replacement.
    • D. This is incorrect because the pandas API on Spark is designed to run on distributed multi-node clusters just like native PySpark, and there is no local-mode-only restriction or automatic fallback behavior.
    • E. This is incorrect because `pyspark.pandas` provides a full set of DataFrame transformation, aggregation, and I/O methods mirroring pandas; plotting is only one small part of its functionality, not its entire purpose.

    Subdomain 7.2: Create and invoke Pandas UDF.

    35.A pipeline needs to filter out rows where `amount` is negative and add a `processed_at` timestamp column, operating on arbitrary batches of the whole DataFrame rather than per group. Which call correctly uses `mapInPandas` to do this?

    1. A.```python def clean_batch(iterator): for pdf in iterator: pdf = pdf[pdf["amount"] >= 0].copy() pdf["processed_at"] = pd.Timestamp.now() yield pdf df.mapInPandas(clean_batch, schema=new_schema) ```, applied directly on the ungrouped DataFrame.
    2. B.The same `clean_batch` function, but called as `df.groupby("amount").mapInPandas(clean_batch, schema=new_schema)` so each batch is scoped to one distinct `amount` value.
    3. C.```python @pandas_udf(new_schema) def clean_batch(pdf: pd.DataFrame) -> pd.DataFrame: pdf = pdf[pdf["amount"] >= 0].copy() pdf["processed_at"] = pd.Timestamp.now() return pdf df.mapInPandas(clean_batch) ```, decorating the function with `pandas_udf` before passing it to `mapInPandas`.
    4. D.```python def clean_batch(pdf): pdf = pdf[pdf["amount"] >= 0].copy() pdf["processed_at"] = pd.Timestamp.now() return pdf df.mapInPandas(clean_batch, schema=new_schema) ```, where the function takes and returns a single pandas DataFrame instead of an iterator.
    5. E.```python def clean_batch(iterator): for pdf in iterator: pdf = pdf[pdf["amount"] >= 0].copy() pdf["processed_at"] = pd.Timestamp.now() yield pdf df.mapInPandas(clean_batch) ```, omitting the `schema` argument since mapInPandas can infer it from the yielded DataFrame.
    Show answer & explanation

    Correct answer: A — ```python def clean_batch(iterator): for pdf in iterator: pdf = pdf[pdf["amount"] >= 0].copy() pdf["processed_at"] = pd.Timestamp.now() yield pdf df.mapInPandas(clean_batch, schema=new_schema) ```, applied directly on the ungrouped DataFrame.

    • A. mapInPandas takes a function that accepts and returns an iterator of pandas DataFrames representing arbitrary batches, and requires an explicit schema argument describing the output; calling it directly on the ungrouped DataFrame matches exactly how this batch-level filtering and column addition is meant to be expressed.
    • B. mapInPandas processes arbitrary batches of the whole DataFrame and is not a GroupedData method, so chaining it after groupby is not a supported call pattern for this operation.
    • C. mapInPandas expects a plain Python function with an iterator-to-iterator signature, not a function decorated with pandas_udf; wrapping the function this way and calling mapInPandas without a schema argument does not match its required calling convention.
    • D. mapInPandas requires the function to accept and yield an iterator of pandas DataFrames, not a single DataFrame passed and returned directly, so this signature does not match what mapInPandas expects to call.
    • E. mapInPandas requires the schema to be passed explicitly as an argument because Spark cannot infer the output schema from Python code before execution; omitting it raises an error rather than being inferred automatically.

    Want the full experience?

    These are just samples. Practice the full Databricks Certified Associate Developer for Apache Spark question bank in quiz mode — free, no signup, with domain practice and exam simulation.