Writing Databricks Parquet That VertiPaq Likes

This post, the repo and every number in it are personal opinion. Nothing here is a Fabric position, benchmark or recommendation.

My colleague Eiki published a post on Microsoft’s white paper on Power BI architecture choices for Azure Databricks. Read it first: it covers the four storage modes, the headline findings, and links to the paper. I assume you know what the paper measured.

I like the paper for a different reason. It benchmarks a lakehouse as designed: open files on object storage, one vendor writes, another reads, Parquet is the only contract. Such benchmarks are rare, hard to reproduce and usually opinionated. This one is deliberately neutral (maybe too neutral) , The dataset is so uniform that V-Order has little room to differentiate the writers, so all writers start on relatively equal ground.

I have written before about VertiPaq reading Parquet that other engines wrote: Optimizing Parquet Layout for Power BI Direct Lake Mode (Dec 2024), Just VOrder, don’t try to understand how VertiPaq works (Nov 2025) and Writing Parquet That VertiPaq Likes (Aug 2026).

The short version: VertiPaq reads Parquet from any producer. It does not need V-Order. It prefers a layout it can transcode fast and hold small, and V-Order is one way to produce it.

This post applies that idea to Spark on Databricks. I am not a Spark expert; I know delta-rs much better. The Parquet concepts carry over, the knobs do not, and Databricks has many settings that interact.

AI helped me navigate them. Where the repo gets Spark wrong, the mistake is mine. Corrections welcome.

What I did

The paper tunes the Databricks write with the mainstream, documented settings, as an industry paper has to. This is a personal blog, so I can be less orthodox and go low level 🙂 I tried a few little-known knobs instead.

All of them are cluster configuration for row-group size and dictionary encoding, the same lines for every table, plus one optional step per fact table: a single clustering key on the column the reports filter on (nothing fancy it is just a global sort).

I used the paper’s protocol: TPC-DS at SF100 and SF1000, its DAX capture, 20 concurrent readers, three load tests back to back on one model. Two changes: the semantic model is deleted after the three runs so the next run starts cold, and the OneLake cache is on because my Databricks tenant is in another region. Only the very first run at a scale factor reads across regions; every later run 1 reads from the cache, so it measures the transcode alone.

Note: I used Fabric layout as it is.

Databricks writes, Direct Lake reads through the mirrored catalog. The configuration, what each line does and what breaks it are in the repo. This post is results only.

The results

Each chart shows the configuration alone, the configuration with the clustering key, and the best layout of each other writer: delta-rs sorted on the date key, and the paper’s Fabric layout of one partition per date with Z-order and V-Order.

Run 1 pays the transcode; runs 2 and 3 rerun the same queries on the loaded model.

At SF100, the configuration alone produces dictionary encoding and large row groups, but is about three times slower than every ordered layout. Adding the clustering key on the date matches the paper’s V-Order layout, with row groups more than ten times larger.

At SF1000, the configuration alone never really warms up: each run is slower than the last, while the three ordered layouts do. That is why the clustering key is not just an optimisation in the recipe; it is what makes the layout usable at this scale. The clustered arm still trails delta-rs and the paper’s layout at SF1000: the clustering did not fully sort the rows at that scale, and the repo carries that as open. Why delta-rs has the best cold run at both scale factors is still a mystery to me.

The interesting part is that VertiPaq does not care who wrote the Parquet. It cares about the physical layout it has to transcode.

The repo

github.com/djouallah/parquet_layout_vertipaq_spark has the configuration, the clustering step, what undoes it, the numbers behind both charts and what is still open.

It also ships a Claude Code skill that applies the recipe.

AI, dbt and Iceberg are already changing data engineering

I have been using AI in VS Code and Onelake, initially trying to make sense of Chat with your data (without much success, that thing is very hard and we need some breakthrough), and more recently for data engineering. I noticed something: AI has become useful enough in the last couple of months that it is changing my workflow.

None of the individual pieces here is new or particularly interesting by itself. But combine them and we have something!!! CI/CD, Iceberg REST catalog, Opus 4.8, me discovering how OIDC GitHub integration works in Fabric, and dbt, and suddenly everything makes more sense.

Let’s take a simple ETL job, read some CSV, clean it, and produce high quality Parquet files that humans and AI can consume.

From a human perspective, and although we like to think our judgement is based purely on cost and performance, that’s never really been the case. It is always influenced by personal bias: Spark people will always use Spark, T-SQL people will always use the DWH, and Pandas people don’t care, they will use their thing.

Now, if we imagine an AI doing that, the incentives are different. Yes, AI is biased by its training data, but it isn’t biased by tribe. Judging by my AI agent (I suspect it is telling me what I want to hear), it doesn’t care. It prefers short loops, and if you tell it “I want the cheapest option”, it is smart enough to try to do that.

1- ETL is just processing raw data into something coherent that can be consumed. It is a deterministic process, and as someone who mainly used GUI tools, it took me a while to get it: data engineering is just code!!! Those tools are writing code, it just happens that we don’t see it 🙂

2- All things considered, an AI agent has no personal attachment to an engine, assuming you give it a strict spec. An agent will not favour an engine because of familiarity, or because it spent so much time using it that it became tribal.

3- Opus 4.8 class AI is good enough for general purpose data engineering. The current issue is the cost. Right now you need a Max subscription in practise and it is just not sustainable, but I am confident the market will figure out a solution. We don’t need AGI, just a super cheap Opus 5 alternative 🙂

4- I am not saying all engines are equal, or that they become just a SQL runtime. That’s not true. But because the cost of switching from one SQL dialect to another is minimal, the human excuse of “I use what I know” (which is pretty tragic, when you think about it) will no longer be relevant.

Selecting an engine will be strictly based on facts: engine 1 can do efficient MERGE, engine 2 cannot, engine 1 wins. That’s fair.

And engines still have a lot to improve, see for example non-trivial incremental processing, async remote scans, partial caching, Dynamic Horizontal scaling etc.

And to be fair, the dialect alone is not the whole story. Subtle differences in engine behaviour, even when using the same SQL text, can be a pain. I learnt this the hard way. I spent too much time trying to debug a result, only to find out that two engines have different behaviour when doing a join using a word with padding: “spain” and “spain ” may or may not mean the same thing.

But the good thing is that a cross-engine parity check catches these subtleties automatically. Either the numbers match or they don’t (you can’t do that with Chat with your data).

5- Cloud storage vendors ultimately want to store more data, regardless of how it was processed — their incentives are aligned with this agentic trend. In my personal opinion, they need to spend less time on MCP and other AI stuff, and just focus on the boring stuff: authentication, documentation, and interoperability.

If you want to deviate from the catalog spec, at least spend some engineering time working with the open source engines to support your own variant.

6- The less a client has to do, the better for everyone: move more stuff to the catalog, including table scanning and server-side planning.

A read-only client has no interest in reading Avro files. Just give me a list of Parquet files and deletion vectors to read. Keep it opaque. Store it in a database, I don’t care. Just give me the data files to read.

7-I added this note because of early feedback on the blog, which was basically: “I don’t care about Iceberg, I use Delta.” That’s a fair reaction. The good news is it won’t really matter — Delta 5 will use the same adaptive metadata as Iceberg V4. They’re still two separate projects with different governance model(or rather, the lack of it), but ultimately they’ll produce the same thing in theory . Just give it another 2 years or so, there is a real table format fatigure, translating from one to another is a waste of everyone time, and to be honest, some vendors find it is much cheaper just to lock users using credential vending instead of weaponising the table format 🙂

Instead of another abstract “thought leadership” blog post 🙂 here’s a concrete example: the same dbt project, unchanged, runs against four different vendors’ Iceberg REST catalogs: OneLake, Cloudflare R2, S3 Tables and Snowflake Horizon, each in about a minute, on a throwaway GitHub Actions runner with DuckDB as the engine.

Switching catalog is literally changing one ATTACH in profiles.yml.

The whole thing is public: testing-iceberg-rest-catalog.

Hopefully, data engineering will move away from configuring and fine-tuning engines and towards talking more to end users and understanding what they need: agreeing on what the numbers are supposed to mean, and writing the tests that catch it when they are wrong.

The machine can build the pipeline, and it will probably be better than us at it, but it doesn’t know what “correct” means. I was tempted to write something about ontology but I am not going there 🙂

And getting access to that ERP will always require a human. That’s an organizational thing, and no AI can fix it.

As a data analyst by trade, I always found data engineering a chore. I never enjoyed it , and to be honest ,it does not bother me if writing transformation becomes fully automated

So I have a bias, to be honest, and maybe this blog is just wishful thinking. But what if it is true? I think, at least, pay attention to this new trend.

How far Python alone can take you on Delta

1. delta-rs is an ACID Delta writer

delta-rs implements the Delta Lake protocol natively. mergeupdate, and delete go through optimistic concurrency control on every commit. No external coordinator, no catalog service. Two writers race for the same version of the log, one wins, the other retries.

All you need is a path. No metastore to provision, no catalog endpoint, no JDBC connection, no warehouse to wake up. A folder on disk (or on ADLS / S3 / GCS) is the whole interface.

Setup: B is a Delta table being fed a series of CSV batches (batch_001.csvbatch_002.csv, …). Each merge should ingest only files B hasn’t seen yet.

A naming note: the project is delta-rs but the Python package is deltalake (pip install deltalake). On Fabric, stick with what’s preinstalled — Python notebooks already ship with deltalake and OneLake access configured.

From the notebook:

# Bootstrap target B with batch_001 already ingested
write_deltalake(Target_PATH, pa.table({...}), mode="overwrite")
vB = DeltaTable(Target_PATH).version() # v0
# Compute the rows to ingest from the target's current state
con.sql(f"ATTACH '{Target_PATH}' AS tgt (TYPE delta, VERSION {vB});")
our_rows = con.sql("""
SELECT s.id, s.value, parse_filename(s.filename) AS filename
FROM read_csv_auto('source_csv/*.csv', filename=true) s
WHERE parse_filename(s.filename) NOT IN (SELECT DISTINCT filename FROM tgt)
""").arrow()
# → 80 new rows from batch_002..005
# First merge: 80 inserts, commits cleanly
DeltaTable(Target_PATH).merge(
source=our_rows,
predicate="t.filename = s.filename",
source_alias="s", target_alias="t",
).when_not_matched_insert_all().execute()
# Same merge re-run: 0 inserts. The predicate is idempotent.
DeltaTable(Target_PATH).merge(...).when_not_matched_insert_all().execute()

Two commits, both correct. The second run does nothing because the predicate already sees the rows. The transaction model travels with the table itself: move the folder, open it from another machine, and the next writer continues from the last commit.

write_deltalake(mode="append") and write_deltalake(mode="overwrite") are blind on purpose. Blind append means N concurrent appenders all succeed and the result is the union of their rows — exactly what you want for event streams or log ingestion. Blind overwrite means the new data wins and whatever was there is gone — what you want when the writer is the authoritative source for the table. OCC only kicks in for operations that actually read the target (mergeupdatedelete), since those are the only ones where a concurrent change can invalidate what you just computed.

2. I want the full read-to-write transaction, Python API is fine

A common pattern: DuckDB or Polars reads, transforms, and hands an Arrow table to delta-rs to commit. The notebook above is exactly that shape — DuckDB computes “filenames not yet in B” and delta-rs merges the result.

Inside delta-rs, OCC still works. What it cannot see is the read on the other side of the engine boundary. delta-rs knows about the merge it is about to commit; it does not know that DuckDB read B at version vB thirty seconds ago.

Carry the snapshot across the boundary by pinning both sides to the same version:

vB = DeltaTable(Target_PATH).version()
import duckdb
con = duckdb.connect()
con.sql(f"ATTACH '{Target_PATH}' AS tgt (TYPE delta, VERSION {vB});")
our_rows = con.sql("SELECT ...").arrow()
DeltaTable(Target_PATH, version=vB).merge( # ← pinned
source=our_rows,
predicate="t.filename = s.filename",
source_alias="s", target_alias="t",
).when_not_matched_insert_all().execute()

The OCC check now compares against vB instead of HEAD. If another process touched B in the meantime — say a parallel job deleted batch_001.csv — the pinned merge raises:

Failed to commit transaction: Commit failed: a concurrent transaction deleted data this operation read.

Catch it, recompute the diff against fresh state, retry. On the Polars side, pl.read_delta(path, version=vB) accepts the same pin, so the pattern works for any reader that exposes versioned reads.

The pin is just a number. No new infrastructure, no shared coordinator, still path-based.

3. I don’t want the Python API, I want SQL only

If you would rather write SQL — say, drive the pipeline from dbt — your options on Delta today are Spark and Fabric Data Warehouse. Both have supported dbt adapters and work great in production. I have to admit, I was hoping DuckDB would fill that gap, since it is a database and SQL-level transactions are what you expect from a database. The market went the other way: investment is going into catalog-based lakehouse formats (DuckLake, Iceberg), and the DuckDB Delta writer that does exist is tied to Unity Catalog and limited to blind appends. I don’t see them investing in a file-based conflict resolver any time soon 🙂 Lakesail seems interested in this use case, but it is still too early to call.

Takeaway

I personally use delta-rs for CSV ingestion, appends, and recording results from high-concurrency performance tests — it is fast, cheap, and bullet-proof in those scenarios. The open source maintainers are very helpful and care deeply about the product, as they use it themselves in production. But it is not the right tool for every case; Data Warehouse and Spark are more appropriate for complex workloads. With time you intuitively pick the tool that makes sense for a particular job and how much compute you can spend. None of that has to be an either/or: at the end of the day it is a lakehouse, and the whole concept of a lakehouse is having the option to choose the engine. That option matters — if we say only one engine (open source or not) is blessed for writes, then there is no point in the concept of a lakehouse.


Notebook: https://github.com/djouallah/Fabric_Notebooks_Demo/blob/main/TableFormat/delta/occ.ipynb

Thanks Raki for keeping me honest:)

Thanks to Ion for explaining how version worked when doing merge: https://www.linkedin.com/in/ionkoutsouris/

Edit : how about Spark

Thanks to Frithjof for explaining Spark behaviour : The merge fixes one snapshot at transaction start (current HEAD = post-delete) and uses it for both its scan and its conflict check. Internally consistent — but bound to HEAD-at-merge-start, which Spark chose, not to the state our read saw, same behaviour when using delta_rs with a lazy dataframe : https://github.com/djouallah/Fabric_Notebooks_Demo/blob/main/TableFormat/delta/occ_spark.ipynb

Ensuring safe single-writer for DuckLake on OneLake using file lease

DuckLake supports multi-writer just fine — but only if your catalog is a real database, like Postgres (there’s some interest in SQL Server support too). But if all you have is object storage and a SQLite or DuckDB file as the catalog, you’re stuck with single-writer: object stores aren’t real filesystems, so the DB file can’t be locked. Nothing stops two processes from writing to it at the same time and corrupting it.

If single-writer is enough for you (one notebook, one pipeline, one user), you don’t need to stand up a database server. You just need accidental concurrent runs to fail fast.

The trick: take a blob lease

OneLake speaks the ADLS API, so you can take a lease on a blob — a mutex for free (it seems S3 needs DynamoDB and GCS needs a homemade lock object). Each run does:

  1. Acquire a lease on metadata.db in abfss://.
  2. Download it to local disk of the notebook.
  3. Point DuckLake at the local copy and do the work.
  4. Upload the modified file under the lease.
  5. Release the lease.

A second notebook that starts while the lease is held fails immediately on acquire_lease. It can’t even read a stale copy. and you can’t delete the file using the UI , I can see already some uses cases here:)

What about crashed runs?

ADLS leases are either 15–60 seconds fixed, or infinite. Fixed leases need a heartbeat — annoying inside a notebook. Infinite leases work until something crashes — then the file is stuck.

The fix: take an infinite lease, but stamp acquired_at = <utc iso> into the blob’s own metadata when you acquire. When the next run hits a lease conflict, read that timestamp. Older than 12 hours? Call break_lease and re-acquire. A crashed run self-heals within 12 hours. You can shorten that window, or break the lease manually with a one-line script if you can’t wait — there’s a snippet in the README.

Code is here.