When is dropping from the DataFrame API to Spark's RDD API still justified?
answer
- assume the answer is no
- some data has no columns yet
- some operators depend on position
- df.rdd is a materialization, not a view
basics
~20 sRarely, and only for work the DataFrame API cannot express: non-tabular or custom-parsed data, explicit control of physical partition placement, position-based operators like zipWithIndex, or interop with RDD-based libraries. Everything expressible as columns should stay in DataFrames.
solid answer
~40 sThe default is DataFrames, because calling `df.rdd` opts out of the optimizer, code generation and the binary row format in one step. Legitimate exceptions exist: data that is not tabular at all (raw binary, irregular text you must parse yourself), needing an explicit custom partitioner so repeated key-based joins stay co-partitioned, position-dependent operators such as `zipWithIndex` or `zipPartitions` with no DataFrame equivalent, and interop with RDD-based libraries such as GraphX. Even then, drop down for the **narrowest possible step**, not the whole pipeline: keep the source read, the joins and the aggregations declarative. In PySpark the RDD escape hatch is usually the wrong tool anyway — `mapInPandas` and `applyInPandas` give you per-batch custom Python with Arrow transfer instead of per-row pickling, and Scala users have `Dataset.mapPartitions` for the same purpose.
code
python · 5 lines# Usually the wrong instinct: leaves the optimizer behind for the whole rest of the job
result = df.rdd.mapPartitions(parse_and_score).toDF(schema)
# Better: keep the plan declarative, do custom Python per Arrow batch
result = df.mapInPandas(parse_and_score_batches, schema=out_schema)go deeper
Know that the DataFrame API is the default choice and that RDDs are the older, lower-level abstraction you would only reach for in unusual cases.
Name concrete cases the structured API cannot express and explain why df.rdd costs you the optimizer, code generation and the binary row format.
Show containment in practice: the narrowest possible drop-down, an explicit schema on the way back, and knowing that mapInPandas or mapPartitions usually removes the need entirely.
Set the policy and the migration order for a legacy estate, judging where a rewrite pays for itself against where opaque custom logic can stay isolated indefinitely.
## Start from the presumption against it RDDs are not deprecated and are not going away — every DataFrame query still executes as tasks over partitions underneath. But writing RDD code today means giving up, at once: predicate pushdown and column pruning, operator reordering, whole-stage code generation, and the compact binary row layout. In PySpark it additionally means shipping every row into a Python worker process. So the burden of proof sits with the RDD, and "I find the functional API more readable" does not discharge it. ## Cases that genuinely justify it **Data that is not tabular.** Raw binary records, protocol dumps, irregular log formats where your parser is the first step. There is no schema to declare until after your code runs. `sc.binaryFiles` or `wholeTextFiles` plus `mapPartitions` is a reasonable front end; convert to a DataFrame the moment rows have a shape. **Explicit control of physical partitioning.** The RDD API lets you attach a custom `Partitioner` to a key-value RDD and keep that partitioning across operations, so repeated joins on the same key avoid a shuffle. The DataFrame API expresses partitioning through `repartition`, bucketing and the planner's own choices; when you have a very specific co-partitioning scheme the planner will not reproduce, the RDD API is the direct route. **Position-based and multi-dataset operators.** `zipWithIndex`, `zipWithUniqueId` and `zipPartitions` have no direct DataFrame equivalent with the same semantics. (`monotonically_increasing_id()` is *not* a dense sequence; claiming it is, is a common wrong answer.) **Library interop.** GraphX is an RDD-based API. Older `spark.mllib` entry points take RDDs, while `spark.ml` is DataFrame-based. Some third-party connectors still expose RDDs. **Fine-grained per-partition resource control.** Opening one database connection or loading one model per partition. This is a real need — but note `Dataset.mapPartitions` in Scala and `mapInPandas` / `applyInPandas` in PySpark cover it without leaving the structured API. ## What the conversion actually costs `df.rdd` is not a free view. Spark executes the plan up to that point and *materializes* every binary row into a JVM `Row` object — allocation, GC pressure, and the end of code generation from there on. In PySpark it goes further and pickles each row into the Python worker. Coming back with `spark.createDataFrame(rdd, schema)` requires you to restate a schema, and inference over an RDD triggers its own sampling job. A pipeline that bounces between the two APIs pays this toll each way. ## The modern alternatives you should reach for first Before `df.rdd`, check whether one of these does the job: the built-in functions in `pyspark.sql.functions` (there are hundreds, including higher-order functions over arrays and maps such as `transform` and `aggregate`); `pandas_udf` for vectorized scalar logic; `mapInPandas` for arbitrary per-batch Python over an iterator of pandas DataFrames; `applyInPandas` with `groupBy` for per-group logic; `Dataset.mapPartitions` or typed `map` in Scala. Most historical RDD use cases in real pipelines are covered by these, and they keep the rest of the plan optimizable. ## How to talk about it in an interview The strong answer has three moves. First, state the presumption: DataFrames by default because the optimizer needs to see intent. Second, give one or two concrete cases where the RDD API is genuinely the right tool, and be specific — a custom partitioner for repeated co-partitioned joins is a better answer than "low-level control". Third, describe containment: drop down for the smallest step, convert back immediately with an explicit schema, and never let an RDD conversion sit upstream of the joins and aggregations that most need optimizing. A weak answer says "RDDs are for unstructured data and DataFrames for structured data" and stops there. It is not wrong, but it does not show you have measured what the conversion costs or that you know the vectorized alternatives that removed most of the historical reasons to reach for RDDs at all. ## Migration reality On a legacy codebase full of RDD jobs, the pragmatic path is not a rewrite. Push the source read and the final write into the DataFrame API first — that alone recovers pushdown and columnar I/O — then replace shuffle-heavy `reduceByKey`/`join` middles, and leave the genuinely opaque custom logic in `mapPartitions` until it earns a rewrite.
- What exactly happens when you call df.rdd in PySpark?Spark executes the plan up to that point and converts each internal binary row into a JVM `Row` object, ending code generation from there on; in PySpark each row is then pickled into a Python worker. It is a materialization boundary, not a cheap view, and everything downstream of it is unoptimized.
- You need a dense sequential row number across a large DataFrame. Is zipWithIndex on the RDD the right tool?It is one correct tool and it is cheap, but be clear about semantics: `monotonically_increasing_id()` is monotonic, not dense, and a `row_number()` window without a partition column funnels every row into one partition. `zipWithIndex` gives a dense index at the cost of an RDD round trip, so choose based on whether density is genuinely required.
- How would you approach a legacy codebase of RDD-only Spark jobs?Not with a rewrite. Convert the source read and final write to the DataFrame API first, which recovers columnar I/O, pushdown and pruning for free. Then replace the shuffle-heavy middles, since that is where the optimizer and adaptive execution pay most. Leave genuinely opaque custom logic in `mapPartitions` until it justifies the change.
saying these in an interview costs you the question
- Says RDDs are always faster because they are lower level
- Treats df.rdd as a free view rather than a materialization
- Claims monotonically_increasing_id produces a dense sequence
- Thinks RDDs are deprecated and must never be used
- Drops a whole pipeline to RDDs for one custom parsing step