Skip to content

Partitioned INSERT computes partition values from the wrong columns when the input plan ends in a projection #35

Description

@andygrove

Describe the bug

For partitioned tables, insert_into adds a ProjectionExec that computes the _partition column with PartitionExpr. PartitionExpr reads its source columns by position: PartitionValueCalculator::calculate indexes batch.columns() with positions precomputed from the table schema. But PartitionExpr::children() returns an empty list (project.rs#L210-L212).

Because the expression appears to reference no columns, DataFusion's ProjectionPushdown rule merges the partition projection into a ProjectionExec beneath it and leaves PartitionExpr unchanged. PartitionExpr then runs against the input batch of the lower projection, which has a different column layout.

The data columns are still written correctly, but each data file gets partition values computed from another column. Reads that prune by partition then silently skip those files.

To Reproduce

-- pt: id INT NOT NULL, val INT NOT NULL, partitioned by identity(id)
-- src (MemTable): a INT NOT NULL, b INT NOT NULL, rows (1, 100), (2, 200)
INSERT INTO pt SELECT b AS id, a AS val FROM src;

Optimized physical plan:

IcebergCommitExec: table=ns.pt
  IcebergWriteExec: table=ns.pt
    ProjectionExec: expr=[b@1 as id, a@0 as val, iceberg_partition_values(id) as _partition]
      DataSourceExec: partitions=1, partition_sizes=[1]

The partition values come from column 0 of the source batch (a) instead of from id (b):

  • data files are written under data/id=1 and data/id=2 instead of data/id=100 and data/id=200
  • SELECT * FROM pt returns (100, 1) and (200, 2)
  • SELECT * FROM pt WHERE id = 100 returns no rows

Expected behavior

Partition values are computed from the columns being written, however the optimizer rearranges the projections.

Additional context

Any INSERT whose physical input ends in a ProjectionExec that renames, reorders or drops columns is exposed (column aliases, inserts from joins, ...). VALUES and SELECT * from a source with the table's column order work only because the positions happen to line up. The schema check in project_with_partition runs before physical optimization, so it doesn't catch this.

Possible fixes:

  • expose the source columns as Column children, so projection rewrites update them
  • check the batch schema in evaluate and resolve source columns by name or field id
  • compute partition values inside IcebergWriteExec instead of in a PhysicalExpr

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    bugSomething isn't working

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions