Skip to main content
Use this recipe when a Delta table ingested with the Delta Lake connector has more rows in Pinot than Databricks reports, and the table is a Databricks streaming table maintained by a CDC pipeline. Lakeflow Connect ingestion pipelines, such as the Salesforce connector, produce tables like this.

Why the counts differ

Databricks keeps older versions of updated rows inside the table’s Parquet files and marks them with hidden columns. The one that matters here is __DeleteVersion. When you query the table through Databricks or Unity Catalog, Databricks hides those superseded rows for you. The Delta Lake connector reads the physical Delta files, and it ingests every row in them. Nothing filters on the hidden columns, so each superseded version lands in Pinot as an extra row. Running OPTIMIZE on the source table doesn’t help, because the superseded rows survive compaction. How to tell this is your problem:
  • SELECT COUNT(*) in Pinot is higher than the same count in Databricks, and the gap grows as the source table takes more updates.
  • The extra rows have distinct primary keys and aren’t flagged as deleted in your business columns, so they don’t look like duplicates.
  • The gap remains on files that have no deletion vector, so deletion vectors aren’t the cause.
You can confirm it in Databricks by reading the physical files (for example with boto3 or a Parquet reader) and counting the rows where __DeleteVersion is set. That number should account for the gap. Direct Spark reads of the backing files may be blocked by your workspace’s security settings.

Fix: filter on __DeleteVersion at ingestion

1

Add __DeleteVersion to the Pinot schema

Add the column as a STRING dimension and turn on column-based null handling:
Keep your existing columns. The snippet shows only what to add.
2

Check which value marks a live row

Before you add the filter, ingest a sample with the new column (a test table is fine) and look at its values:
In the tables we’ve seen, live rows come through as 0 and superseded rows carry a non-zero value. The count for 0 should match the row count in Databricks. If your live rows show a different marker, use that value in the next step.
3

Add the ingestion filter

Add a filterConfig to the table’s ingestionConfig:
Pinot drops every row where the filter function returns true, so the expression has to describe the rows you want to discard. __DeleteVersion = 'null' looks reasonable but does the opposite of what you want.
4

Re-ingest the table

The filter applies only to segments built after the change. Segments that already exist still contain the superseded rows. Drop and recreate the table, or delete its segments, and then trigger the DeltaTableIngestionTask again. See Managing Adhoc Minion Tasks for how to trigger it without waiting for the schedule.
5

Compare counts

Once the task finishes, compare SELECT COUNT(*) in Pinot with the count in Databricks. Pinot can briefly lag behind if the source table changed after the task read its snapshot. The next run catches up.

Limitations

  • __DeleteVersion is an internal Databricks column. It isn’t a documented interface, and its name, type, or values could change in a later Databricks release. Recheck the counts after Databricks upgrades your pipeline.
  • The filter removes superseded row versions. It doesn’t recover deletes that Databricks itself never applied. Lakeflow Connect doesn’t propagate Salesforce hard deletes automatically, and it can miss soft deletes that were purged before the next sync. Those need a full refresh of the Databricks table. See Databricks Salesforce connector limitations.
  • If you can’t change the Pinot table, you can instead publish a filtered copy from Databricks (WHERE __DeleteVersion IS NULL) to a separate Delta location and point delta.ingestion.table.uri at it. It costs an extra pipeline step, but Pinot no longer depends on Databricks internals.