Is your feature request related to a problem or challenge?
A distributed query engine plans on one node and executes on others, so physical
plans must survive serialization. project_with_partition injects a
PartitionExpr into the projection for partitioned writes, and that expression
cannot make the trip.
It wraps a PartitionValueCalculator, which is live non-serializable state. The two
things that are serializable (the PartitionSpec and the table schema) are
consumed at construction and not retained, so nothing on the expression says
what produced it. Exposing the calculator would not help either: its public API is
try_new, partition_type(), partition_arrow_type() and calculate(), with no
accessor for the spec or schema it was built from.
The spec alone is not enough to rebuild the calculator: PartitionSpec stores
only spec_id and fields, referring to columns by source_id. The schema is
what resolves those ids to real columns and determines the partition type.
Nor can the codec supply the missing pieces from elsewhere. In DataFusion 55,
PhysicalExtensionCodec::try_encode_expr receives the expression and a
PhysicalExprEncodeCtx, and that context holds only the child encoder — no
schema, no table, no catalog. Encoding is where the spec and schema have to be
read out of the expression, and there is nothing else to consult, so an
expression that does not carry them cannot be written down in the first place.
The decode side has a little more to work with — PhysicalExprDecodeCtx::schema()
exposes the input schema — but it does not close the gap: that is an
arrow::datatypes::Schema, whereas PartitionValueCalculator::try_new needs an
iceberg::spec::Schema, and no context supplies a PartitionSpec at all.
Describe the solution you'd like
Retain both inputs on PartitionExpr and expose them:
try_new(partition_spec, table_schema) replacing the private
new(calculator, spec), constructing the calculator internally so the
retained inputs and the calculator cannot drift apart.
partition_spec() and table_schema() accessors.
Additive: project_with_partition keeps its signature, and new was private,
so no existing caller changes.
Additional context
Originally filed against apache/iceberg-rust as
apache/iceberg-rust#3002, before the DataFusion
integration moved to this repository.
Willingness to contribute
I can contribute to this feature independently.
Is your feature request related to a problem or challenge?
A distributed query engine plans on one node and executes on others, so physical
plans must survive serialization.
project_with_partitioninjects aPartitionExprinto the projection for partitioned writes, and that expressioncannot make the trip.
It wraps a
PartitionValueCalculator, which is live non-serializable state. The twothings that are serializable (the
PartitionSpecand the table schema) areconsumed at construction and not retained, so nothing on the expression says
what produced it. Exposing the calculator would not help either: its public API is
try_new,partition_type(),partition_arrow_type()andcalculate(), with noaccessor for the spec or schema it was built from.
The spec alone is not enough to rebuild the calculator:
PartitionSpecstoresonly
spec_idandfields, referring to columns bysource_id. The schema iswhat resolves those ids to real columns and determines the partition type.
Nor can the codec supply the missing pieces from elsewhere. In DataFusion 55,
PhysicalExtensionCodec::try_encode_exprreceives the expression and aPhysicalExprEncodeCtx, and that context holds only the child encoder — noschema, no table, no catalog. Encoding is where the spec and schema have to be
read out of the expression, and there is nothing else to consult, so an
expression that does not carry them cannot be written down in the first place.
The decode side has a little more to work with —
PhysicalExprDecodeCtx::schema()exposes the input schema — but it does not close the gap: that is an
arrow::datatypes::Schema, whereasPartitionValueCalculator::try_newneeds aniceberg::spec::Schema, and no context supplies aPartitionSpecat all.Describe the solution you'd like
Retain both inputs on
PartitionExprand expose them:try_new(partition_spec, table_schema)replacing the privatenew(calculator, spec), constructing the calculator internally so theretained inputs and the calculator cannot drift apart.
partition_spec()andtable_schema()accessors.Additive:
project_with_partitionkeeps its signature, andnewwas private,so no existing caller changes.
Additional context
Originally filed against
apache/iceberg-rustasapache/iceberg-rust#3002, before the DataFusion
integration moved to this repository.
Willingness to contribute
I can contribute to this feature independently.