> ## Documentation Index
> Fetch the complete documentation index at: https://docs.startree.ai/llms.txt
> Use this file to discover all available pages before exploring further.

# Fix Row Count Mismatches for Databricks CDC Delta Tables

> Filter superseded row versions out of Databricks streaming tables during Delta Lake ingestion, so Pinot row counts match Unity Catalog.

Use this recipe when a Delta table ingested with the [Delta Lake connector](/corecapabilities/ingestdata/adv-concepts/batch/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

<Steps>
  <Step title="Add __DeleteVersion to the Pinot schema">
    Add the column as a `STRING` dimension and turn on column-based null handling:

    ```json theme={null}
    {
      "schemaName": "<TABLE_NAME>",
      "enableColumnBasedNullHandling": true,
      "dimensionFieldSpecs": [
        {
          "name": "__DeleteVersion",
          "dataType": "STRING"
        }
      ]
    }
    ```

    Keep your existing columns. The snippet shows only what to add.
  </Step>

  <Step title="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:

    ```sql theme={null}
    SELECT __DeleteVersion, COUNT(*) AS rows
    FROM <TABLE_NAME>
    GROUP BY __DeleteVersion
    ORDER BY rows DESC
    LIMIT 10
    ```

    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.
  </Step>

  <Step title="Add the ingestion filter">
    Add a `filterConfig` to the table's `ingestionConfig`:

    ```json theme={null}
    "ingestionConfig": {
      "filterConfig": {
        "filterFunction": "__DeleteVersion != '0'"
      },
      "batchIngestionConfig": {
        "segmentIngestionType": "REFRESH"
      }
    }
    ```

    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.
  </Step>

  <Step title="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](/corecapabilities/manage-data/recipes/adhoc-task-trigger) for how to trigger it without waiting for the schedule.
  </Step>

  <Step title="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.
  </Step>
</Steps>

## 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](https://docs.databricks.com/aws/en/ingestion/lakeflow-connect/salesforce-limits#pipelines).
* 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.
