apache / apache/datafusion

Breaking API Change Idea: reserve scalars information between `ExecutionPlan`

Open
#18,308 1 comment 0 reactions 0 assignees View on GitHub
Dominant language
Rust
Stars
9.3k
Forks
2.4k
Avg merge
3d 7h
Merged PRs (30d)
344

Description

Currently, scalar information is lost when data moves between execution plans. While `PhysicalExpr` and `ScalarUDFImpl` work with `ColumnarValue` (enabling optimized implementations for scalars), the stream of `RecordBatch`es between `ExecutionPlan`s doesn't preserve scalar knowledge.

### Proposed Solution

Allow `ExecutionPlan` to return a stream that preserves scalar information (will return a stream of `ReturnedValue`):

```rust
enum ReturnedValue {
Batch(RecordBatch),
BatchWithScalars(RecordBatchWithScalars)
}

/// Same as ColumnarValue but with Arc-wrapped Scalar to avoid unnecessary copying
enum ColumnarValue {
Array(ArrayRef),
Scalar(Arc),
}

struct RecordBatchWithScalars {
schema: SchemaRef,
columns: Vec,
row_count: usize,
}
```

### Benefits

#### Sort Operations
- **Skip sorting**: When a column contains a scalar value (same value for the entire batch), sorting on that column can be skipped
- **Efficient copying**: When sorting by other columns, scalar columns can be copied more efficiently or only partially
- **Produce scalars**: Sort operations can output scalars for faster downstream operations

#### Aggregation
- **Avoid expensive operations**: When grouping by a scalar column, expensive hashing/comparison can be avoided

### Downsides
1. Significant breaking change
2. Increased code complexity

### Potential Extensions

We could add another variant for row-based encoding:
```rust
ReturnedValue::EncodedRows(Rows)
```

This would allow operators that process data as rows to pass encoded rows directly to the next operator, avoiding unnecessary conversions between columnar and row formats. The receiving operator can then decide whether to:
- Use row-based input directly (avoiding conversion overhead if it would benefit from row based as well)
- Convert to columnar format (same cost as current behavior)

Contributor guide

Open the contributing guide

Research direction

Start by tracing the ExecutionPlan stream boundary and the existing ColumnarValue, PhysicalExpr, ScalarUDFImpl, and RecordBatch representations mentioned in the issue. Evaluate the breaking API and potential EncodedRows extension, then define a design that preserves scalar information without unnecessary conversions and has agreement from maintainers.

Written by the indexing model from the issue text.

Assessment

Tech stack
rust
Domain
backend-api-design
Issue type
Feature
Difficulty
5/5
Estimated time
Over a week
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
25/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.