Why The World Needs Flarion. Read More

Storage Format Lock-in: The Constraint Limiting Query Performance

Processing engines require rigid storage formats, so queries run slower than the data allows. Where Parquet falls short for selective analytics, and why.
By
read time
November 17, 2025

Modern distributed processing engines achieve remarkable scale and reliability, yet they share a fundamental architectural constraint: rigid storage format requirements. This constraint forces organizations to accept significant performance penalties, sometimes 10-100x slower queries than technically possible, simply because their processing engine cannot adapt to optimized storage formats.

Why Parquet Falls Short for Selective Analytics

Parquet excels as a universal columnar format, but its design prioritizes compatibility and compression over query performance. Understanding its limitations requires examining how analytical queries actually execute.

The Row Group Problem

Parquet organizes data into row groups, typically 100,000 to 1 million rows each. This coarse granularity creates a fundamental problem: if even a single row in a row group matches your query filter, you must read ALL requested columns for the ENTIRE row group.

Consider a query filtering for a specific customer ID in a billion-row table. That customer's 100 transactions might be scattered across 100 different row groups. Parquet must read 100 row groups × 100,000 rows × all requested columns, processing 10 million rows to return 100.

Single-Phase Execution

Parquet readers process queries in a single phase:

1. Read all requested columns for relevant row groups

2. Decompress the data

3. Apply filters

4. Return matching rows

This means reading and decompressing massive amounts of data that will immediately be discarded. There's no mechanism to filter first, then read only what's needed.

Limited Predicate Pushdown

While Parquet stores min/max statistics per row group and optional Bloom filters, these only help skip entire row groups. For analytical queries where some matching rows exist in many row groups, this provides minimal benefit. You still read entire 100,000-row chunks to extract perhaps 10 rows from each.

How Query-Optimized Formats Solve These Problems

Formats like ClickHouse's MergeTree take a fundamentally different approach:

Granular Storage

Instead of 100,000-row groups, data is organized into 8,192-row granules. This 12x finer granularity means reading much less unnecessary data when matches are sparse. Finding those same 100 customer transactions requires reading ~12x less data just from granularity alone.

Two-Phase Execution with PREWHERE

Query-optimized formats implement two-phase execution:

1. Phase 1: Read ONLY filter columns for candidate granules

2. Phase 2: For rows that pass filters, read the remaining requested columns

This seemingly simple change has a profound impact. Instead of reading 20 columns for millions of rows, you read 1-2 filter columns first, identify the 100 matching rows, then read the other 18 columns for just those 100 rows.

Sparse Indexing

A sparse primary index stores one entry per granule (every 8,192 rows), creating a tiny index that fits entirely in memory even for billion-row tables. Binary search on this index instantly identifies which granules to read, eliminating 99%+ of data before any I/O occurs.

Lazy Materialization

Column reads are deferred until absolutely necessary. If a query has multiple filters, the engine applies them progressively, reading additional columns only for rows that survive each filter. This minimizes decompression work and memory bandwidth usage.

The Performance Gap: A First-Principles Calculation

Let's calculate the actual performance difference using a realistic analytical query on data stored in S3:

Scenario:

  • 100 million rows, 100 columns
  • Query with 0.01% selectivity returning 10,000 rows
  • Projects 20 columns
  • All data in S3 (network I/O is the same for both formats)

Parquet Execution:

  • Data organized in 100,000-row groups
  • 10,000 matching rows spread across ~100 row groups
  • Must read all 20 requested columns for entire row groups containing any match
  • Compressed data transferred from S3: 160 MB (with 10:1 compression)
  • Data decompressed and processed: 1.6 GB

Query-Optimized Format (like ClickHouse MergeTree):

  • Data organized in 8,192-row granules with sparse index
  • Two-phase execution: read filter columns first, identify matches, then read other columns
  • Compressed data transferred from S3: 46 MB (with 3.5:1 compression)
  • Data decompressed and processed: 162 MB

Despite storing data 2.9x larger on S3 due to less aggressive compression, the query-optimized format transfers 3.5x less data over the network and processes 10x less decompressed data. Processing 1.6 GB vs 162 MB of decompressed data is the difference between 100ms and 10ms query time

Why Processing Engines Don't Support Alternative Formats

While Spark and Ray technically could support arbitrary storage formats through extension APIs, in practice they don't. The ecosystem has converged on Parquet, ORC, and Avro, leaving significant performance opportunities unexplored.

The Spark Reality

Spark provides a DataSource V2 API for custom formats, yet virtually no production deployments use them. The practical barriers:

  • Format implementations must handle Spark's complex internal row representation
  • Custom formats cannot leverage Spark's vectorized execution optimizations
  • The Catalyst optimizer cannot reason about custom format capabilities
  • Maintaining custom format readers requires deep Spark internals knowledge

Ray's Format Limitations

Ray delegates data loading to Pandas or PyArrow:

python

@ray.remote

def process_partition(file_path):

    # Limited to formats Pandas supports

    df = pd.read_parquet(file_path)

    return df.groupby('device_id').mean()

Adding optimized formats would require writing C++ extensions, ensuring serialization compatibility, and maintaining format-specific optimizations Ray doesn't understand.

The Architectural Impedance Mismatch

Even when custom format support exists, fundamental architectural assumptions prevent leveraging format-specific optimizations. Standard engines read all requested columns before filtering and lack the execution primitives for two-phase execution. The granularity mismatch between coarse row groups and fine granules cannot be bridged through format plugins.

Real-World Impact on S3-Based Data Lakes

Production systems using S3 data lakes demonstrate concrete impact:

E-commerce Recommendation Engine

  • Dataset in S3: 500M products × 800 attributes
  • Parquet in Spark: 8.2 seconds average latency
  • Possible with optimized format: 100ms
  • Impact: Real-time recommendations impossible

Financial Risk Analytics

  • Dataset in S3: 10 years of trades, 2,000 columns
  • Parquet in Spark: 45 seconds per portfolio
  • Possible with optimized format: 500ms
  • Impact: Risk managers wait minutes for updates that could take seconds

IoT Monitoring

  • Dataset in S3: 100,000 devices × 1,000 metrics
  • Parquet in Ray: 12 seconds per device group
  • Possible with optimized format: 50ms
  • Impact: Alert latency prevents real-time response

The Compression-Performance Trade-off

Parquet achieves superior compression, typically 2-3x better than query-optimized formats. A 100GB dataset might compress to 10GB in Parquet versus 30GB in an optimized format, translating to $2.30/month versus $6.90/month on S3.

However, compute costs dominate this equation. When queries run 10-100x faster with optimized formats, the required infrastructure shrinks proportionally. A workload requiring a 16-node Spark cluster at $8,000/month might run on 2 nodes at $1,000/month with format optimization. The $7,000/month compute savings overwhelms the $4.60/month additional S3 storage cost by a factor of 1,500x.

Performance Patterns by Query Type

Different query patterns show varying sensitivity to format selection:

Highly Selective Queries (<0.1% of rows returned)

  • Performance difference: 20-100x
  • Parquet's coarse row groups create massive read amplification

Wide Table Projections (many columns, few rows)

  • Performance difference: 10-50x
  • Single-phase execution forces reading all columns upfront

Aggregation Queries

  • Performance difference: 5-20x
  • Processing unnecessary rows dominates computation

Full Table Scans

  • Performance difference: 0.7-1.4x (Parquet often faster)
  • Parquet's superior compression provides advantage when reading everything

The Hidden Costs of Format Lock-in

Beyond direct performance impact, format rigidity creates cascading inefficiencies:

Infrastructure Over-provisioning: Organizations scale clusters to compensate for format inefficiency. A workload naturally requiring 2 nodes might run on 20 nodes to achieve acceptable latency.

Architectural Workarounds: Teams implement pre-aggregation pipelines, materialized views, and caching layers, each adding operational complexity without addressing the root cause.

Opportunity Costs: When queries take 30 seconds instead of 300ms, entire categories of applications become infeasible.

The Path Forward

Parquet remains excellent for data interchange, archival storage, and full-scan workloads. The opportunity lies in enabling processing engines to leverage format diversity based on workload requirements.

The ideal architecture would support understanding the profile of the workloads, adjusting compaction accordingly, and using the best data format for the job. Current processing engines cannot implement this strategy due to their format rigidity. They lack the flexibility to transparently read different formats while applying format-specific optimizations.

This limitation is beginning to change. Next-generation execution engines recognize that format flexibility is essential for modern analytical workloads. By decoupling query execution from storage assumptions, systems like Flarion can deliver the performance that has always been technically possible but practically unreachable.

Organizations no longer need to accept 10x performance penalties as the price of distributed processing. The ability to match storage formats to query patterns represents the difference between queries that frustrate users and analytics that drive real-time decisions.

Related Posts

A few weeks ago AWS shipped the Spark Upgrade Agent, an AI agent that migrates Spark jobs to Spark 4.0. You point it at a repo, it rewrites deprecated APIs, adjusts for behavioral changes, updates the build for Scala 2.13, submits the result to an EMR cluster, and iterates on failures until the job runs. It handles both Scala and PySpark, and it works the way you'd hope an agent would: plan, transform, validate, repeat. The potential payoff is large - newer Spark versions have better performance and years of accumulated bug fixes.

It looks like a good tool and data teams looking into a Spark migration should consider it, but what we’ve found is that in the enterprise, rewriting code isn’t the main impediment to upgrading Spark workloads. Spark programs are part of complex pipelines, parts of which are poorly understood or maintained, and making changes to a sensitive system is inherently risky. There’s no guarantee that the output data will actually remain the same, and data integrity is the fundamental challenge of completing such a migration.

Several companies have written in detail about major Spark upgrades, and the data integrity challenge is a recurring theme.

Slack

Slack's migration from Spark 2 to Spark 3 took about a year across 60+ EMR clusters and 40+ teams. They saw some code-level breakage: `RAND()` in join keys became an `AnalysisException`, some casts that Spark 2 tolerated started failing, the `Greatest` function handled NULLs differently than its Hive counterpart. Validation was a much larger effort. For billing pipelines, Slack required exact matches, building test tables from production data on Spark 3 and running `EXCEPT` and `COUNT` comparisons in Trino against the Spark 2 outputs, with a Python framework for digging into every discrepancy. The discrepancies weren't all bugs. Non-deterministic row ordering, timestamp variations, and genuine semantic differences between Hive and Spark implementations all produce diffs that need to be investigated. Some are noise, some are real regressions.

Uber

Uber's version of this is bigger and more instructive. They migrated from Spark 2.4 to 3.3 with over two million Spark applications running daily. The code transformation was automated with Polyglot Piranha, their structural rewrite tool. It parses the source code into an AST, matches patterns, and applies transformation rules, including inserting legacy flags like `spark.sql.legacy.allowUntypedScalaUDF` where old behavior had to be preserved. This scaled well. The problem that shaped the whole project was stated plainly: "We had over 40,000 Spark apps, so we couldn't decentralize the data validation." No staging environment, no test cases, no way to ask every team to eyeball their own outputs.

So the flagship engineering artifact of Uber's Spark upgrade wasn't actually a code migrator but Iron Dome: a shadow-testing framework which runs the migrated job against production inputs, rewrites output paths at runtime so results land in staging instead of production, puts guardrails at the Hadoop FileSystem interface so a misrouted write can't touch real data, then compares the shadow output against the production run and only marks the job migrated when they agree. 

Facebook

None of this is specific to the Spark 2-to-3 transition, or even to Spark versions. When Facebook moved Hive workloads onto Spark SQL back in 2017, they ran shadow pipelines writing to tables suffixed `_spark_shadow` so downstream jobs were never exposed, used count checks as a cheap first filter, and reached for full hash validation of outputs only reluctantly (because, as they put it, the hash validation was "sometimes even heavier than the query itself.") Funny enough, proving the new engine produced the same answer could cost more compute than producing the answer. They note that non-deterministic UDFs made validation hard, the same diff-adjudication problem Slack hit eight years later.

Takeaway

In the enterprise, migrations are rightfully considered risky projects that take time and incur risk. This is especially true in the age of AI given that the things that AI doesn’t necessarily deliver are also the riskiest parts of the migration - edge cases, data integrity, the long tail, etc. This isn’t to say that AI can’t help with building tooling for a migration, it clearly can, but usually an agent isn’t going to do the trick alone.

At Flarion, a major goal of ours is to give users the best possible performance and access to modern features without requiring a code migration. If you can get the benefits of Spark 4.2 while staying on Spark 3.4 then that’s a huge time save and reduction in risk. Of course, the data integrity problem doesn’t disappear, it’s now Flarion’s responsibility. One we’re happy to shoulder.

Spark 4.2 was released in mid-July, and a lot of the attention went to the headline features: geospatial types, change data capture, vector search. Tucked into the performance section was a smaller item that we found particularly interesting. Arrow-optimized Python UDFs, and Arrow-based data exchange with Python in general, are now on by default.

One can consider this a mere config flip. The feature has existed since Spark 3.5. But defaults are how a platform tells you what it considers normal, and Spark just declared that the normal way to move data between the JVM and Python is Apache Arrow. It's worth walking through why that boundary was slow in the first place, what Arrow does about it, and what it changes for everything sitting underneath.

What a Python UDF Actually Costs

PySpark has a split identity. The engine that executes your job runs on the JVM, and your UDF runs in a separate Python process, because that's where Python code has to run. Every batch of rows that passes through the UDF makes a round trip: out of the JVM, across a socket into the Python worker, through your function, and back.

Until now, the default way to make that trip was pickle. Each row was converted from Spark's internal representation into a Python object, serialized, sent across, deserialized, processed, and then the whole sequence ran again in reverse. One row at a time, one object at a time. The work is pure overhead. It exists because the two sides of the boundary represented the same data differently, and translation was the only way across.

For UDF-heavy jobs the translation regularly cost more than the function being called. It's one of the oldest pieces of PySpark folklore: keep your logic in built-in expressions if you can, because the moment you write a Python UDF, you pay a tax that has nothing to do with what the UDF does.

What Arrow Brings to the Table

Apache Arrow is a specification for how tabular data is laid out in memory: columnar, in large contiguous buffers, with a defined binary layout for every type. The layout is the same regardless of which language or engine produced it. If two systems both hold data in Arrow format, one can hand the other a batch without converting anything, given that the bytes are already in the shape the receiver expects.

Applied to the Python boundary, this removes most of the tax. Instead of serializing rows into Python objects, the JVM sends Arrow batches, and the Python side reads them directly as columnar data. There is still a process boundary and still a copy across the socket, but the expensive part - turning every value into an object and back - is gone. Spark's own benchmarks for Arrow-optimized UDFs showed roughly 2x speedups on chained UDFs when the feature shipped in 3.5, with larger gains the more the workload was dominated by the boundary rather than the function.

The history of this feature tells you something about defaults. Arrow UDFs arrived as an opt-in in Spark 3.5. Spark 4.1 added Arrow-native UDF decorators that skip the pandas conversion entirely. With 4.2, Arrow is the default and pickle is the fallback. Everyone gets the faster boundary, including the large majority of users who never knew there was a flag.

The Ecosystem Keeps Converging on Arrow

The UDF change is one instance of a pattern that has been running for years. `toPandas` and `createDataFrame` now use Arrow by default too, in the same release. Spark Connect streams query results to clients as Arrow batches. Outside of Spark: pandas can be backed by Arrow, Polars is built on it, DuckDB reads and writes it natively, DataFusion uses it as its internal memory model. When these systems exchange data with each other, more and more often no conversion happens, because both ends already speak the same format.

This is what a de facto standard looks like while it's forming. Nobody mandated Arrow; each project adopted it because interoperating through a shared memory layout is cheaper than maintaining pairwise converters. Every Spark release for the past several years has replaced another row-based boundary with an Arrow one, and there's no reason to expect the direction to reverse.

What the Default Doesn't Change

It's worth being precise about what got faster. The boundary between the JVM and Python is now columnar. The engine on the JVM side of that boundary is the same one it was before — predominantly row-oriented, executing on the heap, one row at a time through most operators.

That produces a slightly odd shape for a typical PySpark job. The scan reads Parquet, which is columnar on disk. Spark turns it into rows to execute the joins and aggregations. At the UDF boundary, those rows are batched back into Arrow's columnar form, shipped to Python, processed, returned, and turned back into rows for whatever comes next. The data changes representation multiple times, and the fast columnar format only exists at the edges. Spark 4.2 made the edges cheap. The middle is where most of the job's time goes, and the middle didn't change.

Rewriting the executor is a different scale of undertaking than adopting Arrow at the boundaries, and the boundaries were a reasonable place to start. But the release defines the remaining gap fairly precisely: the format Spark now uses to talk to Python is not the format it uses to compute.

Running the Middle on Arrow Too

That gap is where Flarion sits. Our engine executes Spark's operators - scans, joins, aggregations, and the rest - in native Rust code built on Apache Arrow and DataFusion, replacing the row-at-a-time JVM path for the parts of the plan it supports. Inside the engine, data stays in Arrow's columnar layout the whole way through. It goes in as a plugin on the Spark job you already have; unsupported operations fall back to Spark and run the way they always did.

Spark standardizing its boundaries on Arrow makes this arrangement steadily cleaner. When the engine hands data back to Spark, or Spark hands data to a Python worker, both sides increasingly agree on the memory layout, so the crossings that used to require translation become handoffs. Data moves between Spark's JVM and our engine through Arrow's C Data Interface without copying at all — a pointer to the buffers crosses the boundary, and the data stays where it is.

The UDF change is a preview of what that feels like, applied to one boundary. The excitement around it comes from removing translation overhead at a single crossing point. An Arrow-native engine applies the same idea to the execution itself: the scan produces Arrow, the join consumes Arrow, and the representation never changes because there's nothing to change it into.

Summing Up

Spark 4.2's UDF change is a nice speedup, but the reason it caught our attention is what it says about direction. Spark is a conservative project and it doesn't change defaults lightly, because millions of jobs run on whatever the defaults are. When a project like that decides Arrow is how data should cross the Python boundary, it's acknowledging what the rest of the ecosystem already settled on. In short, the future is Arrow.

Oops! Something went wrong while submitting the form.