These commits are when the Protocol Buffers files have changed: (only the last 100 relevant commits are shown)
| Commit: | 7b1e55d | |
|---|---|---|
| Author: | Anurag Tryambak Raut | |
| Committer: | GitHub | |
fix: preserve dotted column relation qualifiers (#25334) ## Which issue does this PR close? * Closes #24776. ## Rationale for this change `ColumnRelation` currently round-trips `TableReference` qualifiers through a single dotted string. When a catalog, schema, or table identifier itself contains `.`, the qualifier boundaries are lost during deserialization, causing the reconstructed `TableReference` to differ from the original. For example, a schema named `my.schema` can be incorrectly interpreted as separate qualifier components after a protobuf round-trip. ## What changes are included in this PR? * Add a structured `parts` field to `ColumnRelation`. * Serialize `TableReference` components using `TableReference::to_vec()`. * Prefer the structured `parts` representation when deserializing. * Fall back to the existing `relation` string for legacy protobuf messages. * Add round-trip tests for dotted bare, partial, and full table references. * Add a test covering decoding of legacy relation-only messages. * Regenerate the protobuf bindings. ## What is the testing strategy for this PR? Added regression tests covering the `Column -> protobuf -> TableReference` round-trip for: * Bare table references containing `.` * Partial references with dotted schema identifiers * Full references with dotted schema identifiers * Legacy `relation`-only protobuf messages Also ran: `cargo test -p datafusion-proto-common column_relation -- --nocapture` `cargo test -p datafusion-proto-common` All tests pass. ## Are there any user-facing changes? No.
The documentation is generated from this commit.
| Commit: | 532cd03 | |
|---|---|---|
| Author: | Liang-Chi Hsieh | |
| Committer: | GitHub | |
feat: support multi-column NOT IN subqueries with null-aware joins (#19857) ## Which issue does this PR close? - Closes #25737. Follow-up to #10583, which made single-column `NOT IN` a null-aware anti join. ## Rationale for this change Multi-column `IN` / `NOT IN` subqueries fail to plan: ```sql SELECT * FROM t1 WHERE (a, b) NOT IN (SELECT x, y FROM t2); SELECT * FROM t1 WHERE (c2, c3) NOT IN ( SELECT c2, c3 FROM t2 WHERE t1.c1 = t2.c1 ); ``` Supporting them requires SQL three-valued logic over the whole tuple. A NULL in one element does not make every comparison UNKNOWN the way a scalar NULL does: `(NULL, 8) = (1, 2)` is `UNKNOWN AND FALSE`, which is FALSE, so the outer row is kept. `(NULL, 2) = (1, 2)` is UNKNOWN, so it is dropped. ## What changes are included in this PR? A null-aware join already reads its equi-join keys by position: `on[0]` is the `NOT IN` value key and `on[1..]` are the correlation scope keys of a correlated subquery. A tuple needs several value keys, so both the logical `Join` and `HashJoinExec` now record how many leading keys are value keys (`null_aware_value_keys`). The layout becomes `on[..V]` for the tuple elements and `on[V..]` for the correlation keys. - **SQL planning / analyzer**: `(a, b) IN (SELECT x, y ...)` is accepted when the tuple and subquery column counts match. Type coercion and placeholder inference work element by element. A tuple compared with a single struct-typed column (`(a, b) IN (SELECT s ...)`) is still a scalar struct comparison. - **Decorrelation**: emits one equality per tuple element, ahead of the correlation predicates, so the elements become the leading keys in order. Constant elements are projected as outer columns, the same way the scalar constant case already is. A tuple `IN` / `NOT IN` inside a projection or an `OR` goes through a null-aware mark join. - **`HashJoinExec`**: `V > 1` uses the existing per-build-row path of correlated null-aware joins. A row counts as NULL-valued when any value key is NULL. Candidate pairs are then kept only when their value tuples have no definite mismatch, i.e. every element pair is equal or involves a NULL. This works for both `LeftAnti` and `LeftMark`, with or without scope keys and join filters. - **Dynamic filter pushdown**: probe rows with a NULL in any value key are kept. - **Protobuf**: `JoinNode` and `HashJoinExecNode` carry `null_aware_value_keys`. A missing field (`0`) decodes as `1`. ## What is the testing strategy for this PR? - `null_aware_anti_join.slt` covers: - 2- and 3-column `NOT IN` with NULLs on either side, an empty subquery, and correlated equality and non-equality correlations - constant tuple elements - `NOT IN` as a nullable projected value and under `OR` - per-element type coercion, column count errors, a single struct column, and placeholders - `subquery.slt` covers multi-column `IN`. - `hash_join/exec.rs` unit tests run `LeftAnti` and `LeftMark` with probe-side and build-side NULLs, three columns, an empty probe side, and correlated scope keys, across batch sizes. Further unit tests cover value-key validation and the protobuf round trip. `shared_bounds.rs` tests the dynamic filter NULL escape over every value key. - `roundtrip_logical_plan.rs` round-trips a correlated multi-column `NOT IN` plan. ## Are there any user-facing changes? Multi-column `IN` / `NOT IN` subqueries are now supported. API change: `Join` and `HashJoinExec` have a new public `null_aware_value_keys` field, and the generated protobuf `JoinNode` and `HashJoinExecNode` have a new field. Code that builds these with struct literals must set it (`1` keeps the previous behavior). The 56.0.0 upgrade guide describes the change.
| Commit: | d9a9321 | |
|---|---|---|
| Author: | Peter Lee | |
| Committer: | GitHub | |
fix: preserve literal metadata and destructure leaf and unary proto hooks (#25382) ## Which issue does this PR close? Closes #24613. Literal field metadata is lost when physical expressions are serialized. ## Rationale for this change Losing this metadata can strip Arrow extension types and custom metadata from projected columns. ## What changes are included in this PR? Use exhaustive destructuring in all seven leaf and unary expression encoders and decoders. Add a literal protobuf message that carries metadata and regenerate the Rust and JSON bindings. ## What is the testing strategy for this PR? A projection round-trip test covers null and non-null literals through protobuf and JSON and was verified to fail before the fix. Existing tests cover plain literals, and a new test rejects messages missing a value. The extended workspace tests and all-feature Clippy pass. ## Are there any user-facing changes? Literal metadata survives physical plan serialization. Literals with metadata use a new protobuf variant that requires an updated reader.
| Commit: | 5e5d79d | |
|---|---|---|
| Author: | RIchard Baah | |
| Committer: | GitHub | |
introduce optional rle reads from parquet (#24227) ## a large portion of this PR is test! ## Which issue does this PR close? <!-- We generally require a GitHub issue to be filed for all bug fixes and enhancements and this helps us generate change logs for our releases. You can link an issue to this PR using the GitHub syntax. For example `Closes #123` indicates that this PR will close issue #123. --> - Closes #24112. - part of #24111 ## Rationale for this change When DataFusion reads a parquet file with dictionary-encoded string or binary columns, it currently decodes the dictionary and returns plain Utf8/Binary arrays, discarding the encoding. For low-cardinality columns (status, country, category, etc.) this doesn't take full advantage of the compacted format parquet gives the engine Preserving the dictionary encoding reduces memory usage and can improve aggregation performance on these columns. <!-- Why are you proposing this change? If this is already explained clearly in the issue then this section is not needed. Explaining clearly why changes are proposed helps reviewers understand your changes and offer better suggestions for fixes. Please explain the problem you are trying to solve in terms of the user-visible behavior, rather than the implementation. For example, "The code in `foo.rs` doesn't handle nulls" is a symptom of the implementation. "COUNT(DISTINCT) returns wrong results when the column contains nulls" is the user-visible problem. --> ## What changes are included in this PR? Adds `datafusion.execution.parquet.enable_rle_to_dictionary` with a default of `false`. When enabled for inferred-schema Parquet tables, DataFusion inspects Parquet footer metadata and promotes top-level string/binary columns with dictionary pages to Arrow dictionary types. Mixed dictionary/plain files are normalized before schema merge when the value types are compatible. At scan time, the parquet opener passes the promoted schema to arrow-rs so those columns can be read as dictionary arrays directly. *Tables with a user-supplied schema are not promoted because DataFusion does not use footer metadata to infer their schema.* <!-- There is no need to duplicate the description in the issue here, but it is sometimes worth providing a summary of the individual changes in this PR. --> ## Are these changes tested? yes. - `datafusion/sqllogictest/test_files/parquet_rle_to_dictionary.slt` - `datafusion/datasource-parquet/src/schema_coercion.rs` - `uniform_dict_schemas_respects_value_type_families` checks that mixed file schemas are normalized only across compatible string/binary families. - `rle_schema_coercion_respects_dictionary_value_type` checks scan-time coercion into dictionary types, including incompatible cases and the flag-off path. - `datafusion/datasource-parquet/src/opener/mod.rs` - `test_rle_binary_column_promotion` verifies the opener passes a promoted `Dictionary(Int32, Binary)` schema to arrow-rs so binary columns can be read as dictionary arrays directly. <!-- We typically require tests for all PRs in order to: 1. Prevent the code from being accidentally broken by subsequent changes 2. Serve as another way to document the expected behavior of the code If tests are not included in your PR, please explain why (for example, are they covered by existing tests)? --> ## Are there any user-facing changes? New session config option: `SET datafusion.execution.parquet.enable_rle_to_dictionary = true`. Default is false so existing behavior is unchanged. <!-- If there are user-facing changes then we may require documentation to be updated before approving the PR. If there are any breaking changes to public APIs, please add the `api change` label. -->
| Commit: | 1b7f403 | |
|---|---|---|
| Author: | Oleks V | |
| Committer: | GitHub | |
feat: Add option to assign Parquet row groups to file ranges by midpoint (#25940) ## Which issue does this PR close? - Closes #25938. ## Rationale for this change When a Parquet file is split into byte ranges, DataFusion reads a row group in the range that holds its first page, while Spark reads it in the range that holds its midpoint. Engines that hand Spark's splits to DataFusion, like Comet, therefore read row groups in different partitions than Spark planned. That leaves scan tasks idle (apache/datafusion-comet#3817) and reports the wrong split in `_metadata.file_block_*` (apache/datafusion-comet#6512). This PR adds an opt-in rule that matches Spark and keeps the current rule as the default. ## What changes are included in this PR? - New option `datafusion.execution.parquet.row_group_range_assignment`: `start_offset` (default, current behavior) or `midpoint`. It is the `RowGroupRangeAssignment` enum in `datafusion_common::parquet_config`, settable with `SET`, as a `format.` table option, or through `TableParquetOptions`. - `row_group_in_range` takes the assignment. `midpoint` follows parquet-java's `filterFileMetaDataByMidpoint`: the smaller of column 0's dictionary and data page offsets, plus half of the row group's compressed size. Range pruning and the `bytes_processed` credit both use it. - `RowGroupAccessPlanFilter::prune_by_range` is unchanged. - Proto: new `row_group_range_assignment` field in `ParquetOptions`. An empty value decodes as the default, so older plans still load. - Docs: `configs.md`, plus the `FileRange` and `ByteProgress` docs, which described the old rule. ## What is the testing strategy for this PR? - `row_group_filter.rs`: the layout from apache/datafusion-comet#3817 under both rules, every two-way split of a file reading each row group once, and the dictionary offset case. - `opener/mod.rs`: `each_range_of_a_split_file_credits_only_its_own_bytes` now runs under both rules and checks each range's rows and `bytes_processed`. - Proto round trips, `information_schema.slt`, and a `repartition_scan.slt` case that reads a multi-row-group file split into 4 ranges under both rules. ## Are there any user-facing changes? A new config option. Nothing changes unless it is set to `midpoint`. `ParquetOptions` gains a public field.
| Commit: | 54e32bb | |
|---|---|---|
| Author: | Martin Grigorov | |
| Committer: | GitHub | |
feat: Push down offset/skip to TableScan (#25404) ## Which issue does this PR close? - Closes #5562 ## Rationale for this change Currently `TableProvider` provides the `limit` argument to the `scan()` method to request a maximum number of rows. There is no way to tell the implementation to skip some of the rows, for example to fetch the second/third/Nth page of rows (i.e. SQL `... LIMIT 20 OFFSET 40`). Adding an additional field to ScanArgs (named `offset` or `skip`) will make it possible for implementations to override the `scan_with_args()` method and optimize their scan to read and return only the requested rows. ## What changes are included in this PR? * A new field named `offset` is added to `ScanArgs`, with a setter and a getter. * A new method is added to the `TableProvider` trait - `supports_offset_pushdown() -> bool`. By default it returns `false` but any implementation that can support skipping of rows could override it to return `true` and combined with a custom implementation of `scan_with_args()` to optimise its data scan/read. * Update some callers of `TableProvider::scan()` to use `::scan_with_args()` where they could support offset push down * Update the migration guide for 56.0.0 with a section about the offset pushdown support Note: `datafusion-ffi` is **not** updated because it does not expose `scan_with_args()` yet. ## What is the testing strategy for this PR? New unit tests are added for the implementations which support offset pushdown. ## Are there any user-facing changes? The new functionality is opt-in! All currently existing custom implementations of `TableProvider` trait will continue to compile and run without any modifications. Any custom implementation that wants to make use of the new functionality will need to override `TableProvider::supports_offset_pushdown()` to return `true` and make use of `ScanArgs::offset` in its `scan_with_args()` implementation. --------- Signed-off-by: dependabot[bot] <support@github.com> Co-authored-by: Jayant Shrivastava <jayant.shrivastava@datadoghq.com> Co-authored-by: Gabriel <45515538+gabotechs@users.noreply.github.com> Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com> Co-authored-by: Kumar Ujjawal <ujjawalpathak6@gmail.com> Co-authored-by: Dmitri B. <s5dsn-eqee.dev@proton.me> Co-authored-by: Neil Conway <neil.conway@gmail.com> Co-authored-by: Zeren Wang <53075619+Vanzeren@users.noreply.github.com> Co-authored-by: Adrian Garcia Badaracco <1755071+adriangb@users.noreply.github.com> Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com> Co-authored-by: Abhishek <166850903+Abhisheklearn12@users.noreply.github.com> Co-authored-by: Andy Grove <agrove@apache.org> Co-authored-by: xudong.w <wxd963996380@gmail.com>
| Commit: | b73054a | |
|---|---|---|
| Author: | Goutam Adwant | |
| Committer: | GitHub | |
feat: allow scaling RangePartitioning (#24766) ## Which issue does this PR close? - Closes #24712. ## Rationale for this change Distributed and adaptive planners need to change a range-partitioned plan's task count without losing its higher-resolution boundary sample. The optimizer should use a compatible range layout where possible while preserving the requested parallelism when sample capacity is insufficient. ## What changes are included in this PR? - Add validated sample-backed construction, sample/current/maximum partition-count accessors, and `scale(usize) -> Option<Self>`. Zero and above-capacity requests return `None`. - Retain samples across scaling, projection, repartition, join rewrites, protobuf, and FFI round trips. Keep structural equality for interleave safety and use `has_same_layout` for effective-layout comparisons. - Integrate scaling into `EnsureRequirements`, recheck key requirements after singleton scaling, and fall back to key partitioning at the requested parallelism when scaling is unsupported. During co-partitioning, a source already providing the required layout remains native; a newly inserted range exchange does not gain that preference. - Let `RepartitionExec::repartitioned` scale supported ranges with fresh runtime state and metrics, retaining ordering and batch-size settings. Remove the optimizer's redundant outer `Option`. - Deprecate unchecked `new` for 56.0.0, retain validated exact construction and `split_points()`, and document sample sizing and upgrade requirements. - Preserve effective boundaries in protobuf field 2 for older readers, and validate the added sample/count representation against those boundaries. ## What is the testing strategy for this PR? Coverage includes count limits, sample retention, key and sort-option compatibility, structural versus effective-layout equality, interleave safety, sample-preserving rewrites, protobuf/FFI round trips, and aggregate/join results with NULLs. The new exchange regression fails before the follow-up fix and passes afterward. It initializes the original exchange, scales 2 -> 4 -> 1 -> 4, and checks exact row routing, ordering, retained samples, batch configuration, and independent execution state and metrics. The 23 range optimizer integration regressions also pass. A separate join regression fails before the native-reference fix and passes afterward. It checks both input orders and verifies that a larger, newly scaled range exchange does not displace a source already providing the required hash layout. The end-to-end range join test now checks the resulting hash layout and exact output rows when scaling requires exchanges on both sides. Validation with the pinned Rust 1.98.1 toolchain: - Formatting, all-target/all-feature Clippy with warnings denied, and the complete `./dev/rust_lint.sh` suite pass. - The extended workspace suite passes 12,022 Rust tests, with eight ignored, and all 522 SQL logic-test files. - All 155 all-feature FFI tests pass. Performance caveat: six existing range benchmarks were compared with base `925d7f8ffd` using separate Rust 1.98.1 release-nonlto builds and alternating runs. The eight-partition integer routing case was 13.5% slower; the other five differed by -1.1% to +2.2%. A separate before/after check of that case found this follow-up within 0.3% of pre-follow-up head `1d2d59e699`, with both about 18% slower than the base in that run. This indicates an existing full-PR performance concern, not a measured slowdown introduced by this follow-up. Its cause remains unverified; these local results do not establish performance parity with main. ## Are there any user-facing changes? Yes. `RangePartitioning` gains sample-backed construction and scaling APIs. `scale` returns `Option<Self>` rather than a public error type. Unchecked `new` is deprecated; validated exact construction remains supported. Unsupported optimizer scaling falls back to key partitioning rather than retaining lower parallelism. `FFI_RangePartitioning` changes layout and the generated `PhysicalRangePartitioning` protobuf struct gains fields. These API/ABI changes require the `api change` label. This targets main, not a patch release. FFI consumers must rebuild against the compatible release. Existing protobuf payloads remain readable. --------- Signed-off-by: goutamadwant <workwithgoutam@gmail.com>
| Commit: | 714956b | |
|---|---|---|
| Author: | Peter Lee | |
| Committer: | GitHub | |
fix: preserve overflow behavior and exhaustively destructure expression proto hooks (#25034) ## Which issue does this PR close? Closes #24614. ## Rationale for this change This PR fixes serialization losing the setting that makes arithmetic fail on overflow. For example, checked `Int32::MAX + 1` raises an error before serialization but returns `Int32::MIN` after decoding. Preserving this setting ensures that sending an expression through protobuf preserves its arithmetic behavior. ## What changes are included in this PR? The protobuf message now stores `BinaryExpr::fail_on_overflow`, and both supported decoding formats restore it. When the encoder combines nested expressions into a flat list of operands, it requires their operators and overflow settings to match. This preserves the behavior of expressions that mix checked and wrapping arithmetic. The generated Rust and JSON bindings include the new field. All six encoding and decoding hooks for `BinaryExpr`, `LikeExpr`, and `SqlSimilarToPattern` explicitly list every field without a rest pattern. Adding a field to an expression or its protobuf payload will cause a compile error until the corresponding hook handles it. The PR also changes the PostgreSQL SQLLogicTest decimal formatter to borrow its argument, resolving an existing Clippy error that blocked the required checks before committing. ## What is the testing strategy for this PR? The new tests serialize and decode nested additions, then check their evaluated results for all four combinations of checked and wrapping arithmetic. The regression test failed before the fix because an expression that should raise an overflow error returned `Int32(-2147483648)`. Additional tests cover older messages that omit the new field and verify that JSON preserves the overflow setting. All 17 focused expression tests passed. The extended workspace run passed 11,260 Rust tests, with 8 ignored, and completed all 511 SQLLogicTest files. Both conversion tests passed with the PostgreSQL feature enabled. Formatting, Clippy with all targets and features, and the complete `./dev/rust_lint.sh` suite also passed. ## Are there any user-facing changes? Expressions configured to fail on arithmetic overflow now raise the expected error after serialization and decoding, including nested expressions with different overflow settings. --------- Co-authored-by: Andrew Lamb <andrew@nerdnetworks.org>
| Commit: | f0eaf4b | |
|---|---|---|
| Author: | Peter Lee | |
| Committer: | GitHub | |
fix: preserve overflow behavior and exhaustively destructure expression proto hooks (#25034) ## Which issue does this PR close? Closes #24614. ## Rationale for this change This PR fixes serialization losing the setting that makes arithmetic fail on overflow. For example, checked `Int32::MAX + 1` raises an error before serialization but returns `Int32::MIN` after decoding. Preserving this setting ensures that sending an expression through protobuf preserves its arithmetic behavior. ## What changes are included in this PR? The protobuf message now stores `BinaryExpr::fail_on_overflow`, and both supported decoding formats restore it. When the encoder combines nested expressions into a flat list of operands, it requires their operators and overflow settings to match. This preserves the behavior of expressions that mix checked and wrapping arithmetic. The generated Rust and JSON bindings include the new field. All six encoding and decoding hooks for `BinaryExpr`, `LikeExpr`, and `SqlSimilarToPattern` explicitly list every field without a rest pattern. Adding a field to an expression or its protobuf payload will cause a compile error until the corresponding hook handles it. The PR also changes the PostgreSQL SQLLogicTest decimal formatter to borrow its argument, resolving an existing Clippy error that blocked the required checks before committing. ## What is the testing strategy for this PR? The new tests serialize and decode nested additions, then check their evaluated results for all four combinations of checked and wrapping arithmetic. The regression test failed before the fix because an expression that should raise an overflow error returned `Int32(-2147483648)`. Additional tests cover older messages that omit the new field and verify that JSON preserves the overflow setting. All 17 focused expression tests passed. The extended workspace run passed 11,260 Rust tests, with 8 ignored, and completed all 511 SQLLogicTest files. Both conversion tests passed with the PostgreSQL feature enabled. Formatting, Clippy with all targets and features, and the complete `./dev/rust_lint.sh` suite also passed. ## Are there any user-facing changes? Expressions configured to fail on arithmetic overflow now raise the expected error after serialization and decoding, including nested expressions with different overflow settings.
| Commit: | 3b16a3d | |
|---|---|---|
| Author: | Xuanyi Li | |
| Committer: | GitHub | |
fix: preserve MERGE target qualifier bindings (#24429) ## Which issue does this PR close? - Part of #20746. - Follow-up to #22988. ## Rationale for this change #22988 deliberately rejected two valid MERGE forms to avoid silently changing expression meaning. ### Limitation 1: target-correlated subquery with an aliased target ```sql MERGE INTO target AS t USING source AS s ON EXISTS (SELECT 1 FROM source AS x WHERE x.id = t.id) WHEN MATCHED THEN DELETE; ``` `t.id` inside the subquery becomes `outer_ref(t.id)`. The old top-level canonicalizer only rewrote `Expr::Column(t.id)` to `target.id`, leaving the outer reference inconsistent with the schema later rebuilt from `DmlStatement.table_name`. ### Limitation 2: source qualifier equals the target table name ```sql MERGE INTO target AS t USING source AS target ON t.id = target.id WHEN MATCHED THEN DELETE; ``` Canonicalizing target `t.id` to `target.id` collapses both operands onto the source qualifier and can turn the condition into `target.id = target.id`. A recursive string rewrite is not safe: qualifiers are scope-local, so an inner relation can legally shadow `t`. It also cannot solve the second limitation because both relations would still have the same qualifier after rewriting. ## What changes are included in this PR? ### Solution Keep provider identity and the SQL-visible target qualifier as separate plan state: ```text SQL: MERGE INTO target AS t USING source AS target SQL binding +-------------------+ | | target table source relation provider: target qualifier: target qualifier: t | | +---------+---------+ | MERGE expression schema +-----------+-----------+ | t.* | target.* | | index 0.. | index N.. | +-----------+-----------+ | TableProvider::merge_into Logical / protobuf representation DmlStatement.table_name / DmlNode.table_name = target (provider identity) MergeIntoOp.target_qualifier / MergeIntoOpNode.target_qualifier = t (SQL-visible binding) ``` Planning now proceeds as follows: 1. Resolve `target` to the target provider, while retaining alias `t` as the visible qualifier. 2. Store `t` on `MergeIntoOp`; do not canonicalize target columns or recursively rewrite subquery plans. 3. Build one target-first expression schema (`t.*`, then source fields) and use it consistently in SQL planning, analyzer/optimizer rules, physical planning, and programmatic plans. 4. Reject only a source qualifier that collides with the *visible* target qualifier in the outer MERGE scope. A source qualifier matching the real table name remains valid when the target is aliased. 5. Pass the preserved schema and expressions to `TableProvider::merge_into`, where target and source columns resolve to distinct physical indices. ### Protobuf change and 55.0/56.0 compatibility `MergeIntoOpNode` gains optional `target_qualifier` field 3. This field is necessary because `DmlNode.table_name` contains provider identity and cannot also represent alias `t`; without it, a proto round trip loses the binding needed to resolve MERGE expressions. Compatibility is directional: - **55.0 payload -> 56.0 reader: supported.** The field is absent, so the new reader falls back to `DmlNode.table_name`, matching 55.0's canonicalized representation. - **56.0 payload -> 55.0 reader: unsupported.** A 55.0 reader ignores the unknown field but alias-preserving expressions still require it, so expressions can fail resolution or be misbound. ### Rust API compatibility DataFusion 55.0 released `MergeIntoOp` with public struct-literal construction: `MergeIntoOp { on, clauses }`. This PR makes the struct non-exhaustive, adds private target-qualifier state, and requires `MergeIntoOp::new(target_qualifier, on, clauses)`. Therefore existing 55.0 downstream struct literals will not compile unchanged against 56.0. This breaking change is documented in the 56.0 upgrade guide and is appropriate for the next major release. The PR also: - removes alias canonicalization and the recursive target-correlation guard; - documents that providers may receive residual subqueries; and - adds coverage for direct correlations, nested shadowing, lateral and LIMIT scopes, qualifier collisions, quoted/qualified identifiers, proto fallback/round trips, and physical column indices. ## Are these changes tested? - `cargo fmt --all` - `cargo clippy --all-targets --all-features -- -D warnings` - `./ci/scripts/doc_prettier_check.sh --write --allow-dirty` - `RUST_BACKTRACE=1 cargo test --profile ci --exclude datafusion-examples --exclude datafusion-benchmarks --exclude datafusion-cli --workspace --lib --tests --bins --features avro,json,backtrace,extended_tests,recursive_protection,parquet_encryption` ## Are there any user-facing changes? Yes. Both valid MERGE forms above now plan successfully and reach `TableProvider::merge_into`. The Rust and protobuf compatibility constraints for upgrading from 55.0 to 56.0 are documented above. --------- Co-authored-by: kosiew <29057562+kosiew@users.noreply.github.com>
| Commit: | 768fc8a | |
|---|---|---|
| Author: | Burak Şen | |
| Committer: | GitHub | |
fix(proto): preserve Parquet source and sink state (#25057) ## Which issue does this PR close? - Closes #24620. ## Rationale for this change Parquet source and sink protobuf hooks accessed fields individually and silently dropped source metadata prefetch hints and sink sorting-column metadata. This changed performance characteristics and omitted metadata from files written after a plan round trip. ## What changes are included in this PR? - Exhaustively destructure `ParquetSource`, `ParquetSink`, and their protobuf messages. - Preserve `ParquetSource::metadata_size_hint` with a checked integer conversion. - Preserve `ParquetSink::sorting_columns`, including the distinction between `None` and an empty list. - Explicitly document source/sink fields that are reconstructed or runtime-only. - Document the generated Rust API migration in the DataFusion 56.0.0 upgrade guide. #24930 already preserves `reverse_row_groups` and `sort_order_for_reorder`; this PR builds on that behavior rather than duplicating it. The protobuf wire changes are additive and backward compatible. ## What is the testing strategy for this PR? Physical-plan round-trip coverage verifies: - Metadata-size hints preserve `None`, zero, a normal value, and `usize::MAX`. - Parquet sorting columns preserve `None`, an empty list, and populated lists. - Both fields survive binary protobuf and public JSON API round trips. Both regression tests were ablation-checked and fail when their corresponding field is omitted from serialization. Validated with: - `cargo test -p datafusion-proto --test proto_integration --features json` (262 passed) - `cargo clippy -p datafusion-datasource-parquet -p datafusion-proto --all-features --tests -- -D warnings` - `cargo fmt --all -- --check` - `./ci/scripts/doc_prettier_check.sh --write --allow-dirty` ## Are there any user-facing changes? Parquet scan metadata prefetch hints and sink sorting-column metadata now survive protobuf plan round trips. The protobuf wire additions are backward compatible. Adding `metadata_size_hint` to the generated public `ParquetScanExecNode` Rust struct and `sorting_columns` to `ParquetSink` is a source-level API change for callers using exhaustive struct literals. Such callers must initialize the new fields (use `None` to retain previous behavior) or use `..Default::default()`. This migration is documented in the DataFusion 56.0.0 upgrade guide. --------- Co-authored-by: Kumar Ujjawal <ujjawalpathak6@gmail.com>
| Commit: | 2cb6b54 | |
|---|---|---|
| Author: | Burak Şen | |
| Committer: | GitHub | |
fix(proto): preserve CSV sink writer options (#25058) ## Which issue does this PR close? - Closes #24623. ## Rationale for this change The CSV sink protobuf hooks accessed state indirectly, allowing new fields to be omitted silently. Auditing the writer options found that compression level, timezone-aware timestamp format, and line terminator were not represented on the wire, so those settings changed after a physical-plan round trip. ## What changes are included in this PR? - Exhaustively destructure `CsvSink`, `CsvSinkExecNode`, and the nested CSV sink/writer option messages. - Add protobuf fields for compression level, timezone-aware timestamp format, and line terminator. - Make all five optional date/time formats presence-aware so `None` and `Some("")` remain distinct. - Add a DataFusion 56.0 upgrade guide for the generated `CsvWriterOptions` API changes. - Regenerate the protobuf models. The protobuf wire format remains backward compatible. The generated Rust API changes are documented in the upgrade guide. ## What is the testing strategy for this PR? Expanded the CSV sink round-trip test to cover every writer option, optional defaults, explicitly empty formats, and malformed line terminators. Validated with: - `cargo test -p datafusion-proto --test proto_integration roundtrip_csv_sink` - `cargo clippy -p datafusion-datasource-csv -p datafusion-proto-common -p datafusion-proto --all-features --tests -- -D warnings` - `cargo fmt --all --check` ## Are there any user-facing changes? CSV sinks now retain compression level, timezone-aware timestamp formatting, line terminators, and explicitly empty date/time formats across protobuf plan round trips. Invalid wire terminators return an error. The generated public `CsvWriterOptions` gains three fields, and its existing date/time format fields change from `String` to `Option<String>`. Migration is documented in the DataFusion 56.0 upgrade guide. The protobuf wire format remains backward compatible. --------- Co-authored-by: Kumar Ujjawal <ujjawalpathak6@gmail.com>
| Commit: | 82335b4 | |
|---|---|---|
| Author: | Xuanwo | |
| Committer: | GitHub | |
feat: serialize ASOF join plans (#23832) ## Which issue does this PR close? - Part of #318. - Umbrella PR: #23738. - Follows #23829 (merged). ## Rationale for this change This is the serialization layer of the ASOF JOIN stack. It gives logical and physical ASOF plans explicit protobuf representations without coupling wire format review to SQL or DataFrame APIs. #23829 and its prerequisites are merged. After restacking onto the current `main`, this PR now contains only the serialization layer. ## What changes are included in this PR? - Add protobuf messages and enum values for logical and physical ASOF joins. - Encode and decode equality keys, ordered match expressions, match direction, join constraint, and right output indices. - Implement physical serialization through `ExecutionPlan::try_to_proto` and `AsOfJoinExec::try_from_proto`, following the current self-serializing execution-plan pattern. - Regenerate `prost` and `pbjson` sources with the repository generator. The physical-plan oneof uses the next append-only tag after `PiecewiseMergeJoinExec`. - Add logical round trips for all four match directions and a physical `AsOfJoinExec` round trip. ## Are these changes tested? Yes: - `./datafusion/proto-models/regen.sh` - `cargo fmt --all` - `cargo clippy --all-targets --all-features -- -D warnings` - `./dev/rust_lint.sh` - `cargo test -p datafusion-proto roundtrip_asof_join --all-features` - The extended workspace test command from the contributor guide ## Are there any user-facing changes? Logical and physical ASOF join plans can be serialized through `datafusion-proto`. The additions use new messages and append-only oneof/enum tags, so existing wire tags are not reused. Generated public Rust enums gain new variants, however, so downstream exhaustive matches must add arms; this is a Rust source-compatibility break even though the wire additions are compatible. This PR can now be reviewed independently. It does not depend on the optional floating-point follow-up #24375.
| Commit: | 92746a9 | |
|---|---|---|
| Author: | Burak Şen | |
| Committer: | GitHub | |
refactor(proto): destructure FileScanConfig and MemorySourceConfig serde hooks (#24813) ## Which issue does this PR close? - Closes #24624. ## Rationale for this change Proto hooks that access fields individually can silently omit newly added state. Exhaustive destructuring makes such omissions compile errors. `FileScanConfig` is shared by all file sources, so a field omitted here is lost for Parquet, CSV, JSON, Arrow, and Avro. ## What changes are included in this PR? - Exhaustively destructure `FileScanConfig` and `MemorySourceConfig` in their encoders. - Exhaustively destructure their protobuf nodes in the decoders. - Document fields that are serialized indirectly or reconstructed during decoding. - Add an optional `preserve_order` field so explicit values survive file scan round trips while older payloads retain their ordering-derived behavior. - Reject default serialization of non-default expression adapter factories instead of silently discarding their behavior, and update the custom serialization example accordingly. - Decode projected `MemorySourceConfig` sort information against the projected schema without applying projection twice. The protobuf change is additive and backward compatible. ## Are these changes tested? Added or updated coverage for: - explicit `preserve_order` values and legacy payload behavior - default and custom expression adapter factories - projected memory-source sort information, projection, fetch, and display settings ## Are there any user-facing changes? File scans now preserve an explicit `preserve_order` value across protobuf round trips. Default serialization of plans containing custom expression adapter factories now returns an explicit error rather than silently dropping the factory. Projected memory sources now preserve their sort information across round trips. The new protobuf field is additive, and `PhysicalExprAdapterFactory` gains a default method, so existing implementations remain source compatible.
| Commit: | 4cb8c21 | |
|---|---|---|
| Author: | Burak Şen | |
| Committer: | GitHub | |
fix(proto): destructure JSON source and sink serde hooks (#24945) ## Which issue does this PR close? - Closes #24621. ## Rationale for this change Protobuf hooks that access fields individually can silently omit newly added state. For JSON sinks, this caused an explicitly configured compression level to revert to the default after a physical-plan protobuf roundtrip. Exhaustive destructuring makes newly added source, sink, and wire fields compile errors until their serialization behavior is explicitly considered. ## What changes are included in this PR? - Exhaustively destructure `JsonSource` and `JsonSink` in their encoders. - Exhaustively destructure their protobuf nodes in the decoders. - Keep the active `JsonSink` hook exhaustive while centralizing field mapping in its public `TryFrom<&JsonSink>` conversion. - Add `compression_level` to the `JsonWriterOptions` protobuf message and preserve it in both conversion directions. - Strengthen the JSON sink roundtrip test to inspect non-default sink options, file configuration, and sort order directly. The protobuf wire-format change is additive and backward compatible. ## Are these changes tested? Yes. The focused JSON source and sink roundtrip tests pass: ```bash cargo test -p datafusion-proto --test proto_integration roundtrip_json ``` ## Are there any user-facing changes? Physical plans containing JSON sinks now preserve an explicitly configured compression level across protobuf roundtrips. Older payloads without the new field continue to decode with no explicit compression level. Adding `compression_level` to the generated public protobuf `JsonWriterOptions` Rust struct is a source-level API break for callers that construct it with an exhaustive struct literal. Those callers must initialize `compression_level` (use `None` to preserve prior behavior) or use `..Default::default()`. This migration is documented in the DataFusion 56.0.0 upgrade guide.
| Commit: | e1ca94f | |
|---|---|---|
| Author: | Marcelo Tesla | |
| Committer: | GitHub | |
fix: preserve Arrow field metadata on scalar subquery expressions (#24939) ## Which issue does this PR close? - Closes #24933 Related: [apache/sedona-db#1231](https://github.com/apache/sedona-db/issues/1231), [apache/sedona-db#1226](https://github.com/apache/sedona-db/pull/1226). This is the uncorrelated physical-plan counterpart of the earlier outer-reference metadata fix in #17524 / #17422. ## Rationale for this change UDFs that distinguish Arrow extension types from their storage types (for example spatial predicates such as `ST_Intersects`) need `ARROW:extension:name` on every argument. Uncorrelated scalar subqueries kept that metadata in the logical plan, but physical planning built a `ScalarSubqueryExpr` from only the data type and nullability. The synthesized physical field had empty metadata, so queries like `WHERE udf(col, (SELECT geometry FROM t WHERE id = 1))` failed even though the equivalent join form worked. ## What changes are included in this PR? - `ScalarSubqueryExpr` now stores the output `FieldRef` (name `scalar_subquery`, original type/nullability, and metadata) via `new_with_metadata`. The existing `new` constructor is unchanged and still produces a field with empty metadata. - Physical lowering copies metadata from the logical subquery output field while still using `Expr::nullable` so zero-row subqueries remain nullable. - Protobuf encoding adds an additive `metadata` map on `PhysicalScalarSubqueryExprNode` so plan round-trips keep extension metadata. ## What is the testing strategy for this PR? - Unit test `scalar_subquery_preserves_output_field_metadata` in `planner.rs` reproduces the drop during physical lowering. - Unit test `return_field_preserves_extension_metadata` and an updated proto round-trip in `scalar_subquery.rs`. - End-to-end regression `test_extension_metadata_preserve_in_uncorrelated_scalar_subquery` in `user_defined_scalar_functions.rs`, based on the issue reproducer. The existing EXISTS-subquery metadata test still passes. ## Are there any user-facing changes? Additive only: `ScalarSubqueryExpr::new_with_metadata` and an optional protobuf `metadata` map (older payloads decode as empty metadata). Existing `new(data_type, nullable, ...)` keeps working. Queries whose UDFs inspect argument field metadata now see the subquery's original extension metadata in the physical plan. --------- Co-authored-by: Marcelo Tesla <9055877+M-Tesla@users.noreply.github.com>
| Commit: | 5980374 | |
|---|---|---|
| Author: | Jayant Shrivastava | |
| Committer: | GitHub | |
fix: preserve Parquet sort pushdown across proto roundtrips (#24930) ## Which issue does this PR close? - Informs https://github.com/apache/datafusion/issues/24620. ## Rationale for this change `reverse_row_groups` and `sort_order_for_reorder` are not preserved across proto round trips, so those optimizations are lost. ## What changes are included in this PR? In `try_to_proto` and `try_from_proto` in `ParquetSource`, serialize those fields. ## What is the testing strategy for this PR? New unit test `roundtrip_parquet_exec_with_sort_pushdown`
| Commit: | 5e168c9 | |
|---|---|---|
| Author: | Huaijin | |
| Committer: | GitHub | |
fix(proto): preserve AnalyzeExec metric types across serialization (#24669) ## Which issue does this PR close? - Related to #23494. ## Rationale for this change `AnalyzeExec::metric_types` was lost during protobuf round-trips, causing non-default selections such as summary-only metrics to reset to `[Summary, Dev]`. ## What changes are included in this PR? - Serialize and deserialize `AnalyzeExec::metric_types`. - Preserve compatibility with older protobuf messages. - Distinguish an absent field from an explicitly empty metric type list. - Add regression tests for summary, dev, empty, and legacy states. ## Are these changes tested? Yes. - Relevant `datafusion-proto` integration tests - Clippy for the affected crates - Formatting and diff checks ## Are there any user-facing changes? Yes. `EXPLAIN ANALYZE` metric type selections now survive physical-plan protobuf serialization. There are no public Rust API changes.
| Commit: | 408fc6f | |
|---|---|---|
| Author: | Subham Singhal | |
| Committer: | GitHub | |
feat: proto serialization for PiecewiseMergeJoinExec (#24378) ## Which issue does this PR close? Part of https://github.com/apache/datafusion/issues/17427. ## Rationale for this change `PiecewiseMergeJoinExec` has no protobuf representation, so a query the planner turns into a range join cannot be serialized and cannot run on any engine that ships physical plans between processes (Ballista, Flight-based executors): ```text Internal("Unsupported plan and extension codec failed with [This feature is not implemented: PhysicalExtensionCodec is not provided]. Plan: PiecewiseMergeJoinExec { ... }") ``` ## What changes are included in this PR? - New PiecewiseMergeJoinExecNode at oneof PhysicalPlanType tag 40, carrying exactly the six arguments of try_new. Everything else (output schema, sort_options, required input orderings, PlanProperties) is derived inside try_new, so it is not on the wire and the round trip is exact by construction. - ExecutionPlan::try_to_proto + PiecewiseMergeJoinExec::try_from_proto, following the self-serialization pattern (#22419) that SortMergeJoinExec uses. Since the hook is consulted before the central downcast chain, the encode side needs no change in datafusion/proto; only one decode arm is added. - Regenerated pbjson.rs / prost.rs via datafusion/proto-models/regen.sh. ## Are these changes tested? Yes using UT ## Are there any user-facing changes? No breaking changes
| Commit: | 97e1c1c | |
|---|---|---|
| Author: | Burak Şen | |
| Committer: | GitHub | |
fix(proto): preserve CSV/JSON scan options on roundtrip (#24233) ## Which issue does this PR close? - Closes #24180. ## Rationale for this change Physical-plan protobuf serialization does not preserve several CSV and JSON scan options. Custom CSV terminators, JSON newline-delimited mode, and file compression therefore revert to their defaults after a roundtrip, which can cause the decoded plan to read the file incorrectly. ## What changes are included in this PR? - Add `terminator` to `CsvScanExecNode`. - Add `newline_delimited` to `JsonScanExecNode`. - Add `file_compression_type` to the shared `FileScanExecConf`. - Serialize and restore these options in CSV and JSON scans. ## Are these changes tested? Yes. Extended the CSV and JSON physical-plan roundtrip tests to cover custom terminators, non-newline-delimited JSON, compression, and backward-compatible defaults. ## Are there any user-facing changes? CSV and JSON scans now preserve their format and compression options across protobuf roundtrips. The protobuf changes are additive and backward compatible. `JsonSource` also gains an `is_newline_delimited` getter.
| Commit: | 5134a1a | |
|---|---|---|
| Author: | Daniël Heres | |
| Committer: | GitHub | |
feat: add projection support to SortMergeJoinExec (#24517) ## Which issue does this PR close? - Closes https://github.com/apache/datafusion/issues/24518 ## Rationale for this change `SortMergeJoinExec` emits every column of both inputs. We can add projection support to only produce the output columns (saving some compute / copy). ```sql select * from t1 right join t2 on t1.c3 = t2.c3 -- c3 needs a cast ``` ```text ProjectionExec: expr=[c1@0 as c1, ..., c4@8 as c4] SortMergeJoinExec: join_type=Right, on=[(CAST(t1.c3 AS Decimal128(10, 2))@4, c3@2)] ``` The cast column is only there to be joined on, and every operator above the join carries it until the projection removes it. ## What changes are included in this PR? `SortMergeJoinExec` takes an optional projection, set with `with_projection`, the same shape as `HashJoinExec`'s. `try_swapping_with_projection` embeds the projection into the join when it cannot push it into the children, so the query above becomes: ```text SortMergeJoinExec: join_type=Right, on=[(CAST(t1.c3 AS Decimal128(10, 2))@4, c3@2)], projection=[c1@0, c2@1, c3@2, c4@3, c1@5, c2@6, c3@7, c4@8] ``` The change doesn't bring a big speedup but helps aligning with other join types and helping simplify join optimization in other areas (e.g. join enumeration). ## Are these changes tested? Yes: - a projection pushdown test for a projection that interleaves the two sides, which the existing pushdown cannot handle - a serialization round trip - existing sqllogictests, whose plans lose a `ProjectionExec` in four places ## Are there any user-facing changes? Explain / proto changes. --------- Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
| Commit: | 83be77c | |
|---|---|---|
| Author: | 0x70 | |
| Committer: | GitHub | |
fix: translate SQL wildcards in SIMILAR TO patterns (#22263) (#23188) `SIMILAR TO` previously passed the pattern straight to Arrow's regex engine, so SQL wildcards were never translated and matches were unanchored: SELECT 'abc' SIMILAR TO 'a%'; -- returned false SELECT 'x' SIMILAR TO '_'; -- returned false Translate `%` to `(?s:.*)` and `_` to `(?s:.)` (dot-all so they match newlines), then wrap the pattern in `^(?:...)$` so the regex matches the entire string. `.`, `^`, `$`, and `\` are escaped as SQL literals. The POSIX metacharacters that `SIMILAR TO` defines (`| * + ? ( ) { } [ ]`) pass through to the regex unchanged. ## Which issue does this PR close? - Closes #22263. ## Rationale for this change `SIMILAR TO` is a SQL standard operator with well-defined wildcard semantics (`%` = any sequence, `_` = single character, full-string match). DataFusion's previous behavior silently produced wrong results for the most basic patterns, which is a correctness bug for anyone porting queries from Postgres or other SQL engines. Supporting non-literal patterns requires a new physical expression, so that expression also needs protobuf support — otherwise serialized physical plans containing a dynamic `SIMILAR TO` would regress for distributed engines built on `datafusion-proto`. Accepting `LargeUtf8` and `Utf8View` patterns in turn requires type coercion. `Expr::SimilarTo` was previously in the analyzer's no-op list, so nothing reconciled the value type with the pattern type. Since the regex kernel dispatches on the left-hand type and then downcasts the pattern to that same array type, any mismatch aborted the query with a `failed to downcast array` panic. ## What changes are included in this PR? **Pattern translation** - Added a `sql_similar_to_regex` helper that translates `%`/`_` and anchors the pattern with `^(?:...)$`. It tracks bracket state so `^` inside `[...]` is treated as bracket negation, not as a literal, and so SQL wildcards lose their special meaning inside a bracket expression. - Added a new `SqlSimilarToPattern` physical expression in `datafusion/physical-expr/src/expressions/similar_to_pattern.rs` that applies the translation at runtime for non-literal patterns. - `similar_to()` now translates literal `Utf8` / `LargeUtf8` / `Utf8View` patterns at planning time and wraps non-literal patterns in `SqlSimilarToPattern` for runtime translation. - NULL patterns pass through and return `NULL` instead of crashing with an internal error. The translation preserves the pattern's string variant, including for NULL, so the pattern type always matches the value type. - Restored a plan-time type check in `datafusion/sql/src/expr/mod.rs` that rejects non-string patterns with a clean `plan_err!`. - Runtime type errors in `SqlSimilarToPattern` are reported as `exec_err!` rather than `internal_err!`. **Type coercion** - Added an `Expr::SimilarTo` arm to `datafusion/optimizer/src/analyzer/type_coercion.rs` and removed `Expr::SimilarTo` from the no-op list. It uses `regex_coercion`, the same coercion `Operator::RegexMatch` already uses, since that is what `SIMILAR TO` lowers to. Mismatched string types are now cast to a common type instead of reaching the kernel and panicking. **Protobuf support** - Added a `PhysicalSqlSimilarToPatternNode` message and wired it into `PhysicalExprNode.expr_type` as field 24. - Implemented `try_to_proto` / `try_from_proto` for `SqlSimilarToPattern` and added the decode arm in `datafusion/proto/src/physical_plan/from_proto.rs`. - Regenerated `prost.rs` / `pbjson.rs` via `datafusion/proto-models/regen.sh`. ## Are these changes tested? Yes. Pattern semantics (`datafusion/physical-expr/src/expressions/binary.rs`): - `test_similar_to_sql_literal_metachars` confirms that `.`, `^`, `$`, and `\` are treated as SQL literals, not as regex operators. - `test_similar_to_posix_metachars` confirms that `|, *, +, ?, (, ), {, }, [, ], [^...]`, and `[a-z]` behave as `SIMILAR TO` metacharacters. - `test_similar_to_wildcards_match_newlines` confirms that `%` and `_` match newlines. - `test_similar_to` covers basic `%`/`_` semantics, full-string anchoring, and case sensitivity. - `test_similar_to_dynamic_pattern` covers column-based patterns. - `test_similar_to_null_pattern` and `test_similar_to_non_literal_pattern_errors` cover the NULL pattern and non-string literal pattern paths. Translation and coercion: - `test_translate_scalar` / `test_translate_array` cover the translation itself, including that each string variant (and its NULL) is preserved. - `similar_to_for_type_coercion` in `type_coercion.rs` covers matching types, mismatched string types, a NULL pattern cast to the value's type, and the no-common-type error. Protobuf (`similar_to_pattern.rs` `proto_tests`, mirroring the existing `LikeExpr` tests): - Encoding, decoding, rejection of a non-`SqlSimilarToPattern` node, rejection of a missing `expr` field, and error propagation on both the encode and decode paths. - `roundtrip_sql_similar_to_pattern` in `datafusion/proto/tests/cases/roundtrip_physical_plan.rs` covers the full plan roundtrip. End-to-end (`datafusion/sqllogictest/test_files/strings.slt`): - The pre-existing `SIMILAR TO 'p[12].*'` / `NOT SIMILAR TO 'p[12].*'` cases (which only worked because of the bug) were replaced with equivalent cases using standard wildcard syntax. - New cases cover literal metachars, POSIX metachars, newline matching, dynamic column patterns, dynamic function patterns, `SELECT 'a' SIMILAR TO NULL`, rejection of `SELECT 'a' SIMILAR TO 1`, mixed `Utf8` / `LargeUtf8` / `Utf8View` operands, and dictionary-encoded values with both literal and dynamic patterns. ## Are there any user-facing changes? Yes: - `SIMILAR TO` now produces correct results for queries that were previously returning wrong answers. - Queries that relied on the buggy behavior (e.g., treating `.`, `^`, or `$` as regex metacharacters) now follow standard SQL `SIMILAR TO` semantics. - Queries using valid `SIMILAR TO` POSIX metacharacters (`| * + ? ( ) { } [ ]`) now work as expected. - `%` and `_` wildcards now work. - Non-literal patterns (e.g., column patterns) now work instead of returning a `not_impl_err!`, and physical plans containing them can be serialized with `datafusion-proto`. - `LargeUtf8` and `Utf8View` patterns are accepted, and mixing string types between the value and the pattern is coerced rather than failing. - Dictionary-encoded values are coerced to their value type and work with both literal and dynamic patterns. - `SELECT ... SIMILAR TO NULL` now returns `NULL` instead of crashing. - `SELECT ... SIMILAR TO <non-string>` now returns a clean plan error instead of an internal error. `\` is still treated as a literal backslash rather than as the SQL default escape character, and an explicit `ESCAPE` clause is still rejected. Escape support is left as a follow-up.
| Commit: | f956fd1 | |
|---|---|---|
| Author: | Burak Şen | |
| Committer: | GitHub | |
fix(proto): preserve Arrow IPC stream format (#24224) ## Which issue does this PR close? - Closes #24196. ## Rationale for this change `ArrowSource` distinguishes between Arrow IPC file and stream formats, but this format was not serialized. As a result, a stream scan round-tripped through protobuf as a file scan, selecting the wrong opener and potentially allowing unsupported range-based repartitioning. ## What changes are included in this PR? - Add an `ArrowIpcFormat` discriminator to `ArrowScanExecNode`. - Serialize and restore the `ArrowSource` IPC format. - Decode payloads without the new field as file format, preserving the previous behavior. - Regenerate the prost and pbjson models. - Add regression coverage for file, stream, and older payloads without the format field. ## Are these changes tested? Yes: - `cargo test -p datafusion-proto --test proto_integration roundtrip_arrow` - `cargo test -p datafusion-proto --test proto_integration arrow_scan_without_format_field_decodes_as_file_format` - `cargo fmt --all` - Targeted all-feature clippy for `datafusion-datasource-arrow` and `datafusion-proto` ## Are there any user-facing changes? Arrow IPC stream scans now preserve their format across protobuf round trips. The binary protobuf change is additive. Older payloads continue to decode as file format. Stream preservation requires both producer and consumer to include this change because older versions do not carry or read the discriminator. The generated Rust `ArrowScanExecNode` struct gains a `format` field.
| Commit: | 618aaff | |
|---|---|---|
| Author: | Peter Lee | |
| Committer: | GitHub | |
Enable dynamic filters for range-partitioned joins (#23854) ## Which issue does this PR close? - Closes https://github.com/apache/datafusion/issues/23376. ## Rationale for this change Partitioned hash joins build one dynamic filter per build partition. Existing routing uses `hash(key) % N`, which cannot reproduce a Range partitioning layout. Compatible Range co-partitioned joins instead need to route probe rows using their existing ordering and split points. ## What changes are included in this PR? - Enable dynamic filter pushdown for hash joins with compatible Range-partitioned inputs. - Build a searched `CASE` expression that routes probe rows to the corresponding partition filter using the Range ordering and split points. - Move the TopK lexicographic filter builder into the shared ordering module for reuse. ## Are these changes tested? unit test. ## Are there any user-facing changes? Yes. This PR adds `RangeExpr` to the physical-expression protobuf model, which adds the public `ExprType::RangeExpr` enum variant. Downstream Rust consumers that exhaustively match `ExprType` must handle the new variant. It also enables dynamic-filter pushdown for compatible Range-partitioned joins.
| Commit: | a251b94 | |
|---|---|---|
| Author: | Burak Şen | |
| Committer: | GitHub | |
fix(proto): preserve AggregateExec schema and reversed state (#24207) ## Which issue does this PR close? - Closes #24202. ## Rationale for this change Physical-plan deserialization rebuilt `AggregateExec` output schemas from decoded aggregate expression names. `OptimizeAggregateOrder` can reverse aggregate expressions while preserving the original schema, causing protobuf round trips to change output field names and lose the expression's reversed state. ## What changes are included in this PR? - Serialize and restore the preserved `AggregateExec` output schema. - Serialize and restore `AggregateFunctionExpr::is_reversed`. - Fall back to schema reconstruction for older payloads without an output schema. - Regenerate the prost and pbjson models. - Add a byte-level regression test covering optimizer reversal and backward compatibility. ## Are these changes tested? Yes: - `cargo test -p datafusion-proto --test proto_integration roundtrip_aggregate_preserves_optimizer_schema_and_reversed_state` - `cargo fmt --all` - `cargo clippy --all-targets --all-features -- -D warnings` - `./dev/rust_lint.sh` ## Are there any user-facing changes? Physical-plan protobuf round trips now preserve aggregate output field names and reversed state. The protobuf wire format remains backward compatible; generated Rust model structs gain new fields.
| Commit: | f4c8ba1 | |
|---|---|---|
| Author: | Burak Şen | |
| Committer: | GitHub | |
fix(proto): serialize Global/LocalLimitExec required_ordering (#24183) ## Which issue does this PR close? - Closes #24173 ## Rationale for this change `GlobalLimitExec`/`LocalLimitExec` `required_ordering` is not serialized, so a roundtrip loses the only record that a pushed-down `LIMIT` is order-sensitive. Re-running `LimitPushdown` on the decoded plan then sets `preserve_order: false` on the scan, which is free to read files out of order — `ORDER BY ... LIMIT` can return the wrong rows. ## What changes are included in this PR? - Add `required_ordering` to `GlobalLimitExecNode` (field 4) and `LocalLimitExecNode` (field 3); empty list means `None` - New `optional_ordering_try_to_proto`/`optional_ordering_try_from_proto` helpers in `physical-expr-common`, also reused by `SymmetricHashJoinExec` serde which hand-rolled the same pattern - Preserve `required_ordering` in both execs' `with_new_children` ## Are these changes tested? Yes new tests that would fail on main if not fixed ## Are there any user-facing changes? No
| Commit: | fc846dd | |
|---|---|---|
| Author: | Adrian Garcia Badaracco | |
| Committer: | GitHub | |
fix(proto): preserve HashJoinExec fetch across serialization (#24165) ## Which issue does this PR close? <!-- No dedicated issue; found while auditing the plans covered by the `try_to_proto`/`try_from_proto` migration EPIC. --> - Related to #23494 (found while auditing that EPIC's plans for unserialized fields). This is a bug fix, not part of the migration checklist. ## Rationale for this change `HashJoinExec.fetch` was silently dropped by protobuf serialization. `protobuf::HashJoinExecNode` had no `fetch` field, so `HashJoinExec`'s `try_to_proto` never wrote it and `try_from_proto` never restored it: a plan with `fetch = Some(n)` round-tripped to `fetch = None`. This is user-visible. The `limit_pushdown` physical optimizer rule pushes a limit into the join via `ExecutionPlan::with_fetch`, then marks the global state satisfied and drops the enclosing `GlobalLimitExec`. So after a proto round-trip the plan carried no limit at all, and a distributed executor (Ballista/Comet-style, anything that ships physical plans over the wire) returned more rows than the query asked for. ## What changes are included in this PR? - `datafusion.proto`: add `optional uint64 fetch = 12` to `HashJoinExecNode`. The field is **presence-tracked on purpose**, and this is the load-bearing detail for wire compatibility. Messages written by versions predating this field carry no `fetch` at all, and a plain proto3 scalar decodes that absence as `0`. With the negative-sentinel convention used by `SortExecNode`'s `int64 fetch`, `0` would mean "fetch 0 rows" and would silently turn every older plan into an empty result. `optional` gives prost an `Option<u64>` where absent decodes to `None`, which is the correct reading of an older message. A comment in the `.proto` records this. - Regenerated `prost.rs` / `pbjson.rs` via `datafusion/proto-models/regen.sh` (no hand edits). - `hash_join/exec.rs`: write `self.fetch` in the `try_to_proto` hook and restore it in `try_from_proto` via the builder's `with_fetch`, matching how the plan is normally constructed. - New regression test `roundtrip_hash_join_fetch`. The deprecated `PhysicalPlanNodeExt` shims (`try_from_hash_join_exec` / `try_into_hash_join_physical_plan`) delegate straight to these two hooks, so they pick the fix up with no separate change. Verified by reading them rather than assumed. ## Are these changes tested? Yes. `roundtrip_hash_join_fetch` in `datafusion/proto/tests/cases/roundtrip_physical_plan.rs` builds a `HashJoinExec`, applies `with_fetch(Some(7))` the way `limit_pushdown` does, round-trips it through `physical_plan_to_bytes_with_proto_converter` / `physical_plan_from_bytes_with_proto_converter`, and asserts `fetch()` is still `Some(7)`. It also covers `fetch = None`. The assertion deliberately inspects `fetch()` rather than the plan's string form. The existing `roundtrip_test` helper compares `format!("{plan:?}")`, and `HashJoinExec`'s `Debug` output does not include `fetch` — which is exactly why this went unnoticed. I confirmed this empirically: with the encode side reverted, the Debug comparison inside the helper still passes and only the `fetch()` assertion fails (`left: None, right: Some(7)`). Ran locally: - `cargo fmt --all` - `cargo test -p datafusion-proto --test proto_integration` — 215 passed, 0 failed - `cargo test -p datafusion-physical-plan` — 1640 + 9 passed, 0 failed - `cargo clippy --all-targets --all-features` on the touched packages. The changed code is clean; the only two errors reported are pre-existing on an unmodified `main` with my newer local clippy (`uninlined_format_args` in `datafusion/proto-common/src/generated/pbjson.rs` and `needless_pass_by_value` in `datafusion/proto/src/bytes/mod.rs`), in files this PR does not touch. ## Are there any user-facing changes? Yes, a bug fix: a limit pushed into a hash join now survives physical-plan serialization, so distributed executors no longer over-return rows. No API changes. The new proto field is backward and forward compatible in both directions — old readers ignore tag 12, and new readers treat its absence as "no limit". --------- Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
| Commit: | aa38d3c | |
|---|---|---|
| Author: | Qi Zhu | |
| Committer: | GitHub | |
feat(pruning): expose pruning predicate IN-list rewrite size cap as a config option (#24074) ## Which issue does this PR close? - Closes #24059. ## Rationale `PruningPredicate` rewrites `col IN (v1..vn)` into a chain of per-value min/max checks (via `build_predicate_expression`), but only when `n <= MAX_LIST_VALUE_SIZE_REWRITE` — currently a hardcoded `20`. Beyond that, the IN branch falls through to `unhandled_hook`, which by default returns `TRUE`, so row-group and file-range statistics pruning does not fire at all for IN lists longer than 20. This is problematic for query patterns that pass a batch of identifiers as `col IN (...)` — REST endpoints filtering by a page of ~25-100 values, ORM-generated `WHERE id IN (25 items)` queries, batched crawlers. On a table sorted by `col`, the reader is forced to materialize the filter column across every row group instead of skipping row groups whose stats disagree with the IN set. Full context in #24059. ## What changes are included in this PR? - New config option `datafusion.execution.parquet.pruning_max_in_list_size: usize` (default `20`, preserving existing behaviour), placed next to `max_predicate_cache_size` on `TableParquetOptions.global`. - `MAX_LIST_VALUE_SIZE_REWRITE` promoted to `pub const` so callers can reference the historical default explicitly. - `PredicateRewriter::with_max_in_list_size(usize) -> Self` builder, mirroring the existing `with_unhandled_hook`. - `PruningPredicate::try_new_with_max_in_list_size` variant. - `build_pruning_predicate_with_max_in_list_size` variant of the public helper. - Value threaded through `datasource-parquet`: `ParquetSource::pruning_max_in_list_size()` reads from `TableParquetOptions.global`, propagates through `ParquetMorselizer` → `PreparedParquetOpen` → `RowGroupPruner`, then flows into `build_pruning_predicates` at the opener and `build_pruning_predicate_with_max_in_list_size` inside the dynamic row-group pruner. Internal `build_predicate_expression` gains a new `usize` parameter (crate-private). ## Backward compatibility - `PruningPredicate::try_new` and `build_pruning_predicate` are preserved as thin wrappers that pass the historical `MAX_LIST_VALUE_SIZE_REWRITE` default. All existing callers continue to work with unchanged behaviour. - The config option default is `20`, so behaviour is unchanged unless the option is set explicitly. ## Are these changes tested? Two new unit tests in `datafusion-pruning`: - `row_group_predicate_in_list_rewritten_at_raised_cap`: `PredicateRewriter::with_max_in_list_size(32)` rewrites a 25-item IN into per-value min/max checks OR'd together, instead of falling through to `true`. - `row_group_predicate_in_list_disabled_at_zero_cap`: `cap = 0` skips the IN rewrite even for small lists (opt-out path). The existing `row_group_predicate_in_list_to_many_values` continues to pass, guarding the default-20 behaviour. ## Are there any user-facing changes? Yes — one new config option (`datafusion.execution.parquet.pruning_max_in_list_size`, default `20`). Users who want row-group / file-range pruning for IN lists longer than 20 items can raise it (e.g., `SET datafusion.execution.parquet.pruning_max_in_list_size = 128`). New public API on `datafusion-pruning`: - `MAX_LIST_VALUE_SIZE_REWRITE: usize` (re-exported) - `PredicateRewriter::with_max_in_list_size(usize) -> Self` - `PruningPredicate::try_new_with_max_in_list_size(expr, schema, size)` - `build_pruning_predicate_with_max_in_list_size(predicate, schema, errors, size)` --------- Co-authored-by: Andrew Lamb <andrew@nerdnetworks.org>
| Commit: | f2b4835 | |
|---|---|---|
| Author: | Krishna Sudarshan J | |
| Committer: | GitHub | |
feat: Add support for `explode_outer` function for arrays (#22100) ## Which issue does this PR close? <!-- We generally require a GitHub issue to be filed for all bug fixes and enhancements and this helps us generate change logs for our releases. You can link an issue to this PR using the GitHub syntax. For example `Closes #123` indicates that this PR will close issue #123. --> - Closes #19053. ## Rationale for this change <!-- Why are you proposing this change? If this is already explained clearly in the issue then this section is not needed. Explaining clearly why changes are proposed helps reviewers understand your changes and offer better suggestions for fixes. --> DataFusion's 'unnest' had no way to express Spark 'explode_outer' semantics, empty input lists were silently dropped, even with 'preserve_nulls' = true. ## What changes are included in this PR? <!-- There is no need to duplicate the description in the issue here but it is sometimes worth providing a summary of the individual changes in this PR. --> Adds a third unnest behavior that produces a `NULL` row for empty input lists, by replacing `UnnestOptions.preserve_nulls: bool` with a `NullHandling { Drop, Preserve, PreserveAndExpandEmpty }` enum. The `with_preserve_nulls(bool)` builder is kept as a backward-compat shim. ## Are these changes tested? <!-- We typically require tests for all PRs in order to: 1. Prevent the code from being accidentally broken by subsequent changes 2. Serve as another way to document the expected behavior of the code If tests are not included in your PR, please explain why (for example, are they covered by existing tests)? --> Yes, new unit test for the empty-list case, extended longest-length and DataFrame `unnest_column_nulls` tests, and existing proto round-trip coverage. ## Are there any user-facing changes? <!-- If there are user-facing changes then we may require documentation to be updated before approving the PR. --> <!-- If there are any breaking changes to public APIs, please add the `api change` label. --> The `preserve_nulls` field on `UnnestOptions` is renamed to `null_handling`. The `with_preserve_nulls(bool)` builder is preserved, so most callers are unaffected. Add the `api change` label for the field rename.
| Commit: | 5283bc6 | |
|---|---|---|
| Author: | Gabriel | |
| Committer: | Gabriel | |
Decide at planning time whether we should reuse hashes or not
| Commit: | bb670fb | |
|---|---|---|
| Author: | Phoenix | |
| Committer: | GitHub | |
refactor(proto): remove legacy scan field (#23445) ## Which issue does this PR close? <!-- We generally require a GitHub issue to be filed for all bug fixes and enhancements and this helps us generate change logs for our releases. You can link an issue to this PR using the GitHub syntax. For example `Closes #123` indicates that this PR will close issue #123. --> - Closes #N/A. ## Rationale for this change <!-- Why are you proposing this change? If this is already explained clearly in the issue then this section is not needed. Explaining clearly why changes are proposed helps reviewers understand your changes and offer better suggestions for fixes. --> Since DataFusion doesn't typically guartantee wire format compatibility, I remove the backward compatiblity shim that introduced in PR #23189 FYI: https://github.com/apache/datafusion/pull/23189#discussion_r3554855165 ## What changes are included in this PR? <!-- There is no need to duplicate the description in the issue here but it is sometimes worth providing a summary of the individual changes in this PR. --> DataFusion does not guarantee serialized plans across versions. Keeping `partitioned_by_file_group` therefore leaves a dead schema field and decoder path after `output_partitioning` became the source of truth. Reserve the old field number and name to prevent future reuse. ## Are these changes tested? <!-- We typically require tests for all PRs in order to: 1. Prevent the code from being accidentally broken by subsequent changes 2. Serve as another way to document the expected behavior of the code If tests are not included in your PR, please explain why (for example, are they covered by existing tests)? --> Yes ## Are there any user-facing changes? <!-- If there are user-facing changes then we may require documentation to be updated before approving the PR. --> <!-- If there are any breaking changes to public APIs, please add the `api change` label. --> Yes Signed-off-by: Jiawei Zhao <Phoenix500526@163.com> Co-authored-by: Andrew Lamb <andrew@nerdnetworks.org>
| Commit: | 1dc031d | |
|---|---|---|
| Author: | Kumar Ujjawal | |
| Committer: | GitHub | |
feat: Support multiple external table locations (#22695) ## Which issue does this PR close? <!-- We generally require a GitHub issue to be filed for all bug fixes and enhancements and this helps us generate change logs for our releases. You can link an issue to this PR using the GitHub syntax. For example `Closes #123` indicates that this PR will close issue #123. --> Part of #16303. ## Rationale for this change `CREATE EXTERNAL TABLE` can reference only one location today. This adds support for listing multiple explicit locations and reading them as one table. <!-- Why are you proposing this change? If this is already explained clearly in the issue then this section is not needed. Explaining clearly why changes are proposed helps reviewers understand your changes and offer better suggestions for fixes. --> ## What changes are included in this PR? - Adds `LOCATION ('a.parquet', 'b.parquet')` syntax. - Keeps `LOCATION 'a,b.parquet'` as a single path, so literal commas still work. - Carries multiple locations through the logical plan and proto. - Updates listing table creation to scan all listed locations. - Requires all listed locations to use the same object store and matching fields. - Keeps stream tables limited to exactly one location. - Updates docs and upgrade notes. <!-- There is no need to duplicate the description in the issue here but it is sometimes worth providing a summary of the individual changes in this PR. --> ## Are these changes tested? Yes <!-- We typically require tests for all PRs in order to: 1. Prevent the code from being accidentally broken by subsequent changes 2. Serve as another way to document the expected behavior of the code If tests are not included in your PR, please explain why (for example, are they covered by existing tests)? --> ## Are there any user-facing changes? Yes. `CREATE EXTERNAL TABLE` now accepts a parenthesized list of locations. There is also a public API change: `CreateExternalTable.location` is replaced by `CreateExternalTable.locations`. <!-- If there are user-facing changes then we may require documentation to be updated before approving the PR. --> <!-- If there are any breaking changes to public APIs, please add the `api change` label. -->
| Commit: | 0840e5c | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
fix: preserve EmptyExec and PlaceholderRowExec partition count across proto round-trip (#23643) ## Which issue does this PR close? - Closes #23642. ## Rationale for this change `EmptyExecNode` and `PlaceholderRowExecNode` encoded only a schema, so the partition count set by `EmptyExec::with_partitions(n)` was silently lost across a physical-plan round-trip: a plan reporting `n` partitions before serialization reported `1` after. This is silent data loss rather than an error — the decoded plan is well-formed but describes a different plan than the one encoded. It bites any setup that plans in one process and executes in another. In Ballista, an optimizer rule collapses provably-empty sub-plans into an `EmptyExec` inheriting the replaced node's partition count; the scheduler sizes the stage's task count from that count, ships the plan to an executor, and every task above partition 0 fails with: ``` Internal("Assertion failed: partition < self.partitions: EmptyExec invalid partition 1 (expected less than 1)") ``` ## What changes are included in this PR? - Add a `uint32 partitions` field to `EmptyExecNode` and `PlaceholderRowExecNode` in `datafusion.proto`, and regenerate the prost/pbjson code via `regen.sh`. - Write `partitions` on encode and apply it via `with_partitions(...)` on decode. Both `EmptyExec` and `PlaceholderRowExec` mirror their private `partitions` field into `PlanProperties` as `UnknownPartitioning(n)`, so the encoder reads the count via `properties().output_partitioning().partition_count()` on the existing `ExecutionPlan` trait. No new public API is added to `datafusion-physical-plan`. The wire format stays compatible in both directions. Plans encoded before this field existed carry no value for it, which prost surfaces as `0`; decode maps that to the previous default of `1`. Plans encoded after this change add a field that older readers skip. Note for #23501, which migrates these nodes to the `try_to_proto` / `try_from_proto` pattern: that issue specifies "schema only" and a byte-for-byte identical wire format, which would reintroduce this bug. The `partitions` field should be folded into that rewrite. ## Are these changes tested? Yes, three new tests in `roundtrip_physical_plan.rs`: partition-count round-trips for `EmptyExec` and `PlaceholderRowExec`, plus one decoding `partitions: 0` nodes directly to pin the backward-compatibility mapping. The round-trip tests were confirmed to fail against the unfixed decoder with exactly the reported symptom (4 partitions decoding to 1). ## Are there any user-facing changes? No API changes to `datafusion-physical-plan`. Encoded plans gain a new protobuf field, which is backward and forward compatible as described above. The generated `EmptyExecNode` and `PlaceholderRowExecNode` structs gain a `partitions` field, which breaks downstream code constructing them with an exhaustive struct literal; this is documented in the 55.0.0 upgrade guide.
| Commit: | 7ac784f | |
|---|---|---|
| Author: | gstvg | |
| Committer: | GitHub | |
Add protobuf support for lambdas (#22362) ## Which issue does this PR close? Part of #21172 ## Rationale for this change Protobuf support wasn't implemented in main lambda PR to not make it even bigger ## What changes are included in this PR? Protobuf encoding and decoding (~1000 LOC in generated files, ~210 impl, ~400 tests) ## Are these changes tested? Unit tests, similar to the existing ones for scalar functions ## Are there any user-facing changes? Proto `ExprType` has new variants
| Commit: | 63ad991 | |
|---|---|---|
| Author: | Gene Bordegaray | |
| Committer: | GitHub | |
Add `ListingOptions::output_partitioning` and `FileScanConfig::output_partitioning` for pre-defined file partitioning (#22657) ## Which issue does this PR close? <!-- We generally require a GitHub issue to be filed for all bug fixes and enhancements and this helps us generate change logs for our releases. You can link an issue to this PR using the GitHub syntax. For example `Closes #123` indicates that this PR will close issue #123. --> - Closes #22645. ## Rationale for this change <!-- Why are you proposing this change? If this is already explained clearly in the issue then this section is not needed. Explaining clearly why changes are proposed helps reviewers understand your changes and offer better suggestions for fixes. --> This follows up on #22607 by replacing range-partitioning sqllogictest boilerplate with a general file/listing scan API for declared output partitioning. Related: #21992, #22607, https://github.com/apache/datafusion/pull/22607#discussion_r3323904683 ## What changes are included in this PR? <!-- There is no need to duplicate the description in the issue here but it is sometimes worth providing a summary of the individual changes in this PR. --> - Add declared `output_partitioning` to file scan and listing table configuration. - Preserve declared partition counts during listing-table file grouping. - Serialize scan `output_partitioning` through physical plan proto. - Refactor `range_partitioning.slt` to use a CSV `ListingTable` instead of a custom test-only `TableProvider` / `DataSource`. Contract: - Declared partitioning expressions are written against the full table schema before scan projection. For example, `Range([range_key@0], [(10), (20)], 3)` remains valid if the scan projects `range_key` and falls back to `UnknownPartitioning(3)` if `range_key` is not projected. - Listing tables create one file group per declared output partition (which can exceed `target_partitions`). It is up to the user to plan their partitioning. For example, a 4-partition range declaration creates four scan file groups, adding empty trailing groups when fewer files are present. - File group index is part of the contract: file group `i` must contain rows for declared output partition `i`. DataFusion does not validate row placement, matching other user-declared properties such as sortedness. ## Are these changes tested? <!-- We typically require tests for all PRs in order to: 1. Prevent the code from being accidentally broken by subsequent changes 2. Serve as another way to document the expected behavior of the code If tests are not included in your PR, please explain why (for example, are they covered by existing tests)? --> Yes. ## Are there any user-facing changes? <!-- If there are user-facing changes then we may require documentation to be updated before approving the PR. --> Yes. This adds public API for declaring file/listing scan output partitioning. No breaking API changes. <!-- If there are any breaking changes to public APIs, please add the `api change` label. --> --------- Co-authored-by: Gabriel <45515538+gabotechs@users.noreply.github.com>
| Commit: | a1f56b7 | |
|---|---|---|
| Author: | Saad Tajwar | |
| Committer: | GitHub | |
feat: logical plan protobuf representation for range repartitioning (#23030) ## Which issue does this PR close? <!-- We generally require a GitHub issue to be filed for all bug fixes and enhancements and this helps us generate change logs for our releases. You can link an issue to this PR using the GitHub syntax. For example `Closes #123` indicates that this PR will close issue #123. --> - Closes #22787 ## Rationale for this change <!-- Why are you proposing this change? If this is already explained clearly in the issue then this section is not needed. Explaining clearly why changes are proposed helps reviewers understand your changes and offer better suggestions for fixes. --> The range repartitioning scheme for logical plans does not currently have a protobuf representation. ## What changes are included in this PR? <!-- There is no need to duplicate the description in the issue here but it is sometimes worth providing a summary of the individual changes in this PR. --> A protobuf representation of the `RangeRepartition` struct was added to `datafusion.proto`, and the codegened Rust types were created. Added logic for serializing and deserializing to and from the protobuf representation, and a roundtrip test as well! ## Are these changes tested? <!-- We typically require tests for all PRs in order to: 1. Prevent the code from being accidentally broken by subsequent changes 2. Serve as another way to document the expected behavior of the code If tests are not included in your PR, please explain why (for example, are they covered by existing tests)? --> Yes! Added a test in `roundtrip_logical_plan` ## Are there any user-facing changes? No, adding internal protobuf serialization support for an existing logical plan variant <!-- If there are user-facing changes then we may require documentation to be updated before approving the PR. --> <!-- If there are any breaking changes to public APIs, please add the `api change` label. -->
| Commit: | 7bb6e15 | |
|---|---|---|
| Author: | Gabriel | |
| Committer: | GitHub | |
Remove redundant `collect_stat` and `target_partitions` on `ListingOptions` (#22969) ## Which issue does this PR close? - Closes #. ## Rationale for this change Something that was spotted during the review of: - https://github.com/apache/datafusion/pull/22657 `ListingOptions::target_partitions` and `ListingOptions::collect_stat` duplicate `SessionConfig`'s `execution.target_partitions` and `execution.collect_statistics`. After some investigation, I think they only live on `ListingOptions` for historical reasons: when the struct was added (#1010 5 years ago), `TableProvider::scan` had no access to the session, so the values had to be copied onto the table at build time. Once #2660 passed `SessionState` into `scan`, the fields became redundant (and had already drifted — `scan` read them from the session config while `list_files_for_scan` read the stale copy). This PR makes `SessionConfig` the single source of truth. ## What changes are included in this PR? - Remove `target_partitions`/`collect_stat` fields, their builders, and `with_session_config_options` from `ListingOptions`. - `ListingTable` now reads both values from the session config at scan time. - Reserve proto tags 8/9 in `ListingTableScanNode` and drop the related (de)serialization. - Update benchmarks, factory, and test call sites. ## Are these changes tested? Yes, by existing tests ## Are there any user-facing changes? Yes, breaking: the removed fields/builders require configuring `SessionConfig` instead, and the two proto fields no longer round-trip. --------- Co-authored-by: Andrew Lamb <andrew@nerdnetworks.org>
| Commit: | 408dad3 | |
|---|---|---|
| Author: | Xuanyi Li | |
| Committer: | GitHub | |
Add MERGE INTO types to datafusion-expr (#20763) ## Which issue does this PR close? - part of #20746 [EPIC] Complete DML Support (MERGE, INSERT OVERWRITE, TRUNCATE) #19617 As well as task 1 of #20746 ## Rationale for this change Lay the foundation for MERGE INTO support in DataFusion by adding the logical plan types and their proto serialization. Keeping types separate from execution lets reviewers reason about the data model independently of the planner and physical dispatch. ## What changes are included in this PR? **`datafusion/expr` — new types in `dml.rs`** - `MergeIntoOp` — carries the `ON` join condition and ordered list of `WHEN` clauses - `MergeIntoClause` — a single `WHEN` clause: kind + optional predicate + action - `MergeIntoClauseKind` — `Matched` / `NotMatched` / `NotMatchedByTarget` / `NotMatchedBySource`; includes `is_not_matched_by_target()` and `canonical()` helpers because `NotMatched` and `NotMatchedByTarget` are semantically identical and must be treated identically downstream - `MergeIntoAction` — `Update(Vec<(col, expr)>)` / `Insert { columns, values }` / `Delete` - `WriteOp::MergeInto(MergeIntoOp)` variant added to the existing `WriteOp` enum; `WriteOp` is now `#[non_exhaustive]` so future variant additions are not a SemVer break **`datafusion/proto-models` — proto schema** - Extended `DmlNode` with a `MERGE_INTO` type tag and a boxed `MergeIntoOpNode` payload field - Added `MergeIntoOpNode`, `MergeIntoClauseNode`, `MergeIntoActionNode` messages **`datafusion/proto` — serialization** - `from_proto`: `parse_write_op(&DmlNode, ...)` reads the payload when the type tag is `MergeInto`; defensive helpers `parse_merge_into_op/clause/action` with explicit errors for missing fields - `to_proto`: `serialize_merge_into_op/clause/action` helpers; encode path uses an explicit `match` over all `WriteOp` variants producing `(dml_type, merge_into)` pair — no silent payload loss - Cross-crate conversions use `FromProto` (the crate-local trait) rather than `From` to satisfy the Rust orphan rule after the upstream `datafusion-proto-models` refactor **Proto codegen** — after editing `.proto` files, regenerate with: ```bash PROTOC=/tmp/protoc cargo run --manifest-path datafusion/proto-models/gen/Cargo.toml ``` (Install `protoc` from https://github.com/protocolbuffers/protobuf/releases if not present; set `PROTOC` to its path.) ## Are these changes tested? - `datafusion-expr` unit tests: `WriteOp::MergeInto` display, `is_not_matched_by_target`, `canonical` - `datafusion-proto` round-trip test: exercises all four `MergeIntoClauseKind` variants and all three `MergeIntoAction` variants through encode → decode - `datafusion-proto` error-path tests: missing `merge_into` payload, missing `on` expression, unknown clause kind tag, missing clause action, missing action oneof ## Are there any user-facing changes? `WriteOp` gains a `MergeInto` variant and is now `#[non_exhaustive]`. Existing downstream `match` arms need a wildcard arm added (this is intentional and expected for a new DML operation). ## Follow-up A stacking PR that adds the SQL planner, physical planner dispatch, and `TableProvider::merge_into` hook is available at https://github.com/wirybeaver/datafusion/pull/2. If reviewers prefer to review both together in one pass, I'm happy to include that work here instead. Co-authored-by: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
| Commit: | 08da279 | |
|---|---|---|
| Author: | Mithun Chicklore Yogendra | |
| Committer: | GitHub | |
[branch-54] fix: preserve null_aware on logical JoinNode proto round-trip (backport #22104) (#22785) ## Which issue does this PR close? - Backport of #22104 to `branch-54` (for 54.1.0, tracked in #22547). This PR: - Backports #22104 to the `branch-54` line so the `null_aware` proto round-trip fix ships in 54.1.0, as requested in https://github.com/apache/datafusion/issues/22065#issuecomment-4634038807 Clean cherry-pick; `datafusion-proto` builds and both round-trip regression tests pass on `branch-54`.
| Commit: | 84bc876 | |
|---|---|---|
| Author: | Daipayan Mukherjee | |
| Committer: | GitHub | |
feat: add max_row_group_bytes option to ParquetOptions (#22649) ## Which issue does this PR close? <!-- We generally require a GitHub issue to be filed for all bug fixes and enhancements and this helps us generate change logs for our releases. You can link an issue to this PR using the GitHub syntax. For example `Closes #123` indicates that this PR will close issue #123. --> - Closes https://github.com/apache/datafusion/issues/22650. ## Rationale for this change <!-- Why are you proposing this change? If this is already explained clearly in the issue then this section is not needed. Explaining clearly why changes are proposed helps reviewers understand your changes and offer better suggestions for fixes. --> arrow-rs 58.0 added WriterProperties::set_max_row_group_bytes (PR: apache/arrow-rs#9357 Issue: apache/arrow-rs#1213), which flushes a row group when either the row-count or the byte limit is reached, whichever comes first, matching parquet-mr's parquet.block.size. DataFusion already consumes atleast this version of arrow but does not yet expose this new byte-based setter through its config. ## What changes are included in this PR? <!-- There is no need to duplicate the description in the issue here but it is sometimes worth providing a summary of the individual changes in this PR. --> - Add `max_row_group_bytes: Option<usize>` (default None) to ParquetOptions in `datafusion/common/src/config.rs`. - Wire it through `ParquetOptions::into_writer_properties_builder` to `WriterPropertiesBuilder::set_max_row_group_bytes`, with a guard that rejects Some(0) as a configuration error (arrow-rs panics on a zero byte limit). - Plumb the field through protobuf serialization - add it to the ParquetOptions proto message and the proto-common/proto conversions, with regenerated bindings. - Exposed as the max_row_group_bytes COPY / CREATE EXTERNAL TABLE format option alongside max_row_group_size. - Update the generated config docs and the format options table doc. ## Are these changes tested? <!-- We typically require tests for all PRs in order to: 1. Prevent the code from being accidentally broken by subsequent changes 2. Serve as another way to document the expected behavior of the code If tests are not included in your PR, please explain why (for example, are they covered by existing tests)? --> Yes - run locally and passing: Unit (datafusion-common, parquet_writer.rs): - defaults to None, so no byte limit is propagated to WriterProperties. - a configured value propagates to WriterProperties. - Some(0) is rejected with a configuration error. - the existing table_parquet_opts_to_writer_props round-trip and test_defaults_match tests were extended to cover the new field. Protobuf round-trip (datafusion-proto-common): - new test_parquet_options_max_row_group_bytes_round_trip confirms the option survives serialization to protobuf and back. SLTs: - new test_files/parquet_max_row_group_bytes.slt writes Parquet with the option set (via both COPY ... OPTIONS and session config), reads it back, asserts the data round-trips, and asserts a zero value is rejected. - copy.slt exercises the option inside the existing "all supported statement overrides" COPY test. - information_schema.slt updated for the new option in SHOW ALL. Commands run locally (all pass): cargo test -p datafusion-common --features parquet cargo test -p datafusion-proto-common cargo test -p datafusion-proto cargo test --test sqllogictests -- parquet_max_row_group_bytes cargo test --test sqllogictests -- information_schema cargo test --test sqllogictests -- copy ## Are there any user-facing changes? <!-- If there are user-facing changes then we may require documentation to be updated before approving the PR. --> Additive only, does not affect existing options. <!-- If there are any breaking changes to public APIs, please add the `api change` label. --> --------- Co-authored-by: Yongting You <2010youy01@gmail.com>
| Commit: | 2a5be51 | |
|---|---|---|
| Author: | Krisztián Szűcs | |
| Committer: | GitHub | |
[branch-54] refactor: give parquet CDC options an explicit `enabled` flag (backport #22632) (#22648) ## Which issue does this PR close? - Backport of #22632 to `branch-54`. ## Rationale for this change Content-defined chunking (CDC) write options were added in #21110 and are slated for the 54.0.0 release. This backports the refactor in #22632 so the config/proto surface ships in its final form, before the release goes out. The CDC options previously worked as `use_content_defined_chunking: Option<CdcOptions>` with a `ConfigField` impl that accepted a bare `use_content_defined_chunking = true|false` and otherwise enabled CDC implicitly when any sub-field was set. This has a few problems: - **Naming diverges from parquet-rs.** `WriterProperties` exposes `content_defined_chunking()` / `set_content_defined_chunking(Option<CdcOptions>)` with no `use_` prefix. - **Implicit / order-dependent on the SQL side.** Format options in `COPY ... OPTIONS` / `CREATE EXTERNAL TABLE ... OPTIONS` are applied from a `HashMap` (non-deterministic order). With the old bare-boolean form, mixing `... = false` with a sub-field could resolve to enabled or disabled depending on iteration order. - **Extra machinery.** Supporting the bare boolean required hand-written `ConfigField` impls and a `#[expect(clippy::should_implement_trait)]` workaround, plus a zero-sentinel fallback in the proto mapping. Since CDC is unreleased, the config/proto surface can still be changed freely. ## What changes are included in this PR? - Rename the `ParquetOptions` field `use_content_defined_chunking` -> `content_defined_chunking` (matches parquet-rs). - Make `CdcOptions` a plain `config_namespace!` with an explicit `enabled: bool` field alongside the chunking parameters; the field is a bare `CdcOptions` (no longer `Option<CdcOptions>`). CDC is on iff `content_defined_chunking.enabled` is true. Setting a parameter no longer implicitly enables CDC, and the result is independent of key order. - Add `CdcOptions::enabled()` / `CdcOptions::disabled()` shorthand constructors. - Drop the `ConfigField` impls and the `should_implement_trait` workaround — all generated by the macro now. - Add an `enabled` field to the proto `CdcOptions` message so the proto <-> config mapping is a plain field copy in both directions. - Update unit tests, regenerate config docs + the `information_schema` snapshot, and add `parquet_cdc_config.slt` documenting the resolution behavior. ## Are these changes tested? Yes — `datafusion-common` config + writer unit tests, `datafusion-proto-common` proto round-trip tests, `datafusion/core` parquet integration tests, and sqllogictest (`parquet_cdc.slt` + new `parquet_cdc_config.slt`). Cherry-pick applied cleanly onto `branch-54`; affected crates build and the CDC unit tests pass. ## Are there any user-facing changes? Yes, but only to the unreleased CDC options: - Config key `datafusion.execution.parquet.use_content_defined_chunking` -> `datafusion.execution.parquet.content_defined_chunking.enabled` (plus `.min_chunk_size` / `.max_chunk_size` / `.norm_level`). - The bare-boolean form is removed; enable/disable via `content_defined_chunking.enabled = true|false`. No released API is affected. 🤖 Generated with [Claude Code](https://claude.com/claude-code)
| Commit: | d88ab6a | |
|---|---|---|
| Author: | Krisztián Szűcs | |
| Committer: | GitHub | |
refactor: give parquet CDC options an explicit `enabled` flag (#22632) ## Which issue does this PR close? - None ## Rationale for this change The CDC options currently work as `use_content_defined_chunking: Option<CdcOptions>` with a `ConfigField` impl that accepts a bare `use_content_defined_chunking = true|false` and otherwise enables CDC implicitly when any sub-field is set. This has a few problems: - **Naming diverges from parquet-rs.** `WriterProperties` exposes `content_defined_chunking()` / `set_content_defined_chunking(Option<CdcOptions>)` with no `use_` prefix. - **Implicit / order-dependent on the SQL side.** Format options in `COPY ... OPTIONS` / `CREATE EXTERNAL TABLE ... OPTIONS` are applied from a `HashMap` (non-deterministic order). With the old bare-boolean form, mixing `... = false` with a sub-field, or setting a sub-field after `= false`, could resolve to enabled or disabled depending on iteration order. - **Extra machinery.** Supporting the bare boolean required a hand-written `impl ConfigField for CdcOptions` + `impl ConfigField for Option<CdcOptions>` and a `#[expect(clippy::should_implement_trait)]` workaround, plus a zero-sentinel fallback in the proto mapping. Since CDC is unreleased, the config/proto surface can still be changed freely. ## What changes are included in this PR? - Rename the `ParquetOptions` field `use_content_defined_chunking` -> `content_defined_chunking` (matches parquet-rs). - Make `CdcOptions` a plain `config_namespace!` with an explicit `enabled: bool` field alongside the chunking parameters; the field is a bare `CdcOptions` (no longer `Option<CdcOptions>`). CDC is on if `content_defined_chunking.enabled` is true. Setting a parameter no longer implicitly enables CDC, and the result is independent of key order. - Add `CdcOptions::enabled()` / `CdcOptions::disabled()` shorthand constructors. - Drop the `ConfigField` impls and the `should_implement_trait` workaround — all generated by the macro now. - Add an `enabled` field to the proto `CdcOptions` message so the proto <-> config mapping is a plain field copy in both directions (removes the presence-encoding and the zero-sentinel fallback). - Update unit tests, regenerate config docs + the `information_schema` snapshot, and add `parquet_cdc_config.slt` documenting the resolution behavior. ## Are these changes tested? Yes: - `datafusion-common` config + writer unit tests (enable toggle, parameter-does-not-enable, validation, writer round-trip). - `datafusion-proto-common` proto round-trip tests (enabled / disabled / negative norm level). - `datafusion/core` parquet integration tests (data round-trip, page boundaries). - sqllogictest: `parquet_cdc.slt` (end-to-end) and a new `parquet_cdc_config.slt` (config resolution / order independence). ## Are there any user-facing changes? Yes, but only to the unreleased CDC options: - Config key `datafusion.execution.parquet.use_content_defined_chunking` -> `datafusion.execution.parquet.content_defined_chunking.enabled` (plus `.min_chunk_size` / `.max_chunk_size` / `.norm_level`). - The bare-boolean form is removed; enable/disable via `content_defined_chunking.enabled = true|false`. No released API is affected. 🤖 Generated with [Claude Code](https://claude.com/claude-code)
| Commit: | 00c35d0 | |
|---|---|---|
| Author: | Filip Petkovski | |
| Committer: | GitHub | |
Allow specifying an arrow schema for PartitionedFile (#22360) ## Which issue does this PR close? <!-- We generally require a GitHub issue to be filed for all bug fixes and enhancements and this helps us generate change logs for our releases. You can link an issue to this PR using the GitHub syntax. For example `Closes #123` indicates that this PR will close issue #123. --> - Closes https://github.com/apache/datafusion/issues/22200. ## Rationale for this change <!-- Why are you proposing this change? If this is already explained clearly in the issue then this section is not needed. Explaining clearly why changes are proposed helps reviewers understand your changes and offer better suggestions for fixes. --> As described in the linked issue, parsing the arrow schema from parquet metadata can be expensive for point lookups, relative to the rest of the query execution pipeline. If the user knows the arrow schema of the file, they should be able to specify it explicitly. ## What changes are included in this PR? <!-- There is no need to duplicate the description in the issue here but it is sometimes worth providing a summary of the individual changes in this PR. --> * Add a `arrow_schema: SchemaRef` field to `PartitionedFile` * Use the `arrow_schema` field in the parquet opener to bypass schema inference from the `ARROW:schema` metadata field. ## Are these changes tested? Added unit tests for both matching and mismatching schemas. <!-- We typically require tests for all PRs in order to: 1. Prevent the code from being accidentally broken by subsequent changes 2. Serve as another way to document the expected behavior of the code If tests are not included in your PR, please explain why (for example, are they covered by existing tests)? --> ## Are there any user-facing changes? There are no breaking changes, the new field is optional and is set to None by default. <!-- If there are user-facing changes then we may require documentation to be updated before approving the PR. --> <!-- If there are any breaking changes to public APIs, please add the `api change` label. -->
| Commit: | 496f2c2 | |
|---|---|---|
| Author: | Adrian Garcia Badaracco | |
| Committer: | GitHub | |
feat: add pgjson format support for EXPLAIN ANALYZE (#21767) ## Which issue does this PR close? - Closes #. ## Rationale for this change DataFusion already emits PostgreSQL JSON (pgjson) for logical plans via `EXPLAIN (FORMAT pgjson) ...`. This PR extends that support to `EXPLAIN ANALYZE` so the physical plan, along with live execution metrics, can be fed into pgjson visualizers such as [Dalibo](https://explain.dalibo.com/) and PEV2. Today, `EXPLAIN ANALYZE FORMAT pgjson` is explicitly rejected in the planner with `"EXPLAIN ANALYZE with FORMAT is not supported"`. With this PR the restriction is lifted for pgjson. ## What changes are included in this PR? - Add a `format: ExplainFormat` field to the logical `Analyze` node and the physical `AnalyzeExec` operator, threaded through SQL parsing, logical planning, and physical planning. - Accept `EXPLAIN ANALYZE FORMAT pgjson <stmt>`. `Tree` and `Graphviz` with `ANALYZE` still error with a clear message (out of scope for this PR). - Add `DisplayableExecutionPlan::pgjson()` and a new `PgJsonExecutionPlanVisitor` that mirror the logical-plan `PgJsonVisitor`. Per-node output includes: - `Node Type` — `ExecutionPlan::name()` - `Details` — the one-line `DisplayAs::Default` rendering - `Actual Rows` / `Actual Total Time` — PG-canonical metric keys populated from `output_rows` / `elapsed_compute` (emitted as float milliseconds; note DataFusion records compute time, not wall time) - `Extras` — remaining DataFusion metrics keyed by their native name - `Plans` — child nodes - Add an optional `set_summary()` builder so `AnalyzeExec` can attach `Total Rows` and `Duration` at the root in verbose mode. - Honor existing `analyze_level` / `analyze_categories` config exactly as `indent()` does. - Update the `EXPLAIN` user-guide docs (`docs/source/user-guide/sql/explain.md` and `explain-usage.md`) to document pgjson support under `ANALYZE` and lead with the Postgres-style option-list spelling. ### Composes with the `EXPLAIN (...)` option list (#21768) This builds on the now-merged Postgres-style option list (#21768). Because both the keyword form and the parenthesized option list parse into a single `ExplainStatementOptions` that is threaded through `explain_to_plan`, pgjson works with **both** spellings, and the `METRICS` / `LEVEL` knobs from #21768 compose with it in one statement: ```sql EXPLAIN (ANALYZE, FORMAT pgjson) SELECT count(*) FROM t; EXPLAIN (ANALYZE, FORMAT pgjson, METRICS 'rows', LEVEL summary) SELECT count(*) FROM t; ``` The parenthesized form is the idiomatic spelling for pgjson workflows since it mirrors Postgres's `EXPLAIN (ANALYZE, FORMAT json)` — exactly what visualizers like Dalibo / PEV2 document. (Note: `ANALYZE` must go *inside* the parens; a bare `EXPLAIN ANALYZE (FORMAT pgjson)` is invalid, as it is in Postgres.) ## Are these changes tested? - Unit tests in `datafusion/physical-plan/src/display.rs`: - `pgjson_renders_plan_without_metrics` - `pgjson_includes_summary_when_set` - `pgjson_snapshot_of_sample_plan` (insta snapshot) - sqllogictest coverage in `datafusion/sqllogictest/test_files/explain_analyze.slt`: - Structural golden for `EXPLAIN (ANALYZE, FORMAT PGJSON, METRICS 'none')` (option-list form) - `EXPLAIN (ANALYZE, FORMAT PGJSON, METRICS 'rows')` showing `Actual Rows` surfacing - Keyword form `EXPLAIN ANALYZE FORMAT pgjson` still works - Negative tests for `EXPLAIN ANALYZE FORMAT tree` and `EXPLAIN ANALYZE FORMAT graphviz` - `cargo clippy --all-targets --all-features -- -D warnings` clean on the touched crates; `cargo fmt --all` clean. ## Are there any user-facing changes? Yes — `EXPLAIN ANALYZE` now accepts the `pgjson` format, in either spelling: ```sql -- Postgres-style option list (idiomatic; composes with METRICS / LEVEL) EXPLAIN (ANALYZE, FORMAT pgjson) SELECT count(*) FROM t; -- legacy keyword form EXPLAIN ANALYZE FORMAT pgjson SELECT count(*) FROM t; ``` No existing behavior changes: the default (`EXPLAIN ANALYZE ...` with no `FORMAT`) still emits the indent-format plan with metrics, and `EXPLAIN (FORMAT pgjson) ...` on the logical plan is unchanged. 🤖 Generated with [Claude Code](https://claude.com/claude-code) --------- Co-authored-by: Claude Sonnet 4.6 <noreply@anthropic.com>
| Commit: | d5643ae | |
|---|---|---|
| Author: | Adrian Garcia Badaracco | |
| Committer: | GitHub | |
feat(sql): Postgres-style `EXPLAIN (...)` option list (#21768) ## Which issue does this PR close? - Closes #. (Follow-up to #21160, which introduced per-category metric filtering via session config. This PR lets users reach those knobs inline from the EXPLAIN statement.) ## Rationale for this change #21160 added metric categories (`Rows`, `Bytes`, `Timing`, `Uncategorized`) and a verbosity level (`Summary`, `Dev`) to DataFusion's metrics, exposed today only via session config: - `datafusion.explain.analyze_categories` - `datafusion.explain.analyze_level` Users have to `SET` these out-of-band before running `EXPLAIN ANALYZE`, which is awkward for ad-hoc debugging. Postgres solves this with its parenthesized option list: ```sql EXPLAIN (ANALYZE, BUFFERS, VERBOSE, SETTINGS, WAL) SELECT ... ; ``` This PR adds the same ergonomics to DataFusion, mapping option names to DataFusion's existing semantics rather than Postgres's buffer/WAL model. ## What changes are included in this PR? **Parser.** On dialects whose `supports_explain_with_utility_options()` returns true (the default `GenericDialect`, `PostgreSqlDialect`, `DuckDbDialect`, etc.), `DFParser::parse_explain` delegates to sqlparser's `pub fn parse_utility_options()` and feeds the result through a new `ExplainStatementOptions::from_utility_options`. The legacy keyword form (`EXPLAIN ANALYZE VERBOSE FORMAT tree ...`) is unchanged. **Normalized option type.** A new `ExplainStatementOptions` in `datafusion-common` captures the knobs parsed from either form. Argument parsing reuses existing `ExplainFormat::from_str`, `ExplainAnalyzeCategories::from_str`, and `MetricType::from_str`. **Options accepted:** | Option | Argument | Effect | | --------- | ---------------- | --------------------------------------------------------------------- | | `ANALYZE` | bool, default T | Same as keyword `ANALYZE` | | `VERBOSE` | bool, default T | Same as keyword `VERBOSE` | | `FORMAT` | ident/string | `indent` / `tree` / `pgjson` / `graphviz` | | `METRICS` | string | `'all'`, `'none'`, or comma-separated `rows,bytes,timing,uncategorized` | | `LEVEL` | ident/string | `summary` or `dev` | | `TIMING` | bool | Sugar: toggles inclusion of the `timing` category | | `SUMMARY` | bool | Sugar: TRUE → `summary`, FALSE → `dev` | | `COSTS` | bool | Per-statement `show_statistics` override (not valid with `ANALYZE`) | Postgres-only options (`BUFFERS`, `WAL`, `SETTINGS`, `GENERIC_PLAN`, `MEMORY`) return a helpful unsupported-option error. **Logical plan.** `Analyze` gains `analyze_level: Option<MetricType>` and `analyze_categories: Option<ExplainAnalyzeCategories>`. `Explain` gains `show_statistics: Option<bool>`. `None` means "fall back to session config" — existing callers are unchanged. **Physical planner.** `handle_analyze` and `handle_explain` prefer statement-level overrides over session config before constructing `AnalyzeExec` / `ExplainExec`. `AnalyzeExec` itself needs no change — it already accepts the filters from #21160. **Proto.** The new override fields round-trip through `datafusion-proto`: - `datafusion_common.proto` gains `MetricType`, `MetricCategory`, and an `ExplainAnalyzeCategoriesNode` wrapper (`bool all` + `repeated MetricCategory only`, mirroring the Rust enum's `All` / `Only(Vec<…>)` variants). - `AnalyzeNode` gains `optional MetricType analyze_level` and `optional ExplainAnalyzeCategoriesNode analyze_categories`; `ExplainNode` gains `optional bool show_statistics`. - `ExplainOption` is extended with `analyze_level` / `analyze_categories` setters so the proto decode arms construct `LogicalPlan::Analyze` / `LogicalPlan::Explain` through the same `LogicalPlanBuilder::explain_option_format` path as the SQL planner. ## Are these changes tested? Yes: - **Unit tests** in `datafusion/sql/src/parser.rs` cover legacy keyword form on PostgreSQL dialect, each option form (`bare`, `= val`, `ON/OFF`, quoted), unknown-option errors, dialect gating (the parenthesized form is rejected under a dialect that doesn't enable it), and the error path for unsupported Postgres-only options. - **Integration tests** in `datafusion/core/tests/sql/explain_analyze.rs` — `explain_analyze_paren_metrics_filtering`, `explain_analyze_paren_level_overrides_session_config`, `explain_analyze_paren_metrics_overrides_session_config`, `explain_paren_buffers_rejected`. - **sqllogictest** fixtures in `datafusion/sqllogictest/test_files/explain.slt` covering the parenthesized form, round-trip with the legacy form, and each error path. - **Proto round-trip tests** in `datafusion/proto/tests/cases/roundtrip_logical_plan.rs` — `roundtrip_explain_show_statistics_override`, `roundtrip_analyze_level_override`, `roundtrip_analyze_categories_override` — cover each field set and unset, including `All`, `Only(vec![])` (plan-only), and a fully populated `Only` list. Ran `cargo fmt --all` and `cargo clippy --all-targets --all-features -- -D warnings` (clean). Two pre-existing test failures on `main` (`test_display_pg_json` snapshot and a `pgjson` SLT case at `explain.slt:642`) are unrelated to this change — verified by running them against a clean checkout of the same base commit. ## Are there any user-facing changes? Yes — new syntax. User-facing docs updated at `docs/source/user-guide/explain-usage.md` with a new section describing the option list and the dialect gate. No breaking changes: the legacy keyword form continues to work exactly as before. 🤖 Generated with [Claude Code](https://claude.com/claude-code) --------- Co-authored-by: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
| Commit: | 7a6b062 | |
|---|---|---|
| Author: | Gene Bordegaray | |
| Committer: | GitHub | |
Add Physical `Partitioning::Range` enum variant (#22207) ## Which issue does this PR close? - First mechanical PR for `ExprPartitioning` as described in thread: #21992. ## Rationale for this change DataFusion currently cannot truthfully represent range-partitioned physical data. Some sources may be range partitioned, but have to advertise another partitioning shape or fall back to unknown partitioning. This PR introduces the metadata shape for range partitioning without implementing optimizer or execution behavior yet. The goal is to establish the public representation first, then implement planning, compatibility, and execution behavior incrementally in follow-up PRs. ## What changes are included in this PR? - Adds `Partitioning::Range(RangePartitioning)`. - Adds range metadata types: - `RangePartitioning` - `RangePartition` - `RangeInterval` - `RangeBound` - Adds proto serialization/deserialization. - Adds `not_impl_err!` handling for range partitioning at call sites. - Preserves range partitioning through projection only when all partition expressions can be projected, otherwise `UnknownPartitioning`. ## Are these changes tested? Yes. ## Are there any user-facing changes? Yes. This adds new public physical partitioning API and proto for range partitioning. --------- Co-authored-by: Andrew Lamb <andrew@nerdnetworks.org>
| Commit: | cbebc6f | |
|---|---|---|
| Author: | Marc Brinkmann | |
| Committer: | GitHub | |
Fix missing field `partitioned_by_file_group` in serialization (#22365) I'm not super versed in the serialization machinery involved here, please review carefully. ## Which issue does this PR close? - Closes #22363. ## Rationale for this change The partitioned_by_file_group field was introduced in #21351 and #21342 but not added to the protobuf schema, breaking `datafusion-distributed`. ## What changes are included in this PR? - Add optional `bool partitioned_by_file_group = 14` to `FileScanExecConf` in `datafusion.proto` - Serialize the field in `to_proto.rs` - Deserialize the field in `from_proto.rs` - Regenerate prost/pbjson code ## Are these changes tested? Yes, added roundtrip_parquet_exec_partitioned_by_file_group test. ## Are there any user-facing changes? No
| Commit: | 077f08a | |
|---|---|---|
| Author: | Adrian Garcia Badaracco | |
| Committer: | GitHub | |
Split proto serialization to encapsulate private state (#21835) (#21929) ## Which issue does this PR close? - Closes #21835. ## Rationale for this change `datafusion-proto` serializes every built-in `PhysicalExpr` through a single ~300-line `downcast_ref` chain, with a mirror `match` on the decode side. That chain lives outside the crate where each expression is defined, so every field an expression wants to round-trip has to be made `pub`. #21807 is the cautionary tale: it had to add five `pub` "proto-only, not stable" items to `DynamicFilterPhysicalExpr` just to serialize an `RwLock`-wrapped inner. This PR adds the infrastructure so a `PhysicalExpr` can serialize itself and keep its state private. ## What changes are included in this PR? A `PhysicalExpr` can now opt into serializing itself, in both directions: ```rust fn try_to_proto(&self, ctx: &PhysicalExprEncodeCtx) -> Result<Option<PhysicalExprNode>> fn try_from_proto(node: &PhysicalExprNode, ctx: &PhysicalExprDecodeCtx) -> Result<Arc<dyn PhysicalExpr>> ``` `try_to_proto` returning `Ok(None)` (the default) means "fall through to the old downcast chain", so the change is purely additive — nothing is forced to migrate. `Column` and `BinaryExpr` are migrated as working demos; everything else stays on the old path and migrates later, one expression at a time, with no wire-format change. Five stacked commits, each builds green on its own and is independently reviewable (or splittable into its own PR): 1. **Extract `datafusion-proto-models` crate** — move the `.proto` file and prost-generated types into a lightweight crate (mirrors the existing `datafusion-proto-common` split). 2. **Add the `try_to_proto` hook** — feature-gated, off by default. 3. **Migrate `Column` encode.** 4. **Add the decode side and migrate `Column` decode.** 5. **Migrate `BinaryExpr`** (both directions). ## A few design decisions worth flagging - **`FromProto` / `TryFromProto` traits instead of plain `From` / `TryFrom`.** Once the prost types move into their own crate they are *foreign* to `datafusion-proto`, and the orphan rule forbids `impl From<&protobuf::X> for Y` when both `X` and `Y` are foreign. So those conversions become `FromProto` / `TryFromProto` traits in `datafusion_proto::convert`, and callers go from `(&x).into()` to `Y::from_proto(&x)`. This is a known workaround, not the end state — see Future work. - **The ctx is a concrete struct, not `&dyn`.** `PhysicalExprEncodeCtx` / `PhysicalExprDecodeCtx` wrap a sealed dispatch trait. Keeping them concrete keeps `&dyn` out of every expression's signature and gives a stable place to add helpers (UDF encoding, registry hooks) later without churning a public trait. - **`try_from_proto` takes the whole `PhysicalExprNode`**, not the pre-unwrapped variant payload, so every expression's decoder has the same signature and can still see outer-node fields like `expr_id`. ## Are these changes tested? No new behavior, so no new tests. `Column` and `BinaryExpr` produce and consume the same wire format as before; the existing `roundtrip_physical_plan` / `roundtrip_physical_expr` tests already cover both directions and now exercise the new path. ## Are there any user-facing changes? Small API breaks in `datafusion-proto`: - `try_from_physical_plan_with_converter` / `try_into_physical_plan_with_converter` move to a `PhysicalPlanNodeExt` trait — callers add `use datafusion_proto::physical_plan::PhysicalPlanNodeExt;`. - Foreign-foreign `From` / `TryFrom` conversions become `FromProto` / `TryFromProto` (see Design decisions above). - `datafusion_proto::generated::*` is deprecated in favor of `datafusion_proto::protobuf`; it still works. The new `proto` feature on `datafusion-physical-expr(-common)` is off by default, so crates that don't serialize plans pay nothing. ## Future work - Migrate the remaining built-in expressions — including `DynamicFilterPhysicalExpr`, the original motivation — one per follow-up PR. - Apply the same pattern to `ExecutionPlan` serialization. - Drop the `FromProto` / `TryFromProto` workaround: collapse `datafusion-proto-common` into `datafusion-proto-models` and push the conversion impls down to the target-type crates so callers use plain `From` / `TryFrom` again. Full dep-graph analysis and a step-by-step plan are in [#21835 (comment)](https://github.com/apache/datafusion/issues/21835#issuecomment-4348350257). 🤖 Generated with [Claude Code](https://claude.com/claude-code) --------- Co-authored-by: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
| Commit: | 4055e44 | |
|---|---|---|
| Author: | Mithun Chicklore Yogendra | |
| Committer: | GitHub | |
fix: preserve null_aware on logical JoinNode proto round-trip (#22104) ## Summary Closes #22065. `null_aware` was missing from `JoinNode` in the logical proto (it was added to the physical `HashJoinExecNode` in #19635). The encoder dropped it via `..` destructuring and the decoder had no field to restore it from, so any `to_proto` -> `from_proto` round trip silently downgraded a null-aware LeftAnti (NOT IN semantics) to a plain LeftAnti and returned wrong rows. ## Changes - Add `bool null_aware = 9;` to `JoinNode`. - Decoder switches to `Join::try_new`, plumbing `null_aware` and `null_equality` (same bug, same path) from the wire. - Encoder destructure binds `schema: _` instead of `..`, so any future `Join` field is a compile error here instead of a silent drop. - Decoder rejects mismatched `left_join_key`/`right_join_key` lengths via `proto_error`. - Regression tests `roundtrip_join_null_aware` and `roundtrip_join_null_equality`, each exercising one non-default field. ## Test plan - `cargo test -p datafusion-proto --test proto_integration cases::roundtrip_logical_plan` passes. - Clippy clean.
| Commit: | 18c347d | |
|---|---|---|
| Author: | Andy Grove | |
| Committer: | GitHub | |
feat: optional timezone for coerce_int96 (#22318) ## Which issue does this PR close? N/A ## Rationale for this change `coerce_int96_to_resolution` currently produces `Timestamp(unit, None)` for every INT96-derived column. Some downstream readers need the resulting Arrow type to carry a timezone, because the *absence* of a timezone is itself meaningful. The motivating case is Apache DataFusion Comet (a Spark accelerator) trying to enforce [SPARK-36182\: pre-Spark-4 Spark rejects reading a Parquet TimestampLTZ column as TimestampNTZ](https://issues.apache.org/jira/browse/SPARK-36182). Comet's schema adapter pattern-matches `Timestamp(_, Some(_)) -> Timestamp(_, None)` to detect this case, but for INT96 columns the post-coerce type is `Timestamp(unit, None)` — indistinguishable from a true TimestampNTZ source. The LTZ signal is destroyed at the wrong layer. Spark and other systems write INT96 as UTC-adjusted instants, so a caller can ask for the column to surface as `Timestamp(unit, Some(\"UTC\"))`, preserving the LTZ semantic at the Arrow level. ## What changes are included in this PR? - New `TableParquetOptions.global.coerce_int96_tz: Option<String>` config field (defaults to `None`). - `coerce_int96_to_resolution` gains a `timezone: Option<Arc<str>>` parameter and threads it into the constructed `Timestamp` type. - The new option is plumbed through `ParquetSource` -> `ParquetOpener` / `ParquetMorselizer` -> `DFParquetMetadata`. - `with_coerce_int96_tz` builder method on `DFParquetMetadata`. - Default behavior is unchanged when the option is unset. ## Are these changes tested? Yes, see https://github.com/apache/datafusion-comet/pull/4357 ## Are there any user-facing changes? A new \`coerce_int96_tz\` config option. No change in behavior for the default value. --------- Co-authored-by: Oleks V <comphead@users.noreply.github.com>
| Commit: | 47655fd | |
|---|---|---|
| Author: | Jayant Shrivastava | |
| Committer: | GitHub | |
proto: serialize dynamic filters on Sort, Aggregate, HashJoin plan nodes (#22011) ## Which issue does this PR close? <!-- We generally require a GitHub issue to be filed for all bug fixes and enhancements and this helps us generate change logs for our releases. You can link an issue to this PR using the GitHub syntax. For example `Closes #123` indicates that this PR will close issue #123. --> - Closes https://github.com/apache/datafusion/issues/20418 (Looks like this was accidentally closed early) - Informs: https://github.com/apache/datafusion/issues/21207#issuecomment-4254968115 ## Rationale for this change `SortExec`, `AggregateExec`, and `HashJoinExec` do not serialize their dynamic filters, so plans lose dynamic filtering when they are serialized and sent across network boundaries. ## What changes are included in this PR? This change adds `with_dynamic_filter_expr()` and `dynamic_filter_expr()` to `SortExec`, `AggregateExec`, and `HashJoinExec`. ``` pub fn with_dynamic_filter_expr( mut self, filter: Arc<DynamicFilterPhysicalExpr>, ) -> Result<Self> pub fn dynamic_filter_expr(&self) -> Option<&Arc<DynamicFilterPhysicalExpr>> { ``` This are used as getters and setters for the `proto` crate to get and set dynamic filters. ## Are these changes tested? Yes. See `datafusion/datafusion/proto/tests/cases/roundtrip_physical_plan.rs`. There are also tests for the plan nodes in the `physical-plan` crate. ## Are there any user-facing changes? `SortExec`, `AggregateExec`, and `HashJoinExec` now roundtrip serialize dynamic filter expressions. --------- Co-authored-by: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
| Commit: | aca4d13 | |
|---|---|---|
| Author: | Daniel Tu | |
| Committer: | GitHub | |
feat: Add Protobuf support for Explain node (#21994) ## Which issue does this PR close? <!-- We generally require a GitHub issue to be filed for all bug fixes and enhancements and this helps us generate change logs for our releases. You can link an issue to this PR using the GitHub syntax. For example `Closes #123` indicates that this PR will close issue #123. --> - Closes #. ## Rationale for this change <!-- Why are you proposing this change? If this is already explained clearly in the issue then this section is not needed. Explaining clearly why changes are proposed helps reviewers understand your changes and offer better suggestions for fixes. --> `EXPLAIN FORMAT TREE` is supported in logical plans, but protobuf serialization did not preserve the explain format. In Datafusion Ballista, we need the format field to generate corresponding distributed plan. https://github.com/apache/datafusion-ballista/issues/1627#issuecomment-4355988101 ## What changes are included in this PR? <!-- There is no need to duplicate the description in the issue here but it is sometimes worth providing a summary of the individual changes in this PR. --> - Add `ExplainFormat` to protobuf common definitions. - Add the `format` field to protobuf `ExplainNode`. - Regenerate protobuf code. ## Are these changes tested? Yes, we add a roundtrip test for `EXPLAIN FORMAT TREE`. <!-- We typically require tests for all PRs in order to: 1. Prevent the code from being accidentally broken by subsequent changes 2. Serve as another way to document the expected behavior of the code If tests are not included in your PR, please explain why (for example, are they covered by existing tests)? --> ## Are there any user-facing changes? No <!-- If there are user-facing changes then we may require documentation to be updated before approving the PR. --> <!-- If there are any breaking changes to public APIs, please add the `api change` label. --> --------- Co-authored-by: Kumar Ujjawal <ujjawalpathak6@gmail.com>
| Commit: | 948cd09 | |
|---|---|---|
| Author: | Jayant Shrivastava | |
| Committer: | GitHub | |
proto: serialize and dedupe dynamic filters v2 (#21807) ## Which issue does this PR close? <!-- We generally require a GitHub issue to be filed for all bug fixes and enhancements and this helps us generate change logs for our releases. You can link an issue to this PR using the GitHub syntax. For example `Closes #123` indicates that this PR will close issue #123. --> Informs: https://github.com/datafusion-contrib/datafusion-distributed/issues/180 Closes: https://github.com/apache/datafusion/issues/20418 ## Rationale for this change <!-- Why are you proposing this change? If this is already explained clearly in the issue then this section is not needed. Explaining clearly why changes are proposed helps reviewers understand your changes and offer better suggestions for fixes. --> Consider you have a plan with a `HashJoinExec` and `DataSourceExec` ``` HashJoinExec(dynamic_filter_1 on a@0) (...left side of join) ProjectionExec(a := Column("a", source_index)) DataSourceExec ParquetSource(predicate = dynamic_filter_2) ``` You serialize the plan, deserialize it, and execute it. What should happen is that the dynamic filter should "work", meaning: 1. When you deserialize the plan, both the `HashJoinExec` and `DataSourceExec` should have pointers to the same `DynamicFilterPhysicalExpr` 2. The `DynamicFilterPhysicalExpr` should be updated during execution by the `HashJoinExec` and the `DataSourceExec` should filter out rows This does not happen today for a few reasons, a couple of which this PR aims to address 1. `DynamicFilterPhysicalExpr` is not survive round-tripping. The internal exprs get inlined (ex. it may be serialized as `Literal`) due to the `PhysicalExpr::snapshot()` API 2. Even if `DynamicFilterPhysicalExpr` survives round-tripping, the one pushed down to the `DataSourceExec` often has different children. In this case, you have two `DynamicFilterPhysicalExpr` which do not survive deduping, causing referential integrity to be lost. ## What changes are included in this PR? <!-- There is no need to duplicate the description in the issue here but it is sometimes worth providing a summary of the individual changes in this PR. --> This PR aims to fix those problems by: 1. Removing the `snapshot()` call from the serialization process 2. Adding protos for `DynamicFilterPhysicalExpr` so it can be serialized and deserialized 3. Removing `Arc`-based deduplication. We now only dedupe on `expression_id` if the `PhysicalExpr` reports a `expression_id`. After this change, only `DynamicFilterPhysicalExpr` reports an `expression_id` to be deduped. 4. `expression_id` is now just a random u64. Since a given query likely only has a few `DynamicFilterPhysicalExpr` instances, the odds of a collision are very low 5. There's no need for a `DedupingSerializer` anymore since the `expression_id` is already stored in the dynamic filter proto itself Future work: 1. Serialize dynamic filters in `HashJoinExec`, `AggregateExec` and `SortExec` 2. Add tests which actually execute plans after deserialization and assert that dynamic filtering is functional 3. Add proto converters to the `PhysicalExtensionCodec` trait so implementors can utilize deduping logic ## Are these changes tested? - adds tests which roundtrip dynamic filters and assert that referential integrity is maintained - removes tests that test `Arc`-based deduplication and session id rotation since we don't support that anymore ## Are there any user-facing changes? - The default codec does not call `snapshot()` on `PhysicalExpr` during serialization anymore. This means that `DynamicFilterPhysicalExpr` are now serialized and deserialized without snapshotting. - All `PhysicalExpr` are not deduped anymore. Only `DynamicFilterPhysicalExpr` is --------- Co-authored-by: Dmitrii Blaginin <dmitrii@blaginin.me>
| Commit: | f802ed1 | |
|---|---|---|
| Author: | Oleh | |
| Committer: | GitHub | |
Add protobuf serialization/deserialization support for `EmptyTable` scans (#20844) ## Which issue does this PR close? <!-- We generally require a GitHub issue to be filed for all bug fixes and enhancements and this helps us generate change logs for our releases. You can link an issue to this PR using the GitHub syntax. For example `Closes #123` indicates that this PR will close issue #123. --> I figured it will be easier to submit PR right away as change doesn't look controversial. I'm happy to create an issue and link it here if you'd prefer. ## Rationale for this change <!-- Why are you proposing this change? If this is already explained clearly in the issue then this section is not needed. Explaining clearly why changes are proposed helps reviewers understand your changes and offer better suggestions for fixes. --> So short story is: in another project we'd like to use DataFusion's to "build" operations on data and then submit resulting logical plan _somewhere_ to execute (likely not using DF to actually execute the query). Since those plans never meant to be executed by DF we use `EmptyTable` as a base to bring schema to DF without any actual data. `EmptyTable` scans not being serializable prevents us from sending those plans to Python or over the wire. I believe this change makes datafusion's LogicalPlan more portable and more usable outside of datafusion's query executor. Longer story: [VegaFusion](https://github.com/vega/vegafusion) does server-side aggregation for Vega charts and is powered by DataFusion. We recently added option to [use custom query/plan executors](https://github.com/vega/vegafusion/pull/573), which allows user to pass a schema (without data) to VegaFusion which will add all necessary aggregations (but not execute them) and return a logical plan to user. They can then outsource this plan to custom query executor (e.g. Spark). This is already implemented and works. However, since VegaFusion is most commonly used through Python bindings, we'd like to expose this API to Python too (and additionally as part of gPRC API too) , which requires serializing built plans to protobuf. Currently we use `EmptyTable` to bring schema without any data to DataFusion. But since it can't be converted to protobuf, we're unable to expose this API. We considered providing custom decoder/encoder, but that would work only for gRPC case, but not Python as datafusion-python doesn't allow to provide custom decoder as far as I understand. ## What changes are included in this PR? <!-- There is no need to duplicate the description in the issue here but it is sometimes worth providing a summary of the individual changes in this PR. --> * Moved `EmptyTable` from `datafusion-core` into `datafusion-catalog` and added backwards compatibility re-export (following pattern for other table providers moved earlier) * Added new `EmptyTableScanNode` to protobuf definitions * Added encoding and decoding for new entity into `AsLogicalPlan for LogicalPlanNode` implementation ## Are these changes tested? <!-- We typically require tests for all PRs in order to: 1. Prevent the code from being accidentally broken by subsequent changes 2. Serve as another way to document the expected behavior of the code If tests are not included in your PR, please explain why (for example, are they covered by existing tests)? --> I added two roundtrip tests for the new node ## Are there any user-facing changes? <!-- If there are user-facing changes then we may require documentation to be updated before approving the PR. --> `EmptyTable` can be imported from `datafusion-catalog` crate now, but old crate (`datafusion-core`) still re-exports it, so this shouldn't be breaking change <!-- If there are any breaking changes to public APIs, please add the `api change` label. --> P.S. Just to be explicit, code itself was written mostly by LLM (as I'm not that proficient in Rust yet). I did review and test it though
| Commit: | 1bb588e | |
|---|---|---|
| Author: | Neil Conway | |
| Committer: | GitHub | |
perf: Implement physical execution of uncorrelated scalar subqueries (#21240) ## Which issue does this PR close? - Closes #3781. - Closes #18181. ## Rationale for this change Previously, DataFusion evaluated uncorrelated scalar subqueries by transforming them into joins. This has three shortcomings: 1. Scalar subqueries that return > 1 row were allowed, producing incorrect query results. Such queries should instead result in a runtime error. 2. Performance. Evaluating scalar subqueries as a join requires going through the join machinery. More importantly, it means that UDFs that have specialized handling of scalar inputs cannot use those code paths for scalar subqueries, which often results in significantly slower query execution (e.g., #18181). It also makes filter pushdown for scalar subquery filters more difficult (#21324) 3. Uncorrelated scalar subqueries previously did not work in `ORDER BY` or `JOIN ON`, or as arguments to an aggregate function. Those cases are now supported. This PR introduces physical execution of uncorrelated scalar subqueries: * Uncorrelated subqueries are left in the plan by the optimizer, not rewritten into joins * The physical planner collects uncorrelated scalar subqueries and plans them recursively (supporting nested subqueries). We add a `ScalarSubqueryExec` plan node to the top of any physical plan with uncorrelated subqueries: it has N+1 children, N subqueries and its "main" input, which is the rest of the query plan. The subquery expression in the parent plan is replaced with a `ScalarSubqueryExpr`. * `ScalarSubqueryExec` manages the execution of the subqueries. Subquery evaluation is done in parallel (for a given query level), but at present it happens strictly before evaluation of the parent query. This might be improved in the future (#21591). * `ScalarSubqueryExpr` reads its value from a shared slot that `ScalarSubqueryExec` populates when the subquery finishes; the physical planner assigns each subquery its slot index via `ExecutionProps`. This architecture makes it easy to avoid the shortcomings described above. Performance seems roughly unchanged (benchmarks added in this PR), but in situations like #18181, we can now leverage scalar fast-paths; in the case of #18181 specifically, this improves performance from ~800 ms to ~30 ms. ## What changes are included in this PR? * Modify subquery rewriter to not transform subqueries -> joins * Collect and plan uncorrelated scalar subqueries in the physical planner, and wire up `ScalarSubqueryExpr` * Support for subqueries in physical plan serialization/deserialization using `PhysicalProtoConverterExtension` to wire up `ScalarSubqueryExpr` correctly * Support for subqueries in logical plan serialization/deserialization * Add various SLT tests and update expected plan shapes for some tests ## Are these changes tested? Yes. New SLT coverage for cardinality errors, `ORDER BY` / `JOIN ON` / aggregate-arg contexts, nested uncorrelated subqueries, duplicate-subquery deduplication, and partition-pruning filters; new roundtrip tests for logical and physical plan serialization. ## Are there any user-facing changes? SQL: * Uncorrelated scalar subqueries that return more than one row now result in a runtime error, instead of silently producing incorrect results. * Uncorrelated scalar subqueries now work in `ORDER BY`, `JOIN ON`, and as aggregate function arguments. Rust APIs: * In `datafusion-proto`, breaking changes to `Serializeable::from_bytes_with_registry` (renamed to `from_bytes_with_ctx`), `parse_expr` / `parse_sorts` / `parse_exprs`, and the `PhysicalProtoConverterExtension` trait. Plan shape: * `LogicalPlan::Subquery` nodes will now be preserved in the logical plan * Physical plans can now contain `ScalarSubqueryExec` plan node and `ScalarSubqueryExpr` expressions The wire format has also changed to include scalar subqueries. --------- Co-authored-by: Andrew Lamb <andrew@nerdnetworks.org>
| Commit: | 85e75e2 | |
|---|---|---|
| Author: | Xander | |
| Committer: | GitHub | |
Add quote style and trimming to csv writier (#20813) ## Which issue does this PR close? - Closes https://github.com/apache/datafusion/issues/10669 Related arrow-rs PRs https://github.com/apache/arrow-rs/pull/8960 and https://github.com/apache/arrow-rs/pull/9004 ## Rationale for this change The CSV writer was missing support for `quote_style`, `ignore_leading_whitespace`, and `ignore_trailing_whitespace` options that are available on the underlying arrow `WriterBuilder`. This meant users couldn't control quoting behaviour or whitespace trimming when writing CSV files. ## What changes are included in this PR? Adds three new CSV writer options wired through the full stack: - `quote_style` — controls when fields are quoted (`Always`, `Necessary`, `NonNumeric`, `Never`). Modelled as a protobuf enum (`CsvQuoteStyle`). - `ignore_leading_whitespace` — trims leading whitespace from string values on write. - `ignore_trailing_whitespace` — trims trailing whitespace from string values on write. ## Are these changes tested? Yes — sqllogictest coverage added in `csv_files.slt` ## Are there any user-facing changes? Three new `format.*` options available in COPY TO and CREATE EXTERNAL TABLE for CSV: - `format.quote_style` (string: `Always`, `Necessary`, `NonNumeric`, `Never`) - `format.ignore_leading_whitespace` (boolean) - `format.ignore_trailing_whitespace` (boolean)
| Commit: | 8a45d02 | |
|---|---|---|
| Author: | Jeffrey Vo | |
| Committer: | GitHub | |
feat: support `ListView` and `LargeListView` in `ScalarValue` (#21669) ## Which issue does this PR close? <!-- We generally require a GitHub issue to be filed for all bug fixes and enhancements and this helps us generate change logs for our releases. You can link an issue to this PR using the GitHub syntax. For example `Closes #123` indicates that this PR will close issue #123. --> - Closes #18886 - Previous iteration: #18884 ## Rationale for this change <!-- Why are you proposing this change? If this is already explained clearly in the issue then this section is not needed. Explaining clearly why changes are proposed helps reviewers understand your changes and offer better suggestions for fixes. --> More support for listview types in the codebase ## What changes are included in this PR? <!-- There is no need to duplicate the description in the issue here but it is sometimes worth providing a summary of the individual changes in this PR. --> Added `ListView` and `LargeListView` to `ScalarValue` with all accompanying changes Support `ListView` and `LargeListView` in proto, both for the arrow datatype & the newly introduced scalarvalue variants. ## Are these changes tested? <!-- We typically require tests for all PRs in order to: 1. Prevent the code from being accidentally broken by subsequent changes 2. Serve as another way to document the expected behavior of the code If tests are not included in your PR, please explain why (for example, are they covered by existing tests)? --> Yes, added tests ## Are there any user-facing changes? <!-- If there are user-facing changes then we may require documentation to be updated before approving the PR. --> No <!-- If there are any breaking changes to public APIs, please add the `api change` label. --> --------- Co-authored-by: Khanh Duong <dqkqdlot@gmail.com> Co-authored-by: Andrew Lamb <andrew@nerdnetworks.org>
| Commit: | e5966b5 | |
|---|---|---|
| Author: | Huaijin | |
| Committer: | GitHub | |
fix: linearized operands in physical binaryexpr protobuf to avoid recursion limit (#21031) ## Which issue does this PR close? - part of #18602. ## Rationale for this change When a SQL query contains many filter conditions (e.g., 40+ `AND`/`OR` clauses in a `WHERE`), serializing the physical plan to protobuf and deserializing it fails with `DecodeError: recursion limit reached`. [This is because prost has a default recursion limit of 100](https://docs.rs/prost/latest/src/prost/lib.rs.html#30), and each `BinaryExpr` nesting consumes ~2 levels of protobuf recursion depth, so a chain of ~50 AND conditions exceeds the limit. ## What changes are included in this PR? Applied the same **linearization** approach that [logical expressions already use](https://github.com/apache/datafusion/blob/b6b542e87b84f4744096106bea0de755b2e70cc5/datafusion/proto/src/logical_plan/to_proto.rs#L228-L256) that convert a left-deep tree to linearization list. Instead of encoding a chain of same-operator binary expressions as a deeply nested tree, we flatten it into a flat `operands` list: **Before (nested, O(n) recursion depth):** ``` BinaryExpr(AND) { l: BinaryExpr(AND) { l: BinaryExpr(AND) { l: a, r: b }, r: c }, r: d } ``` **After (flat, O(1) recursion depth for the chain):** ``` BinaryExpr(AND) { operands: [a, b, c, d] } ``` ## Are these changes tested? yes, add some test case ## Are there any user-facing changes?
| Commit: | a51971b | |
|---|---|---|
| Author: | Krisztián Szűcs | |
| Committer: | GitHub | |
feat: add support for parquet content defined chunking options (#21110) ## Rationale for this change - closes https://github.com/apache/datafusion/pull/21110 Expose the new Content-Defined Chunking feature from parquet-rs https://github.com/apache/arrow-rs/pull/9450 ## What changes are included in this PR? New parquet writer options for enabling CDC. ## Are these changes tested? In-progress. ## Are there any user-facing changes? New config options. Depends on the 58.1 arrow-rs release.
| Commit: | 2c03881 | |
|---|---|---|
| Author: | Adrian Garcia Badaracco | |
| Committer: | GitHub | |
Add metric category filtering for EXPLAIN ANALYZE (#21160) ## Summary - Adds `MetricCategory` enum (`Rows`, `Bytes`, `Timing`) classifying metrics by what they measure and, critically, their **determinism**: rows/bytes are deterministic given the same plan+data; timing varies across runs. - Each `Metric` can now declare its category via `MetricBuilder::with_category()`. Well-known builder methods (`output_rows`, `elapsed_compute`, `output_bytes`, etc.) set the category automatically. Custom counters/gauges default to "always included". - New session config `datafusion.explain.analyze_categories` accepts `all` (default), `none`, or comma-separated `rows`, `bytes`, `timing`. - This is orthogonal to the existing `analyze_level` (summary/dev) which controls verbosity. ## Motivation Running `EXPLAIN ANALYZE` in `.slt` tests currently requires liberal use of `<slt:ignore>` for every non-deterministic timing metric. With this change, a test can simply: ```sql SET datafusion.explain.analyze_categories = 'rows'; EXPLAIN ANALYZE SELECT ...; -- output contains only row-count metrics — fully deterministic, no <slt:ignore> needed ``` In particular, for dynamic filters we have relatively complex integration tests that exist mostly to assert the plan shapes and state of the dynamic filters after the plan has been executed. For example #21059. With this change I think most of those can be moved to SLT tests. I've also wanted to e.g. make assertions about pruning effectiveness without having timing information included. ## Test plan - [x] New Rust integration test `explain_analyze_categories` covering all combos (rows, none, all, rows+bytes) - [x] New `.slt` tests in `explain_analyze.slt` for `rows`, `none`, `rows,bytes`, and `rows` with dev level - [x] Existing `explain_analyze` integration tests pass (24/24) - [x] Proto roundtrip test updated and passing - [x] `information_schema` slt updated for new config entry - [x] Full `core_integration` suite passes (918 tests) 🤖 Generated with [Claude Code](https://claude.com/claude-code) --------- Co-authored-by: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
| Commit: | 0dfcd97 | |
|---|---|---|
| Author: | Daniël Heres | |
| Committer: | GitHub | |
Replace ahash with foldhash for faster hashing in datafusion-common (#20958) ## Summary - Replace `ahash` with `foldhash`, hashing (`with_hashes`/`create_hashes`) - seems to auto-vectorize much better as it doesn't rely on special instructions. - Use `SeedableRandomState` for rehash paths: fold existing hash into hasher's initial state, eliminating the separate `combine_hashes` step - Add `hash_write` method to `HashValue` trait for writing values into an existing hasher - Use `valid_indices()` iterator for null paths instead of per-element `is_null()` checks - Update some code to be deterministic. Notably `RandomState::default()` now does create random seed every instance, with hash it just reused a single one. Also some group by results changed as the hash function is different, added rowsort. ## Benchmark results (int64, 8192 rows, Apple M1) | Benchmark | Before (ahash) | After (foldhash) | Improvement | |---|---|---|---| | single array, no nulls | 5.65 µs | 3.30 µs | **-42%** | | multiple arrays, no nulls | 22.15 µs | 11.19 µs | **-49%** | | single array, nulls | 11.94 µs | 9.47 µs | **-21%** | | multiple arrays, nulls | 36.92 µs | 29.80 µs | **-19%** | String view improvements (utf8_view, 8192 rows): | Benchmark | Improvement | |---|---| | single, no nulls | **-13%** | | multiple, no nulls | **-28%** | | small strings, single | **-55%** | | small strings, multiple | **-60%** | ## Test plan - [x] All 36 `hash_utils` unit tests pass - [x] Run full CI suite In some hash-heavy benchmarks (clickbench_extended) we clearly see that foldhash is faster! ``` │ QQuery 1 │ 227.74 / 228.63 ±0.76 / 229.71 ms │ 205.48 / 207.44 ±1.72 / 210.35 ms │ +1.10x faster │ │ QQuery 2 │ 541.63 / 543.65 ±1.11 / 544.82 ms │ 499.61 / 502.19 ±1.76 / 504.94 ms │ +1.08x faster │ │ QQuery 3 │ 334.85 / 336.04 ±1.16 / 337.89 ms │ 316.43 / 317.67 ±1.02 / 319.48 ms │ +1.06x faster │ ``` # Are there any user-facing changes? Yes `RandomState::with_seeds` was replaced as `RandomState::with_seeds` and there were protobuf changes for the same function. Also the function will generate different hashes, so in distributed environments it shouldn't use different versions of binaries to run the same query. 🤖 Generated with [Claude Code](https://claude.com/claude-code) --------- Co-authored-by: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
| Commit: | 2c0a38b | |
|---|---|---|
| Author: | Andrew Lamb | |
| Committer: | GitHub | |
[branch-53] ser/de fetch in FilterExec (#20738) (#20883) - Part of https://github.com/apache/datafusion/issues/19692 - Closes https://github.com/apache/datafusion/issues/20737 on branch-53 This PR: - Backports https://github.com/apache/datafusion/pull/20738 from @haohuaijin to the branch-53 line Co-authored-by: Huaijin <haohuaijin@gmail.com>
| Commit: | 4bac1cf | |
|---|---|---|
| Author: | Huaijin | |
| Committer: | GitHub | |
impl ser/de for preserve_order in RepartitionExec (#20798) ## Which issue does this PR close? <!-- We generally require a GitHub issue to be filed for all bug fixes and enhancements and this helps us generate change logs for our releases. You can link an issue to this PR using the GitHub syntax. For example `Closes #123` indicates that this PR will close issue #123. --> - Closes #20797 ## Rationale for this change - see #20797 ## What changes are included in this PR? impl ser/de for preserve_order in RepartitionExec ## Are these changes tested? add one test case ## Are there any user-facing changes? <!-- If there are user-facing changes then we may require documentation to be updated before approving the PR. --> <!-- If there are any breaking changes to public APIs, please add the `api change` label. -->
| Commit: | 15bc6bd | |
|---|---|---|
| Author: | Acfboy | |
| Committer: | GitHub | |
feat: make DefaultLogicalExtensionCodec support serialisation of buil… (#20638) ## Which issue does this PR close? <!-- We generally require a GitHub issue to be filed for all bug fixes and enhancements and this helps us generate change logs for our releases. You can link an issue to this PR using the GitHub syntax. For example `Closes #123` indicates that this PR will close issue #123. --> - Closes #16944. ## Rationale for this change Currently, the `LogicalExtensionCodec` implementation for `DefaultLogicalExtensionCodec` leaves `try_decode_file_format` / `try_encode_file_format` unimplemented (returning "not implemented" errors). However, the actual serialization logic for built-in file formats — arrow, parquet, csv, and json— already exists in their respective codec implementations. All we need to do is tag which format is being used, and delegate to the corresponding format-specific codec to handle the data. <!-- Why are you proposing this change? If this is already explained clearly in the issue then this section is not needed. Explaining clearly why changes are proposed helps reviewers understand your changes and offer better suggestions for fixes. --> ## What changes are included in this PR? Added a `FileFormatKind` enum and a `FileFormatProto` message to `datafusion.proto` to identify the file format type during transmission. Implemented `try_decode_file_format` and `try_encode_file_format` for `DefaultLogicalExtensionCodec`, which dispatch serialization/deserialization to the corresponding format-specific codec based on the format kind. Note that Avro is not covered because the upstream repository has not yet implemented the corresponding Avro codec, so Avro support is not functional at this time. ## Are these changes tested? Yes. Roundtrip tests are included for csv, json, parquet, and arrow. <!-- We typically require tests for all PRs in order to: 1. Prevent the code from being accidentally broken by subsequent changes 2. Serve as another way to document the expected behavior of the code If tests are not included in your PR, please explain why (for example, are they covered by existing tests)? --> ## Are there any user-facing changes? <!-- If there are user-facing changes then we may require documentation to be updated before approving the PR. --> <!-- If there are any breaking changes to public APIs, please add the `api change` label. -->
| Commit: | 4dbb449 | |
|---|---|---|
| Author: | Huaijin | |
| Committer: | GitHub | |
ser/de fetch in FilterExec (#20738) ## Which issue does this PR close? <!-- We generally require a GitHub issue to be filed for all bug fixes and enhancements and this helps us generate change logs for our releases. You can link an issue to this PR using the GitHub syntax. For example `Closes #123` indicates that this PR will close issue #123. --> - Closes https://github.com/apache/datafusion/issues/20737 ## Rationale for this change FilterExec have fetch filed but not impl the ser/de in proto ## What changes are included in this PR? add ser/de for fetch in FilterExec ## Are these changes tested? add one test case ## Are there any user-facing changes? <!-- If there are user-facing changes then we may require documentation to be updated before approving the PR. --> <!-- If there are any breaking changes to public APIs, please add the `api change` label. -->
| Commit: | 88fa0df | |
|---|---|---|
| Author: | Dewey Dunnington | |
| Committer: | GitHub | |
Add `Field` to `Expr::Cast` -- allow logical expressions to express a cast to an extension type (#18136) ## Which issue does this PR close? - Closes #18060. I am sorry that I missed the previous PR implementing this ( https://github.com/apache/datafusion/pull/18120 ) and I'm also happy to review that one instead of updating this! ## Rationale for this change Other systems that interact with the logical plan (e.g., SQL, Substrait) can express types that are not strictly within the arrow DataType enum. ## What changes are included in this PR? For the Cast and TryCast structs, the destination data type was changed from a DataType to a FieldRef. ## Are these changes tested? Yes. ## Are there any user-facing changes? Yes, any code using `Cast { .. }` to create an expression would need to use `Cast::new()` instead (or pass on field metadata if it has it). Existing matches will need to be upated for the `data_type` -> `field` member rename. --------- Co-authored-by: Andrew Lamb <andrew@nerdnetworks.org>
| Commit: | 5fccac1 | |
|---|---|---|
| Author: | Josh Elkind | |
| Committer: | GitHub | |
Add protoc support for ArrowScanExecNode (#20280) (#20284) ## Which issue does this PR close? - Closes #20280. ## Rationale for this change Physical plans that read Arrow files (.arrow / IPC) could not be serialized or deserialized via the proto layer. PhysicalPlanNode already had scan nodes for Parquet, CSV, JSON, Avro, and in-memory sources, but not for Arrow, so a DataSourceExec using ArrowSource was not round-trippable. That blocked use cases like distributing plans that scan Arrow files (e.g. Ballista). This change adds Arrow scan to the proto layer so those plans can be serialized and deserialized like the other file formats. ## What changes are included in this PR? Proto: Added ArrowScanExecNode (with FileScanExecConf base_conf) and arrow_scan = 38 to the PhysicalPlanNode oneof in datafusion.proto. Generated code: Updated prost.rs and pbjson.rs to include ArrowScanExecNode and the ArrowScan variant (manual edits; protoc was not run). To-proto: In try_from_data_source_exec, when the data source is a FileScanConfig whose file source is ArrowSource, it is now serialized as ArrowScanExecNode. From-proto: Implemented try_into_arrow_scan_physical_plan to deserialize ArrowScanExecNode into DataSourceExec with ArrowSource; missing base_conf returns an explicit error (no .unwrap()). Test: Added roundtrip_arrow_scan in roundtrip_physical_plan.rs to assert Arrow scan plans round-trip correctly. ## Are these changes tested? Yes. A new test roundtrip_arrow_scan builds a physical plan that scans Arrow files, serializes it to bytes and deserializes it back, and asserts the round-tripped plan matches the original. The full cargo test -p datafusion-proto suite (150 tests: unit, integration, and doc tests) passes, including all existing roundtrip and serialization tests. ## Are there any user-facing changes? No. This only extends the existing physical-plan proto support to Arrow scan. Callers that already serialize/deserialize physical plans (e.g. for distributed execution) can now round-trip plans that read Arrow files in addition to Parquet, CSV, JSON, and Avro, with no API or behavioral changes for existing usage.
| Commit: | 69d0f44 | |
|---|---|---|
| Author: | Qi Zhu | |
| Committer: | GitHub | |
Support JSON arrays reader/parse for datafusion (#19924) ## Which issue does this PR close? Closes #19920 ## Rationale for this change DataFusion currently only supports line-delimited JSON (NDJSON) format. Many data sources provide JSON in array format `[{...}, {...}]`, which cannot be parsed by the existing implementation. ## What changes are included in this PR? - Add `newline_delimited` option to `JsonOptions` (default `true` for backward compatibility) - Implement streaming JSON array to NDJSON conversion via `JsonArrayToNdjsonReader` - Support both file-based and stream-based (e.g., S3) reading with memory-efficient streaming - Add `ChannelReader` for async-to-sync byte transfer in object store streaming scenarios - Add protobuf serialization support for the new option - Rename `NdJsonReadOptions` to `JsonReadOptions` (with deprecation alias) - SQL support via `OPTIONS ('format.newline_delimited' 'false')` ### Architecture ```text JSON Array File (e.g., 33GB) │ ▼ read chunks via ChannelReader (for streams) or BufReader (for files) ┌───────────────────┐ │ JsonArrayToNdjson │ ← streaming character substitution: │ Reader │ '[' skip, ',' → '\n', ']' stop └───────────────────┘ │ ▼ outputs NDJSON format ┌───────────────────┐ │ Arrow Reader │ ← batch parsing └───────────────────┘ │ ▼ RecordBatch ``` ### Memory Efficiency | Approach | Memory for 33GB file | Parse count | |----------|---------------------|-------------| | Load entire file + serde_json | ~100GB+ | 3x | | Streaming with JsonArrayToNdjsonReader | ~32MB | 1x | ## Are these changes tested? Yes: - Unit tests for `JsonArrayToNdjsonReader` (nested objects, escaped strings, empty arrays, buffer boundaries) - Unit tests for `ChannelReader` - Integration tests for `JsonOpener` (file-based, stream-based, large files, cancellation) - Schema inference tests (normal, empty, nested struct, list types) - End-to-end query tests with SQL - SQLLogicTest for SQL validation ## Are there any user-facing changes? Yes. Users can now read JSON array format files: **Via SQL:** ```sql CREATE EXTERNAL TABLE my_table STORED AS JSON OPTIONS ('format.newline_delimited' 'false') LOCATION 'path/to/array.json'; ``` **Via API:** ```rust let options = JsonReadOptions::default().newline_delimited(false); ctx.register_json("my_table", "path/to/array.json", options).await?; ``` **Note:** `NdJsonReadOptions` is deprecated in favor of `JsonReadOptions`. **Limitation:** JSON array format does not support range-based file scanning (`repartition_file_scans`). Users will see a clear error message if this is attempted.
| Commit: | 81f7a87 | |
|---|---|---|
| Author: | Gabriel | |
| Committer: | GitHub | |
Add BufferExec execution plan (#19760) ## Which issue does this PR close? <!-- We generally require a GitHub issue to be filed for all bug fixes and enhancements and this helps us generate change logs for our releases. You can link an issue to this PR using the GitHub syntax. For example `Closes #123` indicates that this PR will close issue #123. --> - Closes #. ## Rationale for this change This is a PR from a batch of PRs that attempt to improve performance in hash joins: - https://github.com/apache/datafusion/pull/19759 - This PR - https://github.com/apache/datafusion/pull/19761 It adds a building block that allows eagerly collecting data on the probe side of a hash join before the build side is finished. Even if the intended use case is for hash joins, the new execution node is generic and is designed to work anywhere in the plan. ## What changes are included in this PR? > [!NOTE] > The new BufferExec node introduced in this PR is still not wired up automatically <!-- There is no need to duplicate the description in the issue here but it is sometimes worth providing a summary of the individual changes in this PR. --> Adds a new `BufferExec` node that can buffer up to a certain size in bytes for each partition eagerly performing work that otherwise would be delayed. Schematically, it looks like this: ``` ┌───────────────────────────┐ │ BufferExec │ │ │ │┌────── Partition 0 ──────┐│ ││ ┌────┐ ┌────┐││ ┌────┐ ──background poll────────▶│ │ │ ├┼┼───────▶ │ ││ └────┘ └────┘││ └────┘ │└─────────────────────────┘│ │┌────── Partition 1 ──────┐│ ││ ┌────┐ ┌────┐ ┌────┐││ ┌────┐ ──background poll─▶│ │ │ │ │ ├┼┼───────▶ │ ││ └────┘ └────┘ └────┘││ └────┘ │└─────────────────────────┘│ │ │ │ ... │ │ │ │┌────── Partition N ──────┐│ ││ ┌────┐││ ┌────┐ ──background poll───────────────▶│ ├┼┼───────▶ │ ││ └────┘││ └────┘ │└─────────────────────────┘│ └───────────────────────────┘ ``` ## Are these changes tested? <!-- We typically require tests for all PRs in order to: 1. Prevent the code from being accidentally broken by subsequent changes 2. Serve as another way to document the expected behavior of the code If tests are not included in your PR, please explain why (for example, are they covered by existing tests)? --> yes, by new unit tests ## Are there any user-facing changes? <!-- If there are user-facing changes then we may require documentation to be updated before approving the PR. --> users can import a new `BufferExec` execution plan in their codebase, but no internal usage is shipped yet in this PR. <!-- If there are any breaking changes to public APIs, please add the `api change` label. -->
| Commit: | 39da29f | |
|---|---|---|
| Author: | Jeffrey Vo | |
| Committer: | GitHub | |
Add `ScalarValue::RunEndEncoded` variant (#19895) ## Which issue does this PR close? <!-- We generally require a GitHub issue to be filed for all bug fixes and enhancements and this helps us generate change logs for our releases. You can link an issue to this PR using the GitHub syntax. For example `Closes #123` indicates that this PR will close issue #123. --> - Closes #18563 ## Rationale for this change <!-- Why are you proposing this change? If this is already explained clearly in the issue then this section is not needed. Explaining clearly why changes are proposed helps reviewers understand your changes and offer better suggestions for fixes. --> Support RunEndEncoded scalar values, similar to how we support for Dictionary. ## What changes are included in this PR? <!-- There is no need to duplicate the description in the issue here but it is sometimes worth providing a summary of the individual changes in this PR. --> - Add new `ScalarValue::RunEndEncoded` enum variant - Fix `ScalarValue::new_default` to support `Decimal32` and `Decimal64` - Support RunEndEncoded type in proto for both `ScalarValue` message and `ArrowType` message ## Are these changes tested? <!-- We typically require tests for all PRs in order to: 1. Prevent the code from being accidentally broken by subsequent changes 2. Serve as another way to document the expected behavior of the code If tests are not included in your PR, please explain why (for example, are they covered by existing tests)? --> Added tests. ## Are there any user-facing changes? <!-- If there are user-facing changes then we may require documentation to be updated before approving the PR. --> New variant for `ScalarValue` Protobuf changes to support RunEndEncoded type <!-- If there are any breaking changes to public APIs, please add the `api change` label. --> --------- Co-authored-by: Andrew Lamb <andrew@nerdnetworks.org>
| Commit: | 66ee0af | |
|---|---|---|
| Author: | Adrian Garcia Badaracco | |
| Committer: | GitHub | |
Preserve PhysicalExpr graph in proto round trip using Arc pointers as unique identifiers (#20037) Replaces #18192 using the APIs in #19437. Similar to #18192 the end goal here is specifically to enable deduplication of `DynamicFilterPhysicalExpr` so that distributed query engines can get one step closer to using dynamic filters. Because it's actually simpler we apply this deduplication to all `PhysicalExpr`s with the added benefit that we more faithfully preserve the original expression tree (instead of adding new duplicate branches) which will have the immediate impact of e.g. not duplicating large `InListExpr`s.
| Commit: | 7c3ea05 | |
|---|---|---|
| Author: | Nathaniel J. Smith | |
| Committer: | GitHub | |
feat: add AggregateMode::PartialReduce for tree-reduce aggregation (#20019) DataFusion's current `AggregateMode` enum has four variants covering three of the four cells in the input/output matrix: | | Input: raw data | Input: partial state | | - | - | - | | Output: final values | `Single` / `SinglePartitioned` | `Final` / `FinalPartitioned` | | Output: partial state | `Partial` | ??? | This PR adds `AggregateMode::PartialReduce` to fill in the missing cell: it takes partially-reduced values as input, and reduces them further, but without finalizing. This is useful because it's the key component needed to implement distributed tree-reduction (as seen in e.g. the Scuba or Honeycomb papers): a set of worker nodes each perform multithreaded `Partial` aggregations, feed those into a `PartialReduce` to reduce all of this node's values into a single row, and then a head node collects the outputs from all nodes' `PartialReduce` to feed into a `Final` reduction. PR can be reviewed commit by commit: first commit is pure refactor/simplification; most places we were matching on `AggregateMode` we were actually just trying to either check which row of the above table we were in, or else which column. So now we have `is_first_stage` (tells you which column) and `is_last_stage` (tells you which row) and we use them everywhere. Second commit adds `PartialReduce`, and is pretty small because `is_first_stage`/`is_last_stage` do most of the heavy lifting. It also adds a test demonstrating a minimal Partial -> PartialReduce -> Final tree-reduction.
| Commit: | 36c0cda | |
|---|---|---|
| Author: | Kumar Ujjawal | |
| Committer: | GitHub | |
fix: respect DataFrameWriteOptions::with_single_file_output for paths without extensions (#19931) ## Which issue does this PR close? <!-- We generally require a GitHub issue to be filed for all bug fixes and enhancements and this helps us generate change logs for our releases. You can link an issue to this PR using the GitHub syntax. For example `Closes #123` indicates that this PR will close issue #123. --> - Closes #13323. ## Rationale for this change When using `DataFrameWriteOptions::with_single_file_output(true)`, the setting was being ignored if the output path didn't have a file extension. For example: ```rust df.write_parquet("/path/to/output", DataFrameWriteOptions::new().with_single_file_output(true), None).await?; ``` Would create a directory /path/to/output/ with files inside instead of a single file at /path/to/output. This happened because the demuxer used a heuristic based solely on file extension, ignoring the explicit user setting. <!-- Why are you proposing this change? If this is already explained clearly in the issue then this section is not needed. Explaining clearly why changes are proposed helps reviewers understand your changes and offer better suggestions for fixes. --> ## What changes are included in this PR? - New FileOutputMode enum: Uses explicit modes (Automatic, SingleFile, Directory) in FileSinkConfig for clearer output path handling. - The demuxer now uses the user's explicit setting instead of always relying on extension-based heuristics. <!-- There is no need to duplicate the description in the issue here but it is sometimes worth providing a summary of the individual changes in this PR. --> ## Are these changes tested? - New unit test test_single_file_output_without_extension tests the fixed behavior - All sqllogictest pass <!-- We typically require tests for all PRs in order to: 1. Prevent the code from being accidentally broken by subsequent changes 2. Serve as another way to document the expected behavior of the code If tests are not included in your PR, please explain why (for example, are they covered by existing tests)? --> ## Are there any user-facing changes? Breaking for direct FileSinkConfig construction: The struct now requires file_output_mode: FileOutputMode field. Use FileOutputMode::Automatic to preserve existing behavior. <!-- If there are user-facing changes then we may require documentation to be updated before approving the PR. --> <!-- If there are any breaking changes to public APIs, please add the `api change` label. --> --------- Co-authored-by: Andrew Lamb <andrew@nerdnetworks.org>
| Commit: | 35e99b9 | |
|---|---|---|
| Author: | Albert Skalt | |
| Committer: | GitHub | |
preserve FilterExec batch size during ser/de (#19960) ## Rationale for this change Noticed that `FilterExec` batch size is not preserved so it is set to default one after plan serialization + de-serialization. This patch fixes it. ## What changes are included in this PR? <!-- There is no need to duplicate the description in the issue here but it is sometimes worth providing a summary of the individual changes in this PR. -->
| Commit: | e82dc21 | |
|---|---|---|
| Author: | Rosai | |
| Committer: | GitHub | |
Feat : added truncate table support (#19633) ## Which issue does this PR close? - Related to #19617 ## Rationale for this change DataFusion recently added TableProvider hooks for row-level DML operations such as DELETE and UPDATE, but TRUNCATE TABLE was still unsupported. ## What changes are included in this PR? This PR adds planning and integration support for TRUNCATE TABLE in DataFusion, completing another part of the DML surface alongside existing DELETE and UPDATE support. Specifically, it includes: - SQL parsing support for TRUNCATE TABLE - Logical plan support via a new WriteOp::Truncate DML operation - Physical planner routing for TRUNCATE statements - A new TableProvider::truncate() hook for storage-native implementations - Protobuf / DML node support for serializing and deserializing TRUNCATE operations - SQL logic tests validating logical and physical planning behavior The implementation follows the same structure and conventions as the existing DELETE and UPDATE DML support. Execution semantics are delegated to individual TableProvider implementations via the new hook. ## Are these changes tested? Yes. The PR includes: SQL logic tests that verify: - Parsing of TRUNCATE TABLE - Correct logical plan generation - Correct physical planner routing - Clear and consistent errors for providers that do not yet support TRUNCATE These tests mirror the existing testing strategy used for unsupported DELETE and UPDATE operations. ## Are there any user-facing changes? Yes. Users can now execute TRUNCATE TABLE statements in DataFusion for tables whose TableProvider supports the new truncate() hook. Tables that do not support TRUNCATE will return a clear NotImplemented error.
| Commit: | 0aab6a3 | |
|---|---|---|
| Author: | Huaijin | |
| Committer: | GitHub | |
feat: support `SELECT DISTINCT id FROM t ORDER BY id LIMIT n` query use GroupedTopKAggregateStream (#19653) ## Which issue does this PR close? - close https://github.com/apache/datafusion/issues/19638 ## Rationale for this change see issue #19638 ## What changes are included in this PR? 1. Introduced `LimitOptions` struct limit field with both `limit` and optional `descending` ordering direction 2. Extended `TopKAggregation` optimizer rule to DISTINCT queries by recognizing `GROUP BY` queries without aggregates and setting the `descending` flag based on ordering direction 3. Enhanced `GroupedTopKAggregateStream` to handle DISTINCT by using group key as both priority queue key and value for DISTINCT operations 4. Updated Proto definitions to add optional `descending` field to `AggLimit` message for serialization/deserialization ## benchmark result <img width="731" height="475" alt="image" src="https://github.com/user-attachments/assets/05b6eb8c-186d-4b17-84a9-a2897dbcb095" /> ## Are these changes tested? yes, add test case in aggregates_topk.slt ## Are there any user-facing changes? no
| Commit: | 4c67d02 | |
|---|---|---|
| Author: | Liang-Chi Hsieh | |
| Committer: | GitHub | |
feat: Add null-aware anti join support (#19635) ## Which issue does this PR close? <!-- We generally require a GitHub issue to be filed for all bug fixes and enhancements and this helps us generate change logs for our releases. You can link an issue to this PR using the GitHub syntax. For example `Closes #123` indicates that this PR will close issue #123. --> - Closes #10583. ## Rationale for this change <!-- Why are you proposing this change? If this is already explained clearly in the issue then this section is not needed. Explaining clearly why changes are proposed helps reviewers understand your changes and offer better suggestions for fixes. --> ## What changes are included in this PR? <!-- There is no need to duplicate the description in the issue here but it is sometimes worth providing a summary of the individual changes in this PR. --> This patch implements null-aware anti join support for HashJoin LeftAnti operations, enabling correct SQL NOT IN subquery semantics with NULL values. ## Are these changes tested? <!-- We typically require tests for all PRs in order to: 1. Prevent the code from being accidentally broken by subsequent changes 2. Serve as another way to document the expected behavior of the code If tests are not included in your PR, please explain why (for example, are they covered by existing tests)? --> ## Are there any user-facing changes? <!-- If there are user-facing changes then we may require documentation to be updated before approving the PR. --> <!-- If there are any breaking changes to public APIs, please add the `api change` label. --> --------- Co-authored-by: Claude Sonnet 4.5 <noreply@anthropic.com>
| Commit: | 91cfb69 | |
|---|---|---|
| Author: | Adrian Garcia Badaracco | |
| Committer: | GitHub | |
feat(proto): Add protobuf serialization for HashExpr (#19379) ## Summary This PR adds protobuf serialization/deserialization support for `HashExpr`, enabling distributed query execution to serialize hash expressions used in hash joins and repartitioning. This is a followup to #18393 which introduced `HashExpr` but did not add serialization support. This causes errors when serialization is triggered on a query that pushes down dynamic filters from a `HashJoinExec`. As of #18393 `HashJoinExec` produces filters of the form: ```sql CASE (hash_repartition % 2) WHEN 0 THEN a >= ab AND a <= ab AND b >= bb AND b <= bb AND hash_lookup(a,b) WHEN 1 THEN a >= aa AND a <= aa AND b >= ba AND b <= ba AND hash_lookup(a,b) ELSE FALSE END ``` Where `hash_lookup` is an expression that holds a reference to a given partitions hash join hash table and will check for membership. Since we created these new expressions but didn't make any of them serializable any attempt to do a distributed query or similar would run into errors. In https://github.com/apache/datafusion/pull/19300 we fixed `hash_lookup` by replacing it with `true` since it can't be serialized across the wire (we'd have to send the entire hash table). The logic was that this preserves the bounds checks, which as still valuable. This PR handles `hash_repartition` which determines which partition (and hence which branch of the `CASE` expression) the row belongs to. For this expression we *can* serialize it, so that's what I'm doing in this PR. ### Key Changes - **SeededRandomState wrapper**: Added a `SeededRandomState` struct that wraps `ahash::RandomState` while preserving the seeds used to create it. This is necessary because `RandomState` doesn't expose seeds after creation, but we need them for serialization. - **Updated seed constants**: Changed `HASH_JOIN_SEED` and `REPARTITION_RANDOM_STATE` constants to use `SeededRandomState` instead of raw `RandomState`. - **HashExpr enhancements**: - Changed `HashExpr` to use `SeededRandomState` - Added getter methods: `on_columns()`, `seeds()`, `description()` - Exported `HashExpr` and `SeededRandomState` from the joins module - **Protobuf support**: - Added `PhysicalHashExprNode` message to `datafusion.proto` with fields for `on_columns`, seeds (4 `u64` values), and `description` - Implemented serialization in `to_proto.rs` - Implemented deserialization in `from_proto.rs` ## Test plan - [x] Added roundtrip test in `roundtrip_physical_plan.rs` that creates a `HashExpr`, serializes it, deserializes it, and verifies the result - [x] All existing hash join tests pass (583 tests) - [x] All proto roundtrip tests pass 🤖 Generated with [Claude Code](https://claude.com/claude-code) --------- Co-authored-by: Claude Opus 4.5 <noreply@anthropic.com>
| Commit: | 14cd71e | |
|---|---|---|
| Author: | Smotrov Oleksii | |
| Committer: | GitHub | |
feat: add compression level configuration for JSON/CSV writers (#18954) ## Which issue does this PR close? Closes #18947 ## Rationale for this change Currently, DataFusion uses default compression levels when writing compressed JSON and CSV files. For ZSTD, this means level 3, which prioritizes speed over compression ratio. Users working with large datasets who want to optimize for storage costs or network transfer have no way to increase the compression level. This is particularly important for cloud data lake scenarios where storage and egress costs can be significant. ## What changes are included in this PR? - Add `compression_level: Option<u32>` field to `JsonOptions` and `CsvOptions` in `config.rs` - Add `convert_async_writer_with_level()` method to `FileCompressionType` (non-breaking API extension) - Keep original `convert_async_writer()` as a convenience wrapper for backward compatibility - Update `JsonWriterOptions` and `CsvWriterOptions` with `compression_level` field - Update `ObjectWriterBuilder` to support compression level - Update JSON and CSV sinks to pass compression level through the write pipeline - Update proto definitions and conversions for serialization support - Fix unrelated unused import warning in `udf.rs` (conditional compilation for debug-only imports) ## Are these changes tested? The changes follow the existing patterns used throughout the codebase. The implementation was verified by: - Building successfully with `cargo build` - Running existing tests with `cargo test --package datafusion-proto` - All 131 proto integration tests pass ## Are there any user-facing changes? Yes, users can now specify compression level when writing JSON/CSV files: ```rust use datafusion::common::config::JsonOptions; use datafusion::common::parsers::CompressionTypeVariant; let json_opts = JsonOptions { compression: CompressionTypeVariant::ZSTD, compression_level: Some(9), // Higher compression ..Default::default() }; ``` **Supported compression levels:** - ZSTD: 1-22 (default: 3) - GZIP: 0-9 (default: 6) - BZIP2: 1-9 (default: 9) - XZ: 0-9 (default: 6) **This is a non-breaking change** - the original `convert_async_writer()` method signature is preserved for backward compatibility. Co-authored-by: Andrew Lamb <andrew@nerdnetworks.org>
| Commit: | c53a448 | |
|---|---|---|
| Author: | kosiew | |
| Committer: | GitHub | |
Fix panic for `GROUPING SETS(())` and handle empty-grouping aggregates (#19252) ## Which issue does this PR close? * Closes #18974. ## Rationale for this change The DataFusion CLI currently panics with an "index out of bounds" error when executing queries that use `GROUP BY GROUPING SETS(())`, such as: ```sql SELECT SUM(v1) FROM generate_series(10) AS t1(v1) GROUP BY GROUPING SETS(()) ``` This panic originates in the physical aggregation code, which assumes that an empty list of grouping expressions always corresponds to "no grouping". That assumption breaks down in the presence of `GROUPING SETS`, where an empty set is a valid grouping set that should still produce a result row (and `__grouping_id`) rather than crashing. This PR fixes the panic by explicitly distinguishing: * true "no GROUP BY" aggregations, and * `GROUPING SETS`/`CUBE`/`ROLLUP` plans that may have empty grouping expressions but still require grouping-set semantics and a valid `__grouping_id`. The change restores robustness of the CLI and ensures standards-compliant behavior for grouping sets with empty sets. ## What changes are included in this PR? Summary of the main changes: * **Track grouping-set usage explicitly in `PhysicalGroupBy`:** * Add a `has_grouping_set: bool` field to `PhysicalGroupBy`. * Extend `PhysicalGroupBy::new` to accept the `has_grouping_set` flag. * Add helper methods: * `has_grouping_set(&self) -> bool` to expose the flag, and * `is_true_no_grouping(&self) -> bool` to represent the case of genuinely no grouping (no GROUP BY and no grouping sets). * **Correct group state construction for empty grouping with grouping sets:** * Update `PhysicalGroupBy::from_pre_group` so that it only treats `expr.is_empty()` as "no groups" when `has_grouping_set` is `false`. * For `GROUPING SETS(())`, we now build at least one group, avoiding the previous out-of-bounds access on `groups[0]`. * **Clarify when `__grouping_id` should be present:** * Replace the previous `is_single` logic with a clearer distinction based on `has_grouping_set`. * `num_output_exprs`, `output_exprs`, `num_group_exprs`, and `group_schema` now add the `__grouping_id` column only when `has_grouping_set` is `true`. * `is_single` is redefined as "simple GROUP BY" (no grouping sets), i.e. `!self.has_grouping_set`. * **Integrate the new semantics into `AggregateExec`:** * Use `group_by.is_true_no_grouping()` instead of `group_by.expr.is_empty()` when choosing between the specialized no-grouping aggregation path and grouped aggregation. * Ensure that `is_unordered_unfiltered_group_by_distinct` only treats plans as grouped when there are grouping expressions **and** no grouping sets (`!has_grouping_set`). * Preserve existing behavior for regular `GROUP BY` while correctly handling `GROUPING SETS` and related constructs. * **Support `__grouping_id` with the no-grouping aggregation stream:** * Extend `AggregateStreamInner` with an optional `grouping_id: Option<ScalarValue>` field. * Change `AggregateStream::new` to accept a `grouping_id` argument. * Introduce `prepend_grouping_id_column` to prepend a `__grouping_id` column to the finalized accumulator output when needed. * Wire this up so that no-grouping aggregations can still match a schema that includes `__grouping_id` in grouping-set scenarios. * **Planner and execution wiring updates:** * Update all `PhysicalGroupBy::new` call sites to pass the correct `has_grouping_set` value: * `false` for: * ordinary `GROUP BY` or truly no-grouping aggregates. * `true` for: * `GROUPING SETS`, * `CUBE`, and * `ROLLUP` physical planning paths. * Ensure `merge_grouping_set_physical_expr`, `create_cube_physical_expr`, and `create_rollup_physical_expr` correctly mark grouping-set plans. * **Protobuf / physical plan round-trip support:** * Extend `AggregateExecNode` in `datafusion.proto` with a new `bool has_grouping_set = 12;` field. * Update the generated `pbjson` and `prost` code to serialize and deserialize the new field. * When constructing `AggregateExec` from protobuf, pass the decoded `has_grouping_set` into `PhysicalGroupBy::new`. * When serializing an `AggregateExec` back to protobuf, set `has_grouping_set` based on `exec.group_expr().has_grouping_set()`. * Update round-trip physical plan tests to include the new field in their expectations. * **Tests and SQL logic coverage:** * Add sqllogictests for the previously failing cases in `grouping.slt`: * `SELECT COUNT(*) FROM test GROUP BY GROUPING SETS (());` * `SELECT SUM(v1) FROM generate_series(10) AS t1(v1) GROUP BY GROUPING SETS(())` (the original panic case). * Extend or adjust unit tests in `aggregates`, `physical_planner`, `filter_pushdown`, and `coop` modules to account for the `has_grouping_set` flag in `PhysicalGroupBy` and expected debug output. * Update proto round-trip tests to validate `has_grouping_set` is preserved. ## Are these changes tested? Yes. * New sqllogictests covering `GROUPING SETS(())` for both a regular table and `generate_series(10)`: * `grouping.slt` now asserts the expected scalar results (e.g. `2` and `55`), preventing regressions on this edge case. * Updated and existing Rust unit tests: * `physical-plan/src/aggregates` tests updated to include `has_grouping_set` in `PhysicalGroupBy` expectations. * Planner and optimizer tests (e.g. `physical_planner.rs`, `filter_pushdown`) updated to construct `PhysicalGroupBy` with the new flag. * Execution tests in `core/tests/execution/coop.rs` updated to reflect the new constructor and continue to exercise the no-grouping aggregation path. * Protobuf round-trip tests extended to verify that `has_grouping_set` is correctly serialized and deserialized. These tests collectively ensure that: * the panic is fixed, * the aggregation semantics for `GROUPING SETS(())` are correct, and * existing aggregate behavior remains unchanged for non-grouping-set queries. ## Are there any user-facing changes? Yes, but they are bug fixes and behavior clarifications rather than breaking changes: * Queries using `GROUP BY GROUPING SETS(())` no longer cause a runtime panic in the DataFusion CLI. * Instead, they return the expected single aggregate row (e.g. `COUNT(*)` or `SUM(v1)`), consistent with SQL semantics. * For plans using `GROUPING SETS`, `CUBE`, or `ROLLUP`, the internal `__grouping_id` column is now present consistently whenever grouping sets are in use, even when the grouping expressions are empty. * For ordinary `GROUP BY` queries that do not use grouping sets, behavior is unchanged: no unexpected `__grouping_id` column is added. No API signatures were changed in a breaking way for downstream users; the additions are internal flags and protobuf fields to accurately represent the physical plan. ## LLM-generated code disclosure This PR includes LLM-generated code and comments. All LLM-generated content has been manually reviewed and tested.
| Commit: | c1aa1b5 | |
|---|---|---|
| Author: | Adrian Garcia Badaracco | |
| Committer: | GitHub | |
Track column sizes in Statistics; propagate through projections (#19113) Closes #19098, follow up to #19094. Related to #14936
| Commit: | ca67edc | |
|---|---|---|
| Author: | David Stancu | |
| Committer: | GitHub | |
[Proto]: Serialization support for `AsyncFuncExec` (#19118) ## Which issue does this PR close? Closes #19112 ## Rationale for this change Use async functions with Ballista ## What changes are included in this PR? - New `AsyncFuncExecNode` proto definition - to/from glue - A roundtrip test ## Are these changes tested? n/a ## Are there any user-facing changes? n/a
| Commit: | 9af6858 | |
|---|---|---|
| Author: | Andrew Lamb | |
| Committer: | GitHub | |
Add `force_filter_selections` to restore `pushdown_filters` behavior prior to parquet 57.1.0 upgrade (#19003) ~Draft until https://github.com/apache/datafusion/pull/18820 is merged~ ## Which issue does this PR close? - Follow on to https://github.com/apache/datafusion/pull/18820 ## Rationale for this change The parquet 57.1.0 upgrade includes a new adaptive filter from @hhhizzz : - https://github.com/apache/arrow-rs/pull/8733 Our testing shows this is faster in all cases, but I want to have an escape valve for people to turn it off if they hit some issue. I had originally included this in #18820 but @rluvaton suggested it would be easier to understand as its own PR in https://github.com/apache/datafusion/pull/18820#pullrequestreview-3509993052 ## What changes are included in this PR? 1. Add a `force_filter_selections` config setting 2. Add configuration guide 3. Add tests ## Are these changes tested? Yes ## Are there any user-facing changes? A new boolean flag
| Commit: | 9f725d9 | |
|---|---|---|
| Author: | Adrian Garcia Badaracco | |
| Committer: | GitHub | |
move projection handling into FileSource (#18627) - Part of https://github.com/apache/datafusion/issues/14993 This moves ownership of projections from `FileScanConfig` into `FileSource`. Notably we do *not* do anything special with this in Parquet just yet: I leave it for a followup to actually use the projection expressions instead of column indices to e.g. generate the Parquet `ProjectionMask` directly from expressions (in particular to select leaves instead of roots for struct and variant access).
| Commit: | 82b1307 | |
|---|---|---|
| Author: | Dewey Dunnington | |
| Committer: | GitHub | |
Enable placeholders with extension types (#17986) ## Which issue does this PR close? - Closes #17862 ## Rationale for this change Most logical plan expressions now propagate metadata; however, parameters with extension types or other field metadata cannot participate in placeholder/parameter binding. ## What changes are included in this PR? The DataType in the Placeholder struct was replaced with a FieldRef along with anything that stored the "DataType" of a parameter. Strictly speaking one could bind parameters with an extension type by copy/pasting the placeholder replacer, which I figured out towards the end of this change. I still think this change makes sense and opens up the door for things like handling UUID in SQL with full parameter binding support. ## Are these changes tested? Yes ## Are there any user-facing changes? Yes, one new function was added to extract the placeholder fields from a plan. This is a breaking change for code that specifically interacts with the pub fields of the modified structs (ParamValues, Placeholder, and Prepare are the main ones). --------- Co-authored-by: Andrew Lamb <andrew@nerdnetworks.org>
| Commit: | cadf429 | |
|---|---|---|
| Author: | Khanh Duong | |
| Committer: | GitHub | |
feat: support `null_treatment`, `distinct`, and `filter` for window functions in proto (#18024) ## Which issue does this PR close? <!-- We generally require a GitHub issue to be filed for all bug fixes and enhancements and this helps us generate change logs for our releases. You can link an issue to this PR using the GitHub syntax. For example `Closes #123` indicates that this PR will close issue #123. --> - Closes #17417. ## Rationale for this change <!-- Why are you proposing this change? If this is already explained clearly in the issue then this section is not needed. Explaining clearly why changes are proposed helps reviewers understand your changes and offer better suggestions for fixes. --> - Support `null_treatment`, `distinct`, and `filter` for window function in proto. - Support `null_treatment` for aggregate udf in proto. ## What changes are included in this PR? <!-- There is no need to duplicate the description in the issue here but it is sometimes worth providing a summary of the individual changes in this PR. --> - [x] Add `null_treatment`, `distinct`, `filter` fields to `WindowExprNode` message and handle them in `to/from_proto.rs`. - [x] Add `null_treatment` field to `AggregateUDFExprNode` message and handle them in `to/from_proto.rs`. - [ ] Docs update: I'm not sure where to add docs as declared in the issue description. ## Are these changes tested? <!-- We typically require tests for all PRs in order to: 1. Prevent the code from being accidentally broken by subsequent changes 2. Serve as another way to document the expected behavior of the code If tests are not included in your PR, please explain why (for example, are they covered by existing tests)? --> - Add tests to `roundtrip_window` for respectnulls, ignorenulls, distinct, filter. - Add tests to `roundtrip_aggregate_udf` for respectnulls, ignorenulls. ## Are there any user-facing changes? <!-- If there are user-facing changes then we may require documentation to be updated before approving the PR. --> <!-- If there are any breaking changes to public APIs, please add the `api change` label. --> N/A --------- Co-authored-by: Jeffrey Vo <jeffrey.vo.australia@gmail.com>
| Commit: | 980c948 | |
|---|---|---|
| Author: | Andrew Lamb | |
| Committer: | GitHub | |
Upgrade to arrow 56.1.0 (#17275) * Update to arrow/parquet 56.1.0 * Adjust for new parquet sizes, update for deprecated API * Thread through max_predicate_cache_size, add test
| Commit: | da89395 | |
|---|---|---|
| Author: | Jonathan Chen | |
| Committer: | GitHub | |
feat: Add `OR REPLACE` to creating external tables (#17580) * feat: Add `OR REPLACE` to creating external tables * regen * fmt * make more explicit + add tests * clipy fix --------- Co-authored-by: Dmitrii Blaginin <dmitrii@blaginin.me>
| Commit: | 7b16d6b | |
|---|---|---|
| Author: | Qi Zhu | |
| Committer: | GitHub | |
Support csv truncated rows in datafusion (#17465)
| Commit: | 6fd5685 | |
|---|---|---|
| Author: | 张林伟 | |
| Committer: | GitHub | |
Memory datasource protobuf support (#17290) * Add proto * fix proto * gen proto code * exec to proto * gen proto code * impl (de)serialization * Add test * Update submodules * gen proto * Set parquet-testing back to main --------- Co-authored-by: Andrew Lamb <andrew@nerdnetworks.org> Co-authored-by: Tim Saucer <timsaucer@gmail.com>
| Commit: | 2c9f42b | |
|---|---|---|
| Author: | Marko Milenković | |
| Committer: | GitHub | |
feat: Support SortMergeJoin proto serde (#17296) * Implement ser/de part of SortMergeJoin * add round trip tests for sort merge join * add filter test to roundtrip
| Commit: | 944d8c0 | |
|---|---|---|
| Author: | Peter L | |
| Committer: | GitHub | |
Support `distinct` and `ignore_nulls` in window expressions (#17235)
| Commit: | 2ae30af | |
|---|---|---|
| Author: | Peter L | |
| Committer: | GitHub | |
Support serializing `generate_series` in `datafusion-proto` (#17200) * Allow `generate_series` to be serialized via protobuf * Add breaking change to the upgrade guide
| Commit: | 60ac1cc | |
|---|---|---|
| Author: | Jonathan Chen | |
| Committer: | GitHub | |
fix: Remove `datafusion.execution.parquet.cache_metadata` config (#17062) * fix: Remove `datafusion.execution.parquet.cache_metadata` config * prettier * fix prettier? * fix * fix config * fix test behaviour * fix
| Commit: | fa1f8c1 | |
|---|---|---|
| Author: | Andrew Lamb | |
| Committer: | GitHub | |
Upgrade arrow/parquet to 56.0.0 (#16690)
| Commit: | c37dd5e | |
|---|---|---|
| Author: | Nuno Faria | |
| Committer: | GitHub | |
feat: Cache Parquet metadata in built in parquet reader (#16971) * feat: Cache Parquet metadata * Convert FileMetadata and FileMetadataCache to traits * Use as_any to respect MSRV * Use ObjectMeta as the key of FileMetadataCache
| Commit: | 5e0b2d0 | |
|---|---|---|
| Author: | Colin Marc | |
| Committer: | GitHub | |
fix(datafusion-proto): support serializing/deserilizing ArrowFormat tables (#16875) Fixes #16874
| Commit: | a6d4798 | |
|---|---|---|
| Author: | Nga Tran | |
| Committer: | GitHub | |
Fixes 3 bugs during serialization and deserialization of physical plans (#16858)
| Commit: | 8b03e5e | |
|---|---|---|
| Author: | Pepijn Van Eeckhoudt | |
| Committer: | GitHub | |
Use Tokio's task budget consistently, better APIs to support task cancellation (#16398) * Use Tokio's task budget consistently * Rework `ensure_coop` to base itself on evaluation and scheduling properties * Iterating on documentation * Improve robustness of cooperative yielding test cases * Reorganize tests by operator a bit better * Coop documentation * More coop documentation * Avoid Box in temporary CooperativeStream::poll_next implementation * Adapt interleave test cases for range generator * Add temporary `tokio_coop` feature to unblock merging * Extract magic number to constant * Fix documentation error * Push scheduling type down from DataSourceExec to DataSource * Use custom configuration instead of feature to avoid exposing internal cooperation variants * Use dedicated enum for yield results * Documentation improvements from review * More documentation * Change default coop strategy to 'tokio_fallback' * Documentation refinement * Re-enable interleave test cases * fix logical merge conflict --------- Co-authored-by: Andrew Lamb <andrew@nerdnetworks.org>
| Commit: | 11fc52d | |
|---|---|---|
| Author: | Tobias Schwarzinger | |
| Committer: | GitHub | |
Use dedicated NullEquality enum instead of null_equals_null boolean (#16419) * Use dedicated NullEquality enum instead of null_equals_null boolean * Fix wrong operator mapping in hash_join * Add an example to the documentation