cancel
Showing results forย 
Search instead forย 
Did you mean:ย 
Data Engineering
Join discussions on data engineering best practices, architectures, and optimization strategies within the Databricks Community. Exchange insights and solutions with fellow data engineers.
cancel
Showing results forย 
Search instead forย 
Did you mean:ย 

At what data size do you stop reaching for Spark?

Islam_hoti
New Contributor III

Hi everyone,

A question I keep having with my team and I would like to hear how others think about it.

A lot of the jobs we run are not big. Plenty of our pipelines process a few gigabytes, some considerably less. We run them on Spark because that is what the platform is, but a single machine with a columnar engine would handle most of them comfortably, and often faster, because there is no shuffle and no cluster to start.

The counterargument is real too. One engine means one set of skills, one deployment pattern, one place to look when something breaks. Introducing a second execution path for small jobs means maintaining two of everything, and the small job that quietly grows into a large one becomes a migration problem.

So I am curious where people have landed.

Do you have a size threshold below which you deliberately avoid Spark, or do you run everything through it regardless?

For those running small jobs on Spark, is single node compute the answer, or do you find the overhead still dominates?

Has anyone introduced a second engine for small workloads and later regretted it? Or regretted not doing it?

And the honest version of the question: how much of this is a real engineering decision and how much is just the cost of standardisation being worth paying?

Interested in how teams reasoned about it, not just what they picked.

1 ACCEPTED SOLUTION

Accepted Solutions

Khasim_1
New Contributor III

Hi @Islam_hoti ,

This is the 'Silent Architect' dilemmaโ€”when the standard tool is technically suboptimal but organizationally necessary. Weโ€™ve wrestled with this, and hereโ€™s where we landed:

  1. The 'Standardization Premium' is real: We reasoned that the cost of maintaining a second execution engine (e.g., Python/Pandas on serverless functions or a separate DB) is almost always higher than the cost of 'wasted' Spark overhead. You have to maintain two sets of CI/CD pipelines, two monitoring stacks, and two skill sets. Unless your small jobs are truly massive in number, the cognitive load of a second engine rarely pays for itself.
  2. The 'Spark Overhead' Mitigation: We stopped viewing 'Small Spark' as a failure. Instead, we use:
  • Databricks Serverless Compute: This eliminated the 'cluster startup time' pain point. If the job starts in 2 seconds, I no longer care if itโ€™s running on a distributed engine for a small dataset.
  • Single-Node Clusters: For the few jobs that really were 'micro-datasets,' we use single-node clusters. It keeps the same Spark APIs but eliminates the shuffle overhead entirely.
  1. The 'Growth' Trap: You hit the nail on the head: the 'migration problem.' We once had a team move a 500MB job to a local Python script to save costs. Six months later, that job grew to 100GB, hit OOM errors on the local machine, and caused a production outage. Migrating that logic back into a Spark-native framework was far more expensive than just having run it on Spark from day one.

Our Verdict: We now run everything through Spark regardless of size. The 'Cost of Standardization' is an insurance policy against migration debt. If the job is truly tiny, we optimize the compute (Serverless/Single-node) rather than changing the engine.

Itโ€™s an engineering decision disguised as a cost one: standardizing on one engine creates 'velocity' that often outweighs 'efficiency' in the long run.

Data Architect | 13 Years Domain Expertise | Databricks SA Champion Cohort

View solution in original post

3 REPLIES 3

juanlozadab
New Contributor II

I'll answer the size question directly first: no, we don't have a GB threshold, and I've come to think a number there doesn't survive contact with reality.

The reason is that the same 5 GB behaves completely differently depending on shape. Five gigabytes filtered and written out is trivially a single-machine job. Five gigabytes going through a wide join against a dimension table, or a window function over a skewed key, is not. And no threshold expressed in gigabytes tells you which one you're holding. So a size rule ends up either so conservative it never fires, or it fires on the job that turns out to need the shuffle.

What I'd threshold on instead is where the wall clock actually goes. Most "Spark is overkill" complaints I read are startup complaints, and those have a config fix rather than an architecture fix. Serverless jobs in standard mode carry a documented 4-6 minute startup latency by design. Performance optimized mode exists for exactly this case, and instance pools do the same on classic compute. If a 3 GB pipeline takes 7 minutes and 5 are provisioning, a second engine fixes the 2 minutes nobody was complaining about.

Which also answers your single node question, I think: single node removes the shuffle, not the cold start. So it helps if your bottleneck was coordination overhead, and does nothing if your bottleneck was waiting for compute. Worth knowing which one you have before picking.

On the second engine, what decides it for me is the asymmetry. A small job on Spark is a tax: bounded, predictable, shrinking every time start times improve. A job that outgrew the small engine is a rewrite: unbounded, and it lands at the worst possible moment. I'll take a known tax over an unknown rewrite. And the cost people underestimate isn't the tooling, it's that a second execution path doubles the failure modes. That's an on-call conversation, not a performance one.

So yes, a lot of this is the cost of standardisation. I don't think that's a cop-out, just something worth naming rather than dressing up as a performance decision.

 

jlb

shirley54reece
Visitor

Most teams ultimately stick with Spark for small workloads despite the overhead, because the hidden costs of maintaining a dual-engine stackโ€”fragmented skill sets, separate deployment pipelines, and migration headaches when small jobs scaleโ€”usually outweigh raw performance gains. However, if the platform overhead routinely stalls velocity or balloons cloud costs, introducing a lightweight single-node alternative like DuckDB or Polars is justified, provided your team has the operational capacity to manage two distinct execution paths.

Khasim_1
New Contributor III

Hi @Islam_hoti ,

This is the 'Silent Architect' dilemmaโ€”when the standard tool is technically suboptimal but organizationally necessary. Weโ€™ve wrestled with this, and hereโ€™s where we landed:

  1. The 'Standardization Premium' is real: We reasoned that the cost of maintaining a second execution engine (e.g., Python/Pandas on serverless functions or a separate DB) is almost always higher than the cost of 'wasted' Spark overhead. You have to maintain two sets of CI/CD pipelines, two monitoring stacks, and two skill sets. Unless your small jobs are truly massive in number, the cognitive load of a second engine rarely pays for itself.
  2. The 'Spark Overhead' Mitigation: We stopped viewing 'Small Spark' as a failure. Instead, we use:
  • Databricks Serverless Compute: This eliminated the 'cluster startup time' pain point. If the job starts in 2 seconds, I no longer care if itโ€™s running on a distributed engine for a small dataset.
  • Single-Node Clusters: For the few jobs that really were 'micro-datasets,' we use single-node clusters. It keeps the same Spark APIs but eliminates the shuffle overhead entirely.
  1. The 'Growth' Trap: You hit the nail on the head: the 'migration problem.' We once had a team move a 500MB job to a local Python script to save costs. Six months later, that job grew to 100GB, hit OOM errors on the local machine, and caused a production outage. Migrating that logic back into a Spark-native framework was far more expensive than just having run it on Spark from day one.

Our Verdict: We now run everything through Spark regardless of size. The 'Cost of Standardization' is an insurance policy against migration debt. If the job is truly tiny, we optimize the compute (Serverless/Single-node) rather than changing the engine.

Itโ€™s an engineering decision disguised as a cost one: standardizing on one engine creates 'velocity' that often outweighs 'efficiency' in the long run.

Data Architect | 13 Years Domain Expertise | Databricks SA Champion Cohort