From 2d4ff2d4d5539f4c92592ebe37747106de7d31bc Mon Sep 17 00:00:00 2001 From: morningman Date: Mon, 3 Aug 2026 09:04:43 +0800 Subject: [PATCH 1/5] [fix](be) Stop exporting the statically linked RocksDB symbols Exporting them makes this executable the definition every later-loaded library binds to, so a JNI library carrying its own RocksDB runs half on ours. The fluss scanner bundles frocksdbjni, whose librocksdbjni.so defines 2576 rocksdb symbols under names identical to ours but was built against the pre-C++11 libstdc++ string ABI: objects laid out by one copy and used by the other yield a garbage length, an std::bad_alloc that escapes the JNI frame, and an aborted BE. Reading any fluss primary-key table with a kv snapshot killed the process, reproducibly. Scoped to the archive rather than dropping ENABLE_EXPORTS, because what needs the exports is native UDFs (runtime/user_function_cache.cpp dlopens them) and those use the Doris UDF ABI, which has nothing to do with RocksDB. Crash stacks do not need it either -- they are symbolized from debug info, which is why they name even anonymous-namespace functions. 61 rocksdb symbols remain exported: inline and template members the compiler emitted into Doris's own objects, which no archive exclusion can reach. 29 of those still share a name with the JNI library, but none appear in its relocation table -- it never resolves them at load time, so they cannot be interposed. The library also duplicates zstd, lz4, snappy, bzip2 and zlib symbols; those are C ABIs, stable and layout-free, and are left alone. Verified: the fluss primary-key suite passes with BE alive (it aborted before); all three fluss suites green; an internal table survives write, BE restart and read, which is the tablet metadata RocksDB itself round-tripping. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_016VCPjzhwMQuP7nTgdvGVbM --- be/src/service/CMakeLists.txt | 18 ++++++++++++++++++ 1 file changed, 18 insertions(+) diff --git a/be/src/service/CMakeLists.txt b/be/src/service/CMakeLists.txt index b9bcef0b7b821a..64ed8651417748 100644 --- a/be/src/service/CMakeLists.txt +++ b/be/src/service/CMakeLists.txt @@ -49,6 +49,24 @@ if (${MAKE_TEST} STREQUAL "OFF" AND ${BUILD_BENCHMARK} STREQUAL "OFF") # This permits libraries loaded by dlopen to link to the symbols in the program. set_target_properties(doris_be PROPERTIES ENABLE_EXPORTS 1) + # ...but not the symbols of the RocksDB we link statically. Exporting those makes this + # executable the definition every later-loaded library binds to, and a JNI library that + # carries its own RocksDB then runs half on ours: the fluss scanner bundles frocksdbjni, + # whose librocksdbjni.so defines 2576 rocksdb symbols under names identical to ours but + # was built against the pre-C++11 libstdc++ string ABI. Objects laid out by one and used + # by the other yield a garbage length, an std::bad_alloc that escapes the JNI frame, and + # an aborted BE. Hiding this archive lets that library bind to its own copy. + # + # Scoped to the archive rather than dropping ENABLE_EXPORTS: what needs the exports is + # native UDFs (runtime/user_function_cache.cpp dlopens them), and those use the Doris UDF + # ABI, which has nothing to do with RocksDB. Crash stacks do not need it either -- they are + # symbolized from debug info, which is why they name even anonymous-namespace functions. + # + # The same library also duplicates zstd, lz4, snappy, bzip2 and zlib symbols. Those are C + # ABIs, stable across versions and layout-free, so they are left alone until something + # shows otherwise -- unlike RocksDB, whose C++ objects are what actually corrupt. + target_link_options(doris_be PRIVATE "-Wl,--exclude-libs,librocksdb.a") + target_link_libraries(doris_be ${DORIS_LINK_LIBS} ) From 4c4d7eba0ba2d59963fe7f4b12defdce3ccbb748 Mon Sep 17 00:00:00 2001 From: morningman Date: Mon, 3 Aug 2026 13:46:06 +0800 Subject: [PATCH 2/5] [fix](paimon) Claim the table handles this connector produces MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Connector.ownsHandle defaults to false, and this connector never overrode it. That was invisible while paimon was only ever a front-door catalog: the predicate exists so a GATEWAY connector can embed another as a sibling and route a foreign handle back to whoever made it, since the sibling's concrete handle type cannot be named across the plugin classloader split. The fluss connector reads a lake table by delegating to this one, so it asks that question about every handle it gets back — and got "not mine" about handles paimon had just produced. Every guard on the gateway side then falls through, and the first cast throws a ClassCastException naming the GATEWAY's handle type and two class loaders, with nothing to suggest the missing piece is a method here. Same one-liner the iceberg and hudi siblings behind the hms gateway already carry. No unit test could have caught this: a hand-written test double implements ownsHandle precisely because it has to, so the double is more capable than the real connector. It took an end-to-end read to surface. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_016VCPjzhwMQuP7nTgdvGVbM --- .../connector/paimon/PaimonConnector.java | 17 ++++ .../paimon/PaimonConnectorOwnsHandleTest.java | 78 +++++++++++++++++++ 2 files changed, 95 insertions(+) create mode 100644 fe/fe-connector/fe-connector-paimon/src/test/java/org/apache/doris/connector/paimon/PaimonConnectorOwnsHandleTest.java diff --git a/fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonConnector.java b/fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonConnector.java index 79f15eae474fd7..34990f28fc25f7 100644 --- a/fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonConnector.java +++ b/fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonConnector.java @@ -23,6 +23,7 @@ import org.apache.doris.connector.api.ConnectorPartitionInfo; import org.apache.doris.connector.api.ConnectorSession; import org.apache.doris.connector.api.ConnectorValidationContext; +import org.apache.doris.connector.api.handle.ConnectorTableHandle; import org.apache.doris.connector.api.scan.ConnectorScanPlanProvider; import org.apache.doris.connector.cache.ConnectorMetadataCache; import org.apache.doris.connector.metastore.HmsMetaStoreProperties; @@ -254,6 +255,22 @@ public ConnectorMetadata getMetadata(ConnectorSession session) { properties, context, schemaAtMemo, latestSnapshotCache, partitionViewCache); } + /** + * True for a handle this connector produced (a {@link PaimonTableHandle}). Tested against this connector's + * OWN in-loader type, so a gateway connector that embeds this one as a sibling can route a foreign paimon + * handle here without casting it across the plugin classloader split. Returns false for any other + * connector's handle, so the gateway keeps looking. + * + *

The default is {@code false}, which for a sibling means every one of the gateway's guards silently + * fails open and the first cast throws a ClassCastException instead — so this is required of any connector + * used as a sibling, not an optimization. Same implementation as the iceberg and hudi siblings behind the + * hms gateway. + */ + @Override + public boolean ownsHandle(ConnectorTableHandle handle) { + return handle instanceof PaimonTableHandle; + } + @Override public void invalidateTable(String dbName, String tableName) { // REFRESH TABLE (and, via the generic PluginDrivenExternalCatalog DDL hook, a Doris-issued diff --git a/fe/fe-connector/fe-connector-paimon/src/test/java/org/apache/doris/connector/paimon/PaimonConnectorOwnsHandleTest.java b/fe/fe-connector/fe-connector-paimon/src/test/java/org/apache/doris/connector/paimon/PaimonConnectorOwnsHandleTest.java new file mode 100644 index 00000000000000..a3f7955a8e49f3 --- /dev/null +++ b/fe/fe-connector/fe-connector-paimon/src/test/java/org/apache/doris/connector/paimon/PaimonConnectorOwnsHandleTest.java @@ -0,0 +1,78 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package org.apache.doris.connector.paimon; + +import org.apache.doris.connector.api.handle.ConnectorTableHandle; +import org.apache.doris.connector.spi.ConnectorContext; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +import java.util.Collections; + +/** + * Whether this connector claims the handles it produces, which is what lets a gateway connector + * embed it as a sibling: the gateway cannot name this module's handle type across the plugin classloader + * split, so it routes a handle by asking each sibling to test its own in-loader type. + * + *

Asserted rather than left to the SPI default, because the default is {@code false} and the failure it + * causes points away from here. A sibling that disowns its own handles makes every guard on the gateway + * side fail open, and the first cast throws a ClassCastException naming the GATEWAY's handle type and two + * class loaders — with nothing to suggest that the missing piece is a method this connector never + * overrode. That is not hypothetical: it is what happened the first time the fluss connector read a lake + * table through this one end to end. + */ +public class PaimonConnectorOwnsHandleTest { + + @Test + public void claimsItsOwnTableHandle() { + PaimonConnector connector = new PaimonConnector(Collections.emptyMap(), context()); + + Assertions.assertTrue(connector.ownsHandle(new PaimonTableHandle( + "db1", "t1", Collections.emptyList(), Collections.emptyList()))); + } + + @Test + public void disownsAnotherConnectorsHandle() { + // The gateway asks its siblings in turn, so answering yes to a foreign handle would route it to the + // wrong connector instead of leaving the gateway to keep looking. + PaimonConnector connector = new PaimonConnector(Collections.emptyMap(), context()); + + Assertions.assertFalse(connector.ownsHandle(new ForeignHandle())); + } + + /** Stands in for whatever another connector's handle happens to be; only its type matters here. */ + private static final class ForeignHandle implements ConnectorTableHandle { + private static final long serialVersionUID = 1L; + } + + /** The connector wraps whatever context it is given, so it cannot be null; nothing here reads it. */ + private static ConnectorContext context() { + return new ConnectorContext() { + @Override + public String getCatalogName() { + return "test_catalog"; + } + + @Override + public long getCatalogId() { + return 1L; + } + }; + } +} From 7f4d73a7434a1abe637158c4a2b36e87f9c7c816 Mon Sep 17 00:00:00 2001 From: morningman Date: Mon, 3 Aug 2026 13:46:19 +0800 Subject: [PATCH 3/5] [fix](be) Pick the table reader per scan range, not per scan node MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit FileScannerV2 built its table reader once, from the first range, and reused it for every range after that. One scan node can be given ranges of more than one table format: a fluss union read plans the table's lake half through the paimon connector and its log half itself, and both arrive as ranges of the same scan. Whichever range came first then decided the reader for all of them, and the other format's ranges were handed to a reader that does not understand them. That does not fail cleanly — it fails as whatever that reader makes of a foreign range. Here it was paimon's, reporting an unsupported file format for a fluss range that carries no paimon parameters at all. Which ranges share a scanner is up to the engine's assignment, so the same query succeeded or failed by how the ranges happened to be dealt out, and changing the projected columns could flip it either way. The reader now follows the range's table format. The expression contexts are deliberately not rebuilt: they are per-scanner and format-independent, and _init_expr_ctxes is not idempotent. Verified by disabling the rebuild and rerunning the suites: only the union read fails, with exactly the original error, and the four fluss suites that do not mix formats stay green. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_016VCPjzhwMQuP7nTgdvGVbM --- be/src/exec/scan/file_scanner_v2.cpp | 34 ++++++++++++++ be/src/exec/scan/file_scanner_v2.h | 7 +++ be/test/exec/scan/file_scanner_v2_test.cpp | 54 ++++++++++++++++++++++ 3 files changed, 95 insertions(+) diff --git a/be/src/exec/scan/file_scanner_v2.cpp b/be/src/exec/scan/file_scanner_v2.cpp index 4d513acbb9dd4c..d436a8385223a1 100644 --- a/be/src/exec/scan/file_scanner_v2.cpp +++ b/be/src/exec/scan/file_scanner_v2.cpp @@ -403,6 +403,7 @@ Status FileScannerV2::_open_impl(RuntimeState* state) { if (_first_scan_range) { RETURN_IF_ERROR(_create_table_reader_for_format(_current_range, &_table_reader)); DORIS_CHECK(_table_reader != nullptr); + _table_reader_format = table_format_name(_current_range); RETURN_IF_ERROR(_init_expr_ctxes()); RETURN_IF_ERROR(_init_table_reader(_current_range)); } @@ -502,6 +503,14 @@ Status FileScannerV2::_prepare_next_split(bool* eos) { DORIS_CHECK(_table_reader != nullptr); _current_range_path = _current_range.path; + bool reader_rebuilt = false; + RETURN_IF_ERROR(_rebuild_table_reader_if_format_changed(_current_range, &reader_rebuilt)); + if (reader_rebuilt) { + // Same init the first reader got. The expression contexts are NOT rebuilt: they are + // per-scanner and format-independent, and _init_expr_ctxes is not idempotent. + RETURN_IF_ERROR(_init_table_reader(_current_range)); + } + const auto format_type = get_range_format_type(*_params, _current_range); _init_adaptive_batch_size_state(format_type); if (_block_size_predictor != nullptr) { @@ -583,6 +592,31 @@ Status FileScannerV2::_init_table_reader(const TFileRangeDesc& range) { return Status::OK(); } +Status FileScannerV2::_rebuild_table_reader_if_format_changed(const TFileRangeDesc& range, + bool* rebuilt) { + // The reader is chosen by the range's table format, not the node's, because one node can be given + // both: a connector that reads a table as a lake plus the log written after it plans its lake half + // through a sibling connector and its log half itself, and both land here as ranges of the same + // scan. Built once from the first range and never revisited, the reader is then handed a range of + // the other format -- which does not fail cleanly. It fails as whatever that reader makes of a + // foreign range, e.g. paimon's reporting an unsupported file format for a range that carries no + // paimon parameters at all. And which ranges share a scanner is up to the engine's assignment, so + // the same query succeeds or fails by how the ranges happened to be dealt out. + // + // Split out from _prepare_next_split so the decision can be tested on its own: re-initializing the + // new reader needs scan-wide state that choosing it does not, so that step stays with the caller. + auto table_format = table_format_name(range); + if (table_format == _table_reader_format) { + *rebuilt = false; + return Status::OK(); + } + RETURN_IF_ERROR(_create_table_reader_for_format(range, &_table_reader)); + DORIS_CHECK(_table_reader != nullptr); + _table_reader_format = std::move(table_format); + *rebuilt = true; + return Status::OK(); +} + Status FileScannerV2::_create_table_reader_for_format( const TFileRangeDesc& range, std::unique_ptr* reader) const { DORIS_CHECK(reader != nullptr); diff --git a/be/src/exec/scan/file_scanner_v2.h b/be/src/exec/scan/file_scanner_v2.h index 92edbc1a1817f9..d12421be2fea9b 100644 --- a/be/src/exec/scan/file_scanner_v2.h +++ b/be/src/exec/scan/file_scanner_v2.h @@ -129,6 +129,9 @@ class FileScannerV2 final : public Scanner { Status _init_table_reader(const TFileRangeDesc& range); Status _create_table_reader_for_format(const TFileRangeDesc& range, std::unique_ptr* reader) const; + // Replaces _table_reader when {@code range} carries a different table format than the one it was + // built for, reporting whether it did. See the definition for why the reader follows the range. + Status _rebuild_table_reader_if_format_changed(const TFileRangeDesc& range, bool* rebuilt); Status _prepare_table_reader_split(const TFileRangeDesc& range, std::map partition_values); static bool _should_skip_not_found(const Status& status, bool ignore_not_found); @@ -181,6 +184,10 @@ class FileScannerV2 final : public Scanner { std::string _current_range_path; std::unique_ptr _table_reader; + // The table format _table_reader was built for. A scan node may mix table formats -- a fluss + // union read gives one node its lake half as paimon ranges and its log half as fluss ones -- and + // the reader is format-specific, so it is rebuilt whenever this stops matching the range. + std::string _table_reader_format; std::vector _projected_columns; // File formats without embedded schema, such as CSV, still need the FE slot descriptors in // file-column order. This mirrors old FileScanner::_file_slot_descs and is passed only to diff --git a/be/test/exec/scan/file_scanner_v2_test.cpp b/be/test/exec/scan/file_scanner_v2_test.cpp index 353c08043adef4..3506e5db27ae8a 100644 --- a/be/test/exec/scan/file_scanner_v2_test.cpp +++ b/be/test/exec/scan/file_scanner_v2_test.cpp @@ -475,6 +475,60 @@ TEST(FileScannerV2Test, JniCompatibilityShapesUseV2Scanner) { EXPECT_TRUE(FileScannerV2::is_supported(params, legacy_paimon_jni_range_without_reader_type())); } +// Scenario: one scan node is given ranges of two different table formats, which is what a connector +// reading a table as a lake plus the log written after it produces -- its lake half planned by a +// sibling connector, its own half by itself. The reader is format-specific, so it has to follow the +// RANGE. Built once from the first range, it is later handed a foreign one and fails as whatever that +// reader makes of it, not as a clean error; and since which ranges share a scanner is the engine's +// assignment, the same query then succeeds or fails by how the ranges happened to be dealt out. +TEST(FileScannerV2Test, TheTableReaderIsRebuiltWhenARangeChangesTableFormat) { + RuntimeState state {TQueryOptions(), TQueryGlobals()}; + RuntimeProfile profile("file_scanner_v2_reader_per_range"); + TFileScanRangeParams params; + params.__set_format_type(TFileFormatType::FORMAT_PARQUET); + + FileScannerV2 scanner(&state, &profile, nullptr); + scanner._params = ¶ms; + + const auto paimon_range = range_with_format("paimon", TFileFormatType::FORMAT_PARQUET); + const auto hive_range = range_with_format("hive", TFileFormatType::FORMAT_PARQUET); + + // Nothing has been built yet, so the first range always builds. + bool rebuilt = false; + ASSERT_TRUE(scanner._rebuild_table_reader_if_format_changed(paimon_range, &rebuilt).ok()); + EXPECT_TRUE(rebuilt); + EXPECT_EQ(scanner._table_reader_format, "paimon"); + const auto* first_reader = scanner._table_reader.get(); + ASSERT_NE(first_reader, nullptr); + + // A second range of the same format reuses it. Rebuilding here would be wasteful rather than + // wrong, but it would also throw away per-reader state the next split expects to still be there. + ASSERT_TRUE(scanner._rebuild_table_reader_if_format_changed(paimon_range, &rebuilt).ok()); + EXPECT_FALSE(rebuilt); + EXPECT_EQ(scanner._table_reader.get(), first_reader); + + // A range of another format must not be handed to the reader built for the first one. + ASSERT_TRUE(scanner._rebuild_table_reader_if_format_changed(hive_range, &rebuilt).ok()); + EXPECT_TRUE(rebuilt); + EXPECT_EQ(scanner._table_reader_format, "hive"); + EXPECT_NE(scanner._table_reader.get(), first_reader); + + // And back again, because the ranges of a mixed node arrive interleaved rather than grouped. + ASSERT_TRUE(scanner._rebuild_table_reader_if_format_changed(paimon_range, &rebuilt).ok()); + EXPECT_TRUE(rebuilt); + EXPECT_EQ(scanner._table_reader_format, "paimon"); + + // The formats really do get different readers -- otherwise every assertion above would hold + // just as well for a scanner that never rebuilt anything. + std::unique_ptr as_paimon; + std::unique_ptr as_hive; + ASSERT_TRUE(scanner._create_table_reader_for_format(paimon_range, &as_paimon).ok()); + ASSERT_TRUE(scanner._create_table_reader_for_format(hive_range, &as_hive).ok()); + const format::TableReader& paimon_reader = *as_paimon; + const format::TableReader& hive_reader = *as_hive; + EXPECT_STRNE(typeid(paimon_reader).name(), typeid(hive_reader).name()); +} + TEST(FileScannerV2Test, FailedTableReaderCloseCanBeRetriedThroughScanner) { RuntimeState state {TQueryOptions(), TQueryGlobals()}; RuntimeProfile profile("file_scanner_v2_close_retry"); From 648fdc61fd2dd85efb234ed77f4aa08e11627c75 Mon Sep 17 00:00:00 2001 From: morningman Date: Mon, 3 Aug 2026 17:35:58 +0800 Subject: [PATCH 4/5] [feat](paimon) Say which bucket a scan range came from A scan range this connector plans is opaque about its origin: the JNI arm carries a serialized split and nothing else, the native arm a file path and a byte interval. That is fine while the only reader is BE, which just reads what it is handed. It stops being fine once another connector plans splits here on behalf of its own table. The fluss connector does exactly that: a fluss table tiered into paimon keeps a bucket-identical layout, and reading it means pairing the lake data of bucket b with the log tail of bucket b that has not been tiered yet. Nothing on the range says b. Parsing it out of the data-file path would work only on the native arm and only by depending on this connector's directory layout. So carry it: paimon.bucket = DataSplit.bucket(), on the native and JNI arms alike -- which BE reader a split lands on is a session-level escape hatch the sibling does not control, and it must not change what the sibling can learn. FE-only; populateRangeParams does not forward it, so BE sees nothing new. Two ranges deliberately do NOT carry it. The collapsed COUNT(*) range stands for the splits of every bucket, so any single number on it would be a lie. A non-DataSplit system split has no bucket at all. A consumer that needs the binding must fail loud on an absent bucket rather than read it as "no state for this bucket" -- that reading turns a broken contract into duplicated rows. The fixture is two-bucket on purpose: with one bucket every range reads "0" and a hard-coded constant passes. Four mutations checked red -- constant bucket on the native arm, no bucket on the JNI arm, a bucket on the system split, a bucket on the count range. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_016VCPjzhwMQuP7nTgdvGVbM --- .../paimon/PaimonScanPlanProvider.java | 25 +- .../connector/paimon/PaimonScanRange.java | 21 ++ .../paimon/PaimonScanPlanProviderTest.java | 20 +- .../paimon/PaimonScanRangeBucketTest.java | 296 ++++++++++++++++++ 4 files changed, 344 insertions(+), 18 deletions(-) create mode 100644 fe/fe-connector/fe-connector-paimon/src/test/java/org/apache/doris/connector/paimon/PaimonScanRangeBucketTest.java diff --git a/fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonScanPlanProvider.java b/fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonScanPlanProvider.java index ef5e4650d88f9c..3f6213ec339cc6 100644 --- a/fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonScanPlanProvider.java +++ b/fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonScanPlanProvider.java @@ -761,7 +761,8 @@ private List planScanInternal( (optDeletionFiles.isPresent() && i < optDeletionFiles.get().size()) ? optDeletionFiles.get().get(i) : null; ranges.addAll(buildNativeRanges(file, deletionFile, defaultFileFormat, - partitionValues, vendedToken, effectiveSplitSize, weightDenominator)); + partitionValues, vendedToken, effectiveSplitSize, weightDenominator, + dataSplit.bucket())); } } else { // JNI reader path @@ -798,7 +799,8 @@ private List planScanInternal( */ PaimonScanRange buildNativeRange(RawFile file, DeletionFile deletionFile, String defaultFileFormat, Map partitionValues, - Map vendedToken, long start, long length, long weightDenominator) { + Map vendedToken, long start, long length, long weightDenominator, + int bucket) { String fileFormat = getFileFormatBySuffix(file.path()).orElse(defaultFileFormat); // FIX-A1: native sub-split FE weight = the sub-range byte length, + the deletion-vector length when // attached (legacy PaimonSplit(LocationPath,...).selfSplitWeight = length, setDeletionFile += DV). @@ -814,7 +816,8 @@ PaimonScanRange buildNativeRange(RawFile file, DeletionFile deletionFile, .partitionValues(partitionValues) .selfSplitWeight(selfSplitWeight) .targetSplitSize(weightDenominator) - .schemaId(file.schemaId()); + .schemaId(file.schemaId()) + .bucket(bucket); if (deletionFile != null) { builder.deletionFile( normalizeUri(deletionFile.path(), vendedToken), @@ -836,11 +839,12 @@ PaimonScanRange buildNativeRange(RawFile file, DeletionFile deletionFile, */ List buildNativeRanges(RawFile file, DeletionFile deletionFile, String defaultFileFormat, Map partitionValues, - Map vendedToken, long targetSplitSize, long weightDenominator) { + Map vendedToken, long targetSplitSize, long weightDenominator, + int bucket) { List result = new ArrayList<>(); for (long[] offset : computeFileSplitOffsets(file.length(), targetSplitSize)) { result.add(buildNativeRange(file, deletionFile, defaultFileFormat, - partitionValues, vendedToken, offset[0], offset[1], weightDenominator)); + partitionValues, vendedToken, offset[0], offset[1], weightDenominator, bucket)); } return result; } @@ -1372,13 +1376,18 @@ private PaimonScanRange buildJniScanRange(Split split, String defaultFileFormat, String fileFormat = isDataSplit ? dataSplitFileFormat((DataSplit) split, defaultFileFormat) : defaultFileFormat; - return new PaimonScanRange.Builder() + PaimonScanRange.Builder builder = new PaimonScanRange.Builder() .fileFormat(fileFormat) .paimonSplit(serializedSplit) .partitionValues(partitionValues) .selfSplitWeight(splitWeight) - .targetSplitSize(weightDenominator) - .build(); + .targetSplitSize(weightDenominator); + if (isDataSplit) { + // Same bucket property as the native arm: which reader BE ends up using must not change + // what a sibling connector can learn about the split (see PaimonScanRange's props). + builder.bucket(((DataSplit) split).bucket()); + } + return builder.build(); } /** diff --git a/fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonScanRange.java b/fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonScanRange.java index e6097aae4f5920..2a237762115741 100644 --- a/fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonScanRange.java +++ b/fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonScanRange.java @@ -91,6 +91,19 @@ private PaimonScanRange(Builder builder) { if (builder.rowCount != null) { props.put("paimon.row_count", String.valueOf(builder.rowCount)); } + // FE-ONLY (never reaches BE, see populateRangeParams): the paimon bucket this range's data + // belongs to, = DataSplit.bucket(). Read by a sibling connector that plans paimon splits on + // behalf of its own table and has to line them up with its own per-bucket state — today the + // fluss connector, whose lake half is planned here (it binds a fluss log tail to the lake + // splits of the SAME bucket; a lake table tiered from fluss has bucket-identical layout). + // Set on every DataSplit-backed range, native and JNI alike. NOT set on the collapsed + // COUNT(*) range (it stands for splits from many buckets, so any single number would be a + // lie) nor on a non-DataSplit system split (no bucket exists). Consumers must fail loud when + // it is absent on a range they expected to bind — silently treating that as "no state for + // this bucket" is a wrong-results bug, not a degradation. + if (builder.bucket != null) { + props.put("paimon.bucket", String.valueOf(builder.bucket)); + } // FIX-A3: emit the self-split-weight for every JNI split, incl. weight 0. Legacy // PaimonScanNode.setPaimonParams:274 sets it unconditionally on the JNI branch (never on // native); the old `selfSplitWeight > 0` gate was a buggy is-set proxy that dropped a genuine @@ -307,6 +320,9 @@ public static class Builder { // COUNT pushdown private Long rowCount; + // Bucket of the backing DataSplit; null for splits that have none (see the props comment). + private Integer bucket; + public Builder path(String path) { this.path = path; return this; @@ -369,6 +385,11 @@ public Builder rowCount(long rowCount) { return this; } + public Builder bucket(int bucket) { + this.bucket = bucket; + return this; + } + public PaimonScanRange build() { return new PaimonScanRange(this); } diff --git a/fe/fe-connector/fe-connector-paimon/src/test/java/org/apache/doris/connector/paimon/PaimonScanPlanProviderTest.java b/fe/fe-connector/fe-connector-paimon/src/test/java/org/apache/doris/connector/paimon/PaimonScanPlanProviderTest.java index 1c3dc7f6460ff3..3a1a6442a0bd42 100644 --- a/fe/fe-connector/fe-connector-paimon/src/test/java/org/apache/doris/connector/paimon/PaimonScanPlanProviderTest.java +++ b/fe/fe-connector/fe-connector-paimon/src/test/java/org/apache/doris/connector/paimon/PaimonScanPlanProviderTest.java @@ -351,7 +351,7 @@ public void nativeRangeNormalizesBothDataAndDeletionVectorPaths() { "oss://bkt/warehouse/db/t/index/dv-0.index", 8L, 16L, 4L); PaimonScanRange range = provider.buildNativeRange( - file, dv, "parquet", Collections.emptyMap(), Collections.emptyMap(), 0L, 100L, 64L * 1024 * 1024); + file, dv, "parquet", Collections.emptyMap(), Collections.emptyMap(), 0L, 100L, 64L * 1024 * 1024, 0); // WHY: BE's scheme-dispatched S3 file factory only opens canonical s3://. An un-normalized // oss:// DATA-file path fails the native ORC/Parquet read outright; an un-normalized oss:// DV @@ -376,7 +376,7 @@ public void nativeRangeWithoutDeletionVectorNormalizesOnlyDataPath() { PaimonScanRange range = provider.buildNativeRange( parquetRawFile("oss://bkt/a/part-0.parquet"), null, "parquet", - Collections.emptyMap(), Collections.emptyMap(), 0L, 100L, 64L * 1024 * 1024); + Collections.emptyMap(), Collections.emptyMap(), 0L, 100L, 64L * 1024 * 1024, 0); // WHY: a DV-less native split must still normalize its data-file path and must NOT emit a DV // descriptor. MUTATION: emitting a deletion_file for a null DV, or skipping data normalization -> red. @@ -396,7 +396,7 @@ public void nativeRangeWithoutContextPreservesRawPath() { PaimonScanRange range = provider.buildNativeRange( parquetRawFile("oss://bkt/a/part-0.parquet"), null, "parquet", - Collections.emptyMap(), Collections.emptyMap(), 0L, 100L, 64L * 1024 * 1024); + Collections.emptyMap(), Collections.emptyMap(), 0L, 100L, 64L * 1024 * 1024, 0); // MUTATION: NPE on null context, or fabricating a normalized path from nothing -> red. Assertions.assertEquals("oss://bkt/a/part-0.parquet", range.getPath().orElse(null)); @@ -421,7 +421,7 @@ public void buildNativeRangeThreadsVendedTokenToBothPaths() { "oss://bkt/warehouse/db/t/index/dv-0.index", 8L, 16L, 4L); PaimonScanRange range = provider.buildNativeRange( - file, dv, "parquet", Collections.emptyMap(), vendedToken, 0L, 100L, 64L * 1024 * 1024); + file, dv, "parquet", Collections.emptyMap(), vendedToken, 0L, 100L, 64L * 1024 * 1024, 0); // WHY: the engine seam normalizes against the VENDED map (the REST static map is empty). If the // connector dropped the token (reverting to the 1-arg seam) or substituted an empty map, a REST @@ -1733,7 +1733,7 @@ public void buildNativeRangesAttachesSameDeletionVectorToEverySubRange() { long target = Math.max(1L, file.length() / 3); // force the file to sub-split into >=2 ranges List ranges = provider.buildNativeRanges( - file, dv, "parquet", Collections.emptyMap(), Collections.emptyMap(), target, 64L * 1024 * 1024); + file, dv, "parquet", Collections.emptyMap(), Collections.emptyMap(), target, 64L * 1024 * 1024, 0); // WHY: the load-bearing correctness claim of FIX-NATIVE-SUBSPLIT — a paimon deletion vector is a // bitmap of GLOBAL file row positions, so EVERY sub-range of a DV-bearing file must carry the @@ -1761,7 +1761,7 @@ public void buildNativeRangesKeepsFileWholeWhenTargetNonPositive() { RawFile file = parquetRawFile("oss://bkt/a/part-0.parquet"); List ranges = provider.buildNativeRanges( - file, null, "parquet", Collections.emptyMap(), Collections.emptyMap(), 0L, 64L * 1024 * 1024); + file, null, "parquet", Collections.emptyMap(), Collections.emptyMap(), 0L, 64L * 1024 * 1024, 0); Assertions.assertEquals(1, ranges.size(), "a non-positive target (COUNT(*) pushdown) must keep the file as one whole-file range"); @@ -2572,14 +2572,14 @@ public void buildNativeRangeSetsProportionalWeightFromLengthAndDv() { DeletionFile dv = new DeletionFile("/data/dv-0.index", 8L, 16L, 4L); PaimonScanRange withDv = provider.buildNativeRange( - file, dv, "parquet", Collections.emptyMap(), Collections.emptyMap(), 0L, 64L, 64 * MB); + file, dv, "parquet", Collections.emptyMap(), Collections.emptyMap(), 0L, 64L, 64 * MB, 0); Assertions.assertEquals(64L + dv.length(), withDv.getSelfSplitWeight(), "native weight = sub-range length + the deletion-vector length"); Assertions.assertEquals(64 * MB, withDv.getTargetSplitSize(), "native range must carry the weight denominator"); PaimonScanRange noDv = provider.buildNativeRange( - file, null, "parquet", Collections.emptyMap(), Collections.emptyMap(), 0L, 70L, 64 * MB); + file, null, "parquet", Collections.emptyMap(), Collections.emptyMap(), 0L, 70L, 64 * MB, 0); Assertions.assertEquals(70L, noDv.getSelfSplitWeight(), "a DV-less native range weight is just the sub-range length"); } @@ -2597,7 +2597,7 @@ public void buildNativeRangesThreadsDenominatorDistinctFromFileSplitTarget() { List ranges = provider.buildNativeRanges( file, null, "parquet", Collections.emptyMap(), Collections.emptyMap(), - fileSplitTarget, denominator); + fileSplitTarget, denominator, 0); Assertions.assertEquals( PaimonScanPlanProvider.computeFileSplitOffsets(file.length(), fileSplitTarget).size(), @@ -2620,7 +2620,7 @@ public void buildNativeRangesCarriesDenominatorEvenWhenFileSplitSizeZero() { RawFile file = parquetRawFile("/data/part-0.parquet"); List ranges = provider.buildNativeRanges( - file, null, "parquet", Collections.emptyMap(), Collections.emptyMap(), 0L, 64 * MB); + file, null, "parquet", Collections.emptyMap(), Collections.emptyMap(), 0L, 64 * MB, 0); Assertions.assertEquals(1, ranges.size(), "a non-positive target keeps the file whole"); Assertions.assertEquals(64 * MB, ranges.get(0).getTargetSplitSize(), diff --git a/fe/fe-connector/fe-connector-paimon/src/test/java/org/apache/doris/connector/paimon/PaimonScanRangeBucketTest.java b/fe/fe-connector/fe-connector-paimon/src/test/java/org/apache/doris/connector/paimon/PaimonScanRangeBucketTest.java new file mode 100644 index 00000000000000..24de90b2ad4571 --- /dev/null +++ b/fe/fe-connector/fe-connector-paimon/src/test/java/org/apache/doris/connector/paimon/PaimonScanRangeBucketTest.java @@ -0,0 +1,296 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package org.apache.doris.connector.paimon; + +import org.apache.doris.connector.api.ConnectorSession; +import org.apache.doris.connector.api.handle.ConnectorColumnHandle; +import org.apache.doris.connector.api.scan.ConnectorScanRange; +import org.apache.doris.connector.api.scan.ConnectorScanRequest; + +import org.apache.paimon.catalog.Catalog; +import org.apache.paimon.catalog.FileSystemCatalog; +import org.apache.paimon.catalog.Identifier; +import org.apache.paimon.data.GenericRow; +import org.apache.paimon.fs.local.LocalFileIO; +import org.apache.paimon.schema.Schema; +import org.apache.paimon.table.Table; +import org.apache.paimon.table.sink.BatchTableCommit; +import org.apache.paimon.table.sink.BatchTableWrite; +import org.apache.paimon.table.sink.BatchWriteBuilder; +import org.apache.paimon.table.sink.CommitMessage; +import org.apache.paimon.table.source.DataSplit; +import org.apache.paimon.table.source.Split; +import org.apache.paimon.types.DataTypes; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +import java.nio.file.Path; +import java.util.Collections; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.TreeSet; + +/** + * Pins the {@code paimon.bucket} scan-range property (P4-2-a0). + * + *

WHY this property exists: a sibling connector can plan paimon splits on behalf of its OWN table + * and then has to line each split up with its own per-bucket state. Today that sibling is the fluss + * connector: a fluss table tiered into paimon keeps bucket-identical layout, and the fluss connector + * binds the un-tiered log tail of bucket b to the lake splits of bucket b. Without the + * bucket on the range there is nothing in a {@link PaimonScanRange} that says which bucket it came + * from (the JNI arm carries only an opaque serialized split; the native arm only a file path), so the + * sibling would have to fall back to whole-table binding — correct but wasteful — or parse the + * sibling's internal directory layout, which the JNI arm does not even expose. + * + *

The fixture is deliberately MULTI-bucket: with a single bucket every range would carry "0" and a + * hard-coded constant would pass. Each test therefore asserts the ranges reproduce the split-side + * bucket SET, which a constant cannot. + */ +public class PaimonScanRangeBucketTest { + + /** + * A two-bucket PK table with rows in BOTH buckets. PK {@code id} hashes into + * {@code bucket = hash(id) % 2}; ids 1..8 cover both buckets for paimon's hash function. + */ + private static Table createTwoBucketTable(Catalog catalog) throws Exception { + catalog.createDatabase("db", false); + Identifier id = Identifier.create("db", "t"); + catalog.createTable(id, Schema.newBuilder() + .column("id", DataTypes.INT()) + .column("val", DataTypes.BIGINT()) + .primaryKey("id") + .option("bucket", "2") + .build(), false); + Table table = catalog.getTable(id); + + BatchWriteBuilder wb = table.newBatchWriteBuilder(); + try (BatchTableWrite write = wb.newWrite()) { + for (int i = 1; i <= 8; i++) { + write.write(GenericRow.of(i, (long) i * 100)); + } + List messages = write.prepareCommit(); + try (BatchTableCommit commit = wb.newCommit()) { + commit.commit(messages); + } + } + return table; + } + + /** The buckets paimon's own read plan reports — the reference the ranges must reproduce. */ + private static Set planBuckets(Table table) throws Exception { + Set buckets = new TreeSet<>(); + for (Split s : table.newReadBuilder().newScan().plan().splits()) { + if (s instanceof DataSplit) { + buckets.add(((DataSplit) s).bucket()); + } + } + return buckets; + } + + /** The buckets the planned ranges claim, as ints. Fails the test if any range omits the property. */ + private static Set rangeBuckets(List ranges) { + Set buckets = new TreeSet<>(); + for (ConnectorScanRange r : ranges) { + String bucket = r.getProperties().get("paimon.bucket"); + Assertions.assertNotNull(bucket, + "every DataSplit-backed range must carry paimon.bucket; missing on " + r); + buckets.add(Integer.parseInt(bucket)); + } + return buckets; + } + + private static PaimonScanPlanProvider providerFor(Table table) { + RecordingPaimonCatalogOps ops = new RecordingPaimonCatalogOps(); + ops.table = table; + return new PaimonScanPlanProvider(Collections.emptyMap(), ops); + } + + private static PaimonTableHandle handleFor(String tableName) { + return new PaimonTableHandle("db", tableName, + Collections.emptyList(), Collections.emptyList()); + } + + @Test + public void nativeRangesCarryTheBucketOfTheSplitTheyCameFrom(@TempDir Path warehouse) + throws Exception { + try (Catalog catalog = new FileSystemCatalog(LocalFileIO.create(), + new org.apache.paimon.fs.Path(warehouse.toUri()))) { + Table table = createTwoBucketTable(catalog); + Set expected = planBuckets(table); + Assertions.assertTrue(expected.size() >= 2, + "fixture precondition: the table must really span >=2 buckets, got " + expected); + + List ranges = providerFor(table).planScan( + sessionWithProps(Collections.emptyMap()), + ConnectorScanRequest.builder(handleFor("t"), noColumns()).build()); + + Assertions.assertFalse(ranges.isEmpty(), "the fixture must plan at least one range"); + for (ConnectorScanRange r : ranges) { + Assertions.assertTrue(((PaimonScanRange) r).isNativeReadRange(), + "fixture precondition: this arm must exercise the NATIVE range builder"); + } + // WHY: a native range is one sub-range of one raw file of one DataSplit, so the bucket has + // to be threaded down from the DataSplit loop through buildNativeRanges/buildNativeRange. + // MUTATION: hard-coding 0 (or dropping .bucket() from the native builder) -> {0} instead of + // {0, 1} -> red; the multi-bucket fixture is what makes the constant detectable. + Assertions.assertEquals(expected, rangeBuckets(ranges), + "native ranges must reproduce exactly the buckets paimon's own plan reports"); + } + } + + @Test + public void jniRangesCarryTheBucketOfTheSplitTheyCameFrom(@TempDir Path warehouse) + throws Exception { + try (Catalog catalog = new FileSystemCatalog(LocalFileIO.create(), + new org.apache.paimon.fs.Path(warehouse.toUri()))) { + Table table = createTwoBucketTable(catalog); + Set expected = planBuckets(table); + Assertions.assertTrue(expected.size() >= 2, + "fixture precondition: the table must really span >=2 buckets, got " + expected); + + List ranges = providerFor(table).planScan( + sessionWithProps(Collections.singletonMap("force_jni_scanner", "true")), + ConnectorScanRequest.builder(handleFor("t"), noColumns()).build()); + + Assertions.assertFalse(ranges.isEmpty(), "the fixture must plan at least one range"); + for (ConnectorScanRange r : ranges) { + Assertions.assertTrue(r.getProperties().containsKey("paimon.split"), + "fixture precondition: this arm must exercise the JNI range builder"); + } + // WHY: which BE reader a split ends up on (native vs JNI, a session-level escape hatch the + // sibling does not control) must not change what the sibling can learn about the split. + // If only the native arm carried the bucket, turning on force_jni_scanner would silently + // break the sibling's binding. MUTATION: setting .bucket() only on the native arm -> the + // rangeBuckets assertNotNull fires -> red. + Assertions.assertEquals(expected, rangeBuckets(ranges), + "JNI ranges must reproduce exactly the buckets paimon's own plan reports"); + } + } + + @Test + public void collapsedCountRangeCarriesNoBucket(@TempDir Path warehouse) throws Exception { + try (Catalog catalog = new FileSystemCatalog(LocalFileIO.create(), + new org.apache.paimon.fs.Path(warehouse.toUri()))) { + Table table = createTwoBucketTable(catalog); + Assertions.assertTrue(planBuckets(table).size() >= 2, + "fixture precondition: >=2 buckets, so the collapse really does span buckets"); + + List ranges = providerFor(table).planScan( + sessionWithProps(Collections.emptyMap()), + ConnectorScanRequest.builder(handleFor("t"), noColumns()) + .countPushdown(true).build()); + + // WHY: the count collapse folds the splits of ALL buckets into ONE range carrying the summed + // total, so no single bucket number is true of it. Stamping the representative split's bucket + // would hand a sibling a range that claims to be bucket b while actually standing for every + // bucket — it would suppress/join against the wrong state. Absent is the honest answer, and + // the sibling is required to fail loud rather than guess (it never forwards count pushdown, + // so it must never see one of these). + // MUTATION: adding .bucket() to buildCountRange -> the count range carries one -> red. + int countRanges = 0; + for (ConnectorScanRange r : ranges) { + if (r.getProperties().containsKey("paimon.row_count")) { + ++countRanges; + Assertions.assertFalse(r.getProperties().containsKey("paimon.bucket"), + "the collapsed count range spans every bucket, so it must claim none"); + } + } + Assertions.assertEquals(1, countRanges, + "fixture precondition: count pushdown must produce exactly one collapsed range"); + } + } + + @Test + public void systemTableSplitCarriesNoBucket(@TempDir Path warehouse) throws Exception { + try (Catalog catalog = new FileSystemCatalog(LocalFileIO.create(), + new org.apache.paimon.fs.Path(warehouse.toUri()))) { + createTwoBucketTable(catalog); + Table snapshots = catalog.getTable(Identifier.create("db", "t$snapshots")); + + List ranges = providerFor(snapshots).planScan( + sessionWithProps(Collections.emptyMap()), + ConnectorScanRequest.builder(handleFor("t$snapshots"), noColumns()).build()); + + Assertions.assertFalse(ranges.isEmpty(), "a snapshots system table must plan >=1 range"); + // WHY: a system-table split is not a DataSplit and has no bucket at all — fabricating one + // (say 0) would be a lie a sibling could act on. MUTATION: setting .bucket() unconditionally + // in buildJniScanRange (dropping the isDataSplit gate) -> red. This also documents the shape + // the sibling must reject: it plans only data reads, so a bucket-less range reaching its + // wrapper means the contract broke. + for (ConnectorScanRange r : ranges) { + Assertions.assertFalse(r.getProperties().containsKey("paimon.bucket"), + "a non-DataSplit system split has no bucket, so it must not claim one"); + } + } + } + + private static List noColumns() { + return Collections.emptyList(); + } + + private static ConnectorSession sessionWithProps(Map sessionProps) { + return new ConnectorSession() { + @Override + public String getQueryId() { + return "q"; + } + + @Override + public String getUser() { + return "u"; + } + + @Override + public String getTimeZone() { + return "UTC"; + } + + @Override + public String getLocale() { + return "en_US"; + } + + @Override + public long getCatalogId() { + return 0; + } + + @Override + public String getCatalogName() { + return "c"; + } + + @Override + public T getProperty(String name, Class type) { + return null; + } + + @Override + public Map getCatalogProperties() { + return Collections.emptyMap(); + } + + @Override + public Map getSessionProperties() { + return sessionProps; + } + }; + } +} From c5ffeef95db258bd55b4617849619f9a55076557 Mon Sep 17 00:00:00 2001 From: morningman Date: Mon, 3 Aug 2026 19:35:25 +0800 Subject: [PATCH 5/5] [feat](connector) Let a connector name the columns its reader must read A connector whose BE-side reader merges or suppresses rows by key needs that key read whether or not the query selected it. Doris already keeps those columns for its own aggregate and merge-on-read unique-key tables -- preserveExtraStorageKeySlots, four lines above where the scan's slots are pruned -- for exactly that reason. A plugin connector had no way to say the same thing, and BE cannot read a column the plan never asked for. So ask it: getMustReadColumns, answered per scan, empty by default, so nothing changes for a connector that needs only what the query projects. The answer arrives during plan translation, after the scan node is initialized and before splits are planned, and widens the scan's tuple only -- the project above it already has its own output tuple, so the column is read and then dropped rather than returned. The question goes through the same memoized provider that will plan the splits, because the two have to come from one decision: a connector that answers "no extra columns" here and then plans a read that needs them leaves BE looking for a column that is not in the projection. A name that matches no slot fails the query and says which name, rather than being skipped -- skipping turns a disagreement about the table into silently wrong rows. Checked red by six mutations: dropping the branch, skipping unknown names, stopping after the first match, dropping the null answer guard, resolving a fresh provider to ask, and a non-empty SPI default. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_016VCPjzhwMQuP7nTgdvGVbM --- .../api/scan/ConnectorScanPlanProvider.java | 37 ++++ ...orScanPlanProviderMustReadColumnsTest.java | 79 ++++++++ .../datasource/scan/PluginDrivenScanNode.java | 23 +++ .../translator/PhysicalPlanTranslator.java | 39 ++++ ...uginDrivenScanNodeMustReadColumnsTest.java | 129 +++++++++++++ ...ysicalPlanTranslatorMustReadSlotsTest.java | 169 ++++++++++++++++++ 6 files changed, 476 insertions(+) create mode 100644 fe/fe-connector/fe-connector-api/src/test/java/org/apache/doris/connector/api/scan/ConnectorScanPlanProviderMustReadColumnsTest.java create mode 100644 fe/fe-core/src/test/java/org/apache/doris/datasource/scan/PluginDrivenScanNodeMustReadColumnsTest.java create mode 100644 fe/fe-core/src/test/java/org/apache/doris/nereids/glue/translator/PhysicalPlanTranslatorMustReadSlotsTest.java diff --git a/fe/fe-connector/fe-connector-api/src/main/java/org/apache/doris/connector/api/scan/ConnectorScanPlanProvider.java b/fe/fe-connector/fe-connector-api/src/main/java/org/apache/doris/connector/api/scan/ConnectorScanPlanProvider.java index acd44d5dfc9136..a13af026a759b0 100644 --- a/fe/fe-connector/fe-connector-api/src/main/java/org/apache/doris/connector/api/scan/ConnectorScanPlanProvider.java +++ b/fe/fe-connector/fe-connector-api/src/main/java/org/apache/doris/connector/api/scan/ConnectorScanPlanProvider.java @@ -30,6 +30,7 @@ import java.util.Map; import java.util.Optional; import java.util.OptionalLong; +import java.util.Set; /** * Plans the set of scan ranges (splits) needed to read a connector table. @@ -148,6 +149,42 @@ default TFileCompressType adjustFileCompressType(TFileCompressType inferred) { return inferred; } + /** + * The columns BE must READ for this scan even when the query references none of them, by Doris-side + * column name. The engine keeps their slots in the scan's tuple instead of pruning them away; the + * projection above the scan still removes them from the query's output, so the answer changes what is + * read, never what is returned. + * + *

This exists for a connector whose BE-side reader needs a column to produce CORRECT ROWS rather than + * to answer the query — a merge key, a suppression key, a row identity. Doris does the same thing for its + * own aggregate / merge-on-read unique-key tables ({@code PhysicalPlanTranslator.preserveExtraStorageKeySlots}): + * BE merges by key whether or not the user selected the key. Trino has no counterpart because its + * connectors own the page source and can add such columns privately; here the reader is BE, so the columns + * have to reach it through the plan.

+ * + *

Answer per SCAN, not per table: a connector that only sometimes needs the column (e.g. only when it + * decides to combine two sources) must return it only for those scans, and must reach the SAME decision + * when it later plans the splits — the engine asks this during plan translation, strictly before + * {@link #planScan}. Memoize that decision on the provider instance (the engine keeps one per scan node) + * rather than deciding twice: two independent decisions can disagree, and then BE is asked to read a + * column the tuple does not carry.

+ * + *

Every name returned must be a column of the scanned table, spelled as Doris knows it (the same + * identifier-mapped name {@link #classifyColumn} receives). A name that matches no slot in the scan's + * tuple fails the query loud: it means the connector and the engine disagree about the table, and reading + * on would silently produce whatever the connector's reader does without that column.

+ * + *

The default returns an empty set — every connector whose reader needs nothing beyond the projection + * is untouched, and its scans prune exactly as before.

+ * + * @param session the current session + * @param handle the table handle being scanned + * @return Doris-side names of the columns to read regardless of the projection (default: empty) + */ + default Set getMustReadColumns(ConnectorSession session, ConnectorTableHandle handle) { + return Collections.emptySet(); + } + /** * Plans the scan described by {@code request}, returning the ranges that cover the requested data. * diff --git a/fe/fe-connector/fe-connector-api/src/test/java/org/apache/doris/connector/api/scan/ConnectorScanPlanProviderMustReadColumnsTest.java b/fe/fe-connector/fe-connector-api/src/test/java/org/apache/doris/connector/api/scan/ConnectorScanPlanProviderMustReadColumnsTest.java new file mode 100644 index 00000000000000..f5e6cdb1d8d4dc --- /dev/null +++ b/fe/fe-connector/fe-connector-api/src/test/java/org/apache/doris/connector/api/scan/ConnectorScanPlanProviderMustReadColumnsTest.java @@ -0,0 +1,79 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package org.apache.doris.connector.api.scan; + +import org.apache.doris.connector.api.ConnectorSession; +import org.apache.doris.connector.api.handle.ConnectorTableHandle; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +import java.util.Collections; +import java.util.List; +import java.util.Set; + +/** + * Guards the additive {@code getMustReadColumns} SPI default on {@link ConnectorScanPlanProvider}. + * + *

WHY: the engine consults this on EVERY plugin-table scan that has a projection above it, and widens the + * scan's tuple by whatever comes back. The default must therefore be empty, or every connector that never + * asked for anything would start reading extra columns — and, worse, would fail the query loud when a name + * it never returned matches no slot. This is the zero-break guard for es/jdbc/paimon/iceberg/hive/maxcompute, + * none of which override it.

+ */ +public class ConnectorScanPlanProviderMustReadColumnsTest { + + /** Bare provider: only the abstract planScan implemented; everything else inherits SPI defaults. */ + private static final class BareProvider implements ConnectorScanPlanProvider { + @Override + public List planScan(ConnectorSession session, ConnectorScanRequest request) { + return Collections.emptyList(); + } + } + + /** A connector whose BE-side reader needs a merge key the query may not have selected. */ + private static final class KeyReadingProvider implements ConnectorScanPlanProvider { + @Override + public List planScan(ConnectorSession session, ConnectorScanRequest request) { + return Collections.emptyList(); + } + + @Override + public Set getMustReadColumns(ConnectorSession session, ConnectorTableHandle handle) { + return Collections.singleton("id"); + } + } + + @Test + public void defaultAsksForNoExtraColumns() { + ConnectorScanPlanProvider provider = new BareProvider(); + + // MUTATION: a default returning anything non-empty would widen every connector's scans and fail + // loud on the first name that matches no slot -> red here first. + Assertions.assertEquals(Collections.emptySet(), provider.getMustReadColumns(null, null), + "a connector that never opted in must ask for no extra columns"); + } + + @Test + public void connectorThatOptsInIsObeyed() { + ConnectorScanPlanProvider provider = new KeyReadingProvider(); + + Assertions.assertEquals(Collections.singleton("id"), provider.getMustReadColumns(null, null), + "the engine must read back exactly what the connector asked for"); + } +} diff --git a/fe/fe-core/src/main/java/org/apache/doris/datasource/scan/PluginDrivenScanNode.java b/fe/fe-core/src/main/java/org/apache/doris/datasource/scan/PluginDrivenScanNode.java index 910aa25e85d610..c5cc0cbaab7e43 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/datasource/scan/PluginDrivenScanNode.java +++ b/fe/fe-core/src/main/java/org/apache/doris/datasource/scan/PluginDrivenScanNode.java @@ -627,6 +627,29 @@ ConnectorColumnCategory classifyColumnByConnector(String columnName) { return onPluginClassLoader(scanProvider, () -> scanProvider.classifyColumn(columnName)); } + /** + * Asks the connector which columns BE must read for this scan even when the query references none of them + * ({@link ConnectorScanPlanProvider#getMustReadColumns}), so the translator can keep their slots instead of + * pruning them ({@code PhysicalPlanTranslator.preserveConnectorMustReadSlots}, the plugin-table counterpart + * of {@code preserveExtraStorageKeySlots} for aggregate / merge-on-read unique-key OLAP tables). + * + *

Asked through the SAME memoized provider the rest of the scan uses, so a connector that memoizes the + * decision on its provider instance answers this question and plans its splits from one decision — the + * whole point, since a column preserved here and a split plan that assumes otherwise disagree silently. + * A connector with no scan provider (no scan capability) needs nothing. Public + overridable because the + * caller is the translator, in another package, and so the preservation is unit-testable without a live + * connector (mirrors {@link #classifyColumnByConnector}, whose caller is this class).

+ */ + public Set mustReadColumnsFromConnector() { + ConnectorScanPlanProvider scanProvider = resolveScanProvider(); + if (scanProvider == null) { + return Collections.emptySet(); + } + Set columns = onPluginClassLoader(scanProvider, + () -> scanProvider.getMustReadColumns(connectorSession, currentHandle)); + return columns == null ? Collections.emptySet() : columns; + } + /** * Lets the owning connector adjust the compression type this node inferred from the split's file path * before it is shipped to BE, WITHOUT any source-specific code here: the base inference runs first, then diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/glue/translator/PhysicalPlanTranslator.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/glue/translator/PhysicalPlanTranslator.java index c09072174ce2d8..60c9a7bfb61914 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/glue/translator/PhysicalPlanTranslator.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/glue/translator/PhysicalPlanTranslator.java @@ -3041,6 +3041,8 @@ private void updateScanSlotsMaterialization(ScanNode scanNode, } if (scanNode instanceof OlapScanNode) { preserveExtraStorageKeySlots((OlapScanNode) scanNode, requiredWithVirtualColumns); + } else if (scanNode instanceof PluginDrivenScanNode) { + preserveConnectorMustReadSlots((PluginDrivenScanNode) scanNode, requiredWithVirtualColumns); } // Find the smallest column, for count(*) or other situation that slot is empty after prune SlotDescriptor smallest = getSmallestSlot(scanNode.getTupleDesc().getSlots()); @@ -3078,6 +3080,43 @@ private void preserveExtraStorageKeySlots(OlapScanNode scanNode, Set req } } + /** + * Keeps the slots of the columns a plugin connector must read for this scan even when the query + * references none of them — the plugin-table counterpart of {@link #preserveExtraStorageKeySlots}, and + * for the same reason: a reader that merges or suppresses rows by key needs the key whether or not the + * user selected it. The connector answers per scan + * ({@code ConnectorScanPlanProvider.getMustReadColumns}, empty by default), so every connector that needs + * nothing beyond the projection prunes exactly as before. + * + *

Only the scan's tuple is widened. The project above it was already given its own output tuple and + * project list a few lines up, so a column preserved here is read and then dropped — it never reaches the + * query's output.

+ * + *

A name that matches no slot fails the query loud rather than being skipped: it means the connector + * and the engine disagree about the table's columns, and the connector's reader would then be handed a + * scan missing a column it said it needs — silently wrong rows, not an error. Static + visible for testing + * so the ask-and-preserve step is pinned without a live connector.

+ */ + @VisibleForTesting + static void preserveConnectorMustReadSlots(PluginDrivenScanNode scanNode, Set requiredSlotIds) { + Set mustRead = scanNode.mustReadColumnsFromConnector(); + if (mustRead.isEmpty()) { + return; + } + Set missing = Sets.newLinkedHashSet(mustRead); + for (SlotDescriptor slot : scanNode.getTupleDesc().getSlots()) { + Column column = slot.getColumn(); + if (column != null && mustRead.contains(column.getName())) { + requiredSlotIds.add(slot.getId()); + missing.remove(column.getName()); + } + } + if (!missing.isEmpty()) { + throw new AnalysisException("connector requires column(s) " + missing + + " to be read, but the scan has no such column"); + } + } + private boolean shouldPreserveStorageKeySlots(OlapScanNode scanNode) { long selectedIndexId = scanNode.getSelectedIndexId() == -1 ? scanNode.getOlapTable().getBaseIndexId() diff --git a/fe/fe-core/src/test/java/org/apache/doris/datasource/scan/PluginDrivenScanNodeMustReadColumnsTest.java b/fe/fe-core/src/test/java/org/apache/doris/datasource/scan/PluginDrivenScanNodeMustReadColumnsTest.java new file mode 100644 index 00000000000000..a5ff8203b9fe73 --- /dev/null +++ b/fe/fe-core/src/test/java/org/apache/doris/datasource/scan/PluginDrivenScanNodeMustReadColumnsTest.java @@ -0,0 +1,129 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package org.apache.doris.datasource.scan; + +import org.apache.doris.common.jmockit.Deencapsulation; +import org.apache.doris.connector.api.Connector; +import org.apache.doris.connector.api.ConnectorSession; +import org.apache.doris.connector.api.handle.ConnectorTableHandle; +import org.apache.doris.connector.api.scan.ConnectorScanPlanProvider; + +import com.google.common.collect.ImmutableSet; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.mockito.Mockito; + +import java.util.Collections; +import java.util.Set; + +/** + * Guards {@link PluginDrivenScanNode#mustReadColumnsFromConnector()} — the seam the translator asks before it + * prunes a plugin table's scan slots ({@code PhysicalPlanTranslator.preserveConnectorMustReadSlots}), so a + * connector whose BE-side reader merges or suppresses rows by key gets that key read even when the query + * never mentions it. + * + *

WHY this matters (Rule 9): the answer is the whole contract between plan translation and split + * planning. Answering for the WRONG handle, or resolving a FRESH provider to ask, would let the connector + * decide "combine the two sources" at split time while the tuple was pruned as if it had said no — and BE + * would be told to suppress rows by a column that is not there. The memo assertion below is what pins + * "asked and planned through one provider instance".

+ * + *

Driven on a {@code CALLS_REAL_METHODS} node with only the connector/session/handle fields injected — + * the same technique as {@link PluginDrivenScanNodeScanProviderSelectionTest}.

+ */ +public class PluginDrivenScanNodeMustReadColumnsTest { + + private static PluginDrivenScanNode nodeWith(ConnectorScanPlanProvider provider, + ConnectorTableHandle handle, ConnectorSession session) { + PluginDrivenScanNode node = Mockito.mock(PluginDrivenScanNode.class, Mockito.CALLS_REAL_METHODS); + Connector connector = Mockito.mock(Connector.class); + Mockito.when(connector.getScanPlanProvider(handle)).thenReturn(provider); + Deencapsulation.setField(node, "connector", connector); + Deencapsulation.setField(node, "currentHandle", handle); + Deencapsulation.setField(node, "connectorSession", session); + return node; + } + + @Test + public void forwardsTheConnectorsAnswerForTheScannedHandle() { + ConnectorTableHandle handle = Mockito.mock(ConnectorTableHandle.class); + ConnectorSession session = Mockito.mock(ConnectorSession.class); + ConnectorScanPlanProvider provider = Mockito.mock(ConnectorScanPlanProvider.class); + Mockito.when(provider.getMustReadColumns(session, handle)).thenReturn(ImmutableSet.of("id", "part")); + PluginDrivenScanNode node = nodeWith(provider, handle, session); + + Set mustRead = node.mustReadColumnsFromConnector(); + + // WHY: the connector answers PER SCAN, from the handle this scan holds. MUTATION: passing a + // different handle (or the table's original one after pushdown refined it) makes the connector + // answer about another read -> the stub returns empty -> red. + Assertions.assertEquals(ImmutableSet.of("id", "part"), mustRead); + Mockito.verify(provider).getMustReadColumns(session, handle); + } + + @Test + public void connectorWithoutScanCapabilityNeedsNothing() { + ConnectorTableHandle handle = Mockito.mock(ConnectorTableHandle.class); + PluginDrivenScanNode node = nodeWith(null, handle, Mockito.mock(ConnectorSession.class)); + + // WHY: getScanPlanProvider() is null for a connector with no scan capability; every other resolver + // in this node degrades to its default rather than throwing. MUTATION: dropping the null check -> + // NPE during plan translation for such a catalog -> red. + Assertions.assertEquals(Collections.emptySet(), node.mustReadColumnsFromConnector()); + } + + @Test + public void nullAnswerIsReadAsNoExtraColumns() { + ConnectorTableHandle handle = Mockito.mock(ConnectorTableHandle.class); + ConnectorSession session = Mockito.mock(ConnectorSession.class); + ConnectorScanPlanProvider provider = Mockito.mock(ConnectorScanPlanProvider.class); + Mockito.when(provider.getMustReadColumns(session, handle)).thenReturn(null); + PluginDrivenScanNode node = nodeWith(provider, handle, session); + + // WHY: a third-party connector may return null where the SPI says "empty". Turning that into an + // NPE inside plan translation would blame the engine for a connector's slip. MUTATION: returning + // the raw answer -> NPE in the translator's isEmpty() -> red. + Assertions.assertEquals(Collections.emptySet(), node.mustReadColumnsFromConnector()); + } + + @Test + public void asksThroughTheSameProviderInstanceThatWillPlanTheSplits() { + ConnectorTableHandle handle = Mockito.mock(ConnectorTableHandle.class); + ConnectorSession session = Mockito.mock(ConnectorSession.class); + ConnectorScanPlanProvider provider = Mockito.mock(ConnectorScanPlanProvider.class); + Mockito.when(provider.getMustReadColumns(session, handle)).thenReturn(ImmutableSet.of("id")); + Connector connector = Mockito.mock(Connector.class); + Mockito.when(connector.getScanPlanProvider(handle)).thenReturn(provider); + PluginDrivenScanNode node = Mockito.mock(PluginDrivenScanNode.class, Mockito.CALLS_REAL_METHODS); + Deencapsulation.setField(node, "connector", connector); + Deencapsulation.setField(node, "currentHandle", handle); + Deencapsulation.setField(node, "connectorSession", session); + + node.mustReadColumnsFromConnector(); + Object providerForSplits = Deencapsulation.invoke(node, "resolveScanProvider"); + + // WHY: the connector is allowed to memoize "do I combine two sources?" on its provider instance, + // and MUST reach the same answer when it plans the splits later — the columns kept here and the + // splits planned there have to come from one decision. A fresh provider per question loses that + // memo and lets the two disagree. MUTATION: asking via connector.getScanPlanProvider(...) directly + // instead of the memoized resolveScanProvider() -> two instances + a second resolve -> red. + Assertions.assertSame(provider, providerForSplits, + "the must-read question must go through the same memoized provider as split planning"); + Mockito.verify(connector, Mockito.times(1)).getScanPlanProvider(handle); + } +} diff --git a/fe/fe-core/src/test/java/org/apache/doris/nereids/glue/translator/PhysicalPlanTranslatorMustReadSlotsTest.java b/fe/fe-core/src/test/java/org/apache/doris/nereids/glue/translator/PhysicalPlanTranslatorMustReadSlotsTest.java new file mode 100644 index 00000000000000..5f1caa9b39d6d2 --- /dev/null +++ b/fe/fe-core/src/test/java/org/apache/doris/nereids/glue/translator/PhysicalPlanTranslatorMustReadSlotsTest.java @@ -0,0 +1,169 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package org.apache.doris.nereids.glue.translator; + +import org.apache.doris.analysis.DescriptorTable; +import org.apache.doris.analysis.SlotDescriptor; +import org.apache.doris.analysis.SlotId; +import org.apache.doris.analysis.TupleDescriptor; +import org.apache.doris.catalog.Column; +import org.apache.doris.catalog.PrimitiveType; +import org.apache.doris.common.jmockit.Deencapsulation; +import org.apache.doris.datasource.scan.PluginDrivenScanNode; +import org.apache.doris.nereids.exceptions.AnalysisException; + +import com.google.common.collect.ImmutableList; +import com.google.common.collect.ImmutableSet; +import com.google.common.collect.Sets; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.mockito.Mockito; + +import java.util.Collections; +import java.util.Set; +import java.util.stream.Collectors; + +/** + * Guards {@code PhysicalPlanTranslator.preserveConnectorMustReadSlots} — the plugin-table branch next to + * {@code preserveExtraStorageKeySlots}, which keeps the slots a connector says its BE-side reader must read + * even when the query references none of them. + * + *

WHY this matters (Rule 9): the scan's tuple is where the columns BE reads are decided + * ({@code FileQueryScanNode.updateRequiredSlots} rebuilds the required-slot list from exactly these slots + * after planning). A reader that suppresses or merges rows by key and does not get the key back reads the + * key-less rows and emits duplicates — no error anywhere. Doris does the same preservation for its own + * aggregate / merge-on-read unique-key tables; this is that mechanism for plugin connectors.

+ * + *

These tests drive the extracted static entry point directly with a real {@link TupleDescriptor} and a + * {@code CALLS_REAL_METHODS} node whose connector answer is stubbed — building a translator over a live + * plugin catalog needs a harness this module does not have. What they do NOT cover, because it is decided + * before this branch runs and by generic code: the project above the scan gets its own output tuple, so a + * column preserved here is read and then dropped rather than returned. The fluss suites' row baselines are + * the end-to-end guard for that.

+ */ +public class PhysicalPlanTranslatorMustReadSlotsTest { + + private static final DescriptorTable DESC_TABLE = new DescriptorTable(); + + /** A scan tuple holding one slot per named column, in order. */ + private static TupleDescriptor tupleOf(String... columnNames) { + TupleDescriptor tuple = DESC_TABLE.createTupleDescriptor(); + for (String name : columnNames) { + SlotDescriptor slot = DESC_TABLE.addSlotDescriptor(tuple); + slot.setColumn(new Column(name, PrimitiveType.INT)); + } + return tuple; + } + + private static PluginDrivenScanNode nodeAnswering(TupleDescriptor tuple, Set mustRead) { + PluginDrivenScanNode node = Mockito.mock(PluginDrivenScanNode.class, Mockito.CALLS_REAL_METHODS); + Mockito.doReturn(tuple).when(node).getTupleDesc(); + Mockito.doReturn(mustRead).when(node).mustReadColumnsFromConnector(); + return node; + } + + private static SlotId slotIdOf(TupleDescriptor tuple, String columnName) { + for (SlotDescriptor slot : tuple.getSlots()) { + if (slot.getColumn().getName().equals(columnName)) { + return slot.getId(); + } + } + throw new IllegalStateException("no slot for " + columnName); + } + + @Test + public void connectorNamedColumnsSurvivePruning() { + TupleDescriptor tuple = tupleOf("id", "name", "amount"); + PluginDrivenScanNode node = nodeAnswering(tuple, ImmutableSet.of("id")); + // "select name": only that slot is required by the project above the scan. + Set required = Sets.newHashSet(slotIdOf(tuple, "name")); + + PhysicalPlanTranslator.preserveConnectorMustReadSlots(node, required); + + // WHY: 'id' is what the connector's reader needs to suppress rows; without it in the required set + // the removeIf below this call drops it from the tuple and BE reads key-less rows. MUTATION: + // dropping the branch (or the add) -> 'id' absent -> red. + Assertions.assertEquals(ImmutableSet.of(slotIdOf(tuple, "name"), slotIdOf(tuple, "id")), required); + } + + @Test + public void connectorThatNeedsNothingChangesNothing() { + TupleDescriptor tuple = tupleOf("id", "name", "amount"); + PluginDrivenScanNode node = nodeAnswering(tuple, Collections.emptySet()); + Set required = Sets.newHashSet(slotIdOf(tuple, "name")); + + PhysicalPlanTranslator.preserveConnectorMustReadSlots(node, required); + + // WHY: this is the gate that keeps the branch inert for every connector that never opted in — and + // for an opted-in connector on a scan it decided NOT to combine (a fluss table read from fluss + // alone), which is exactly the "only when it is really needed" requirement. MUTATION: preserving + // unconditionally (e.g. the whole primary key regardless of the decision) -> an extra slot -> red. + Assertions.assertEquals(Collections.singleton(slotIdOf(tuple, "name")), required); + } + + @Test + public void everyNamedColumnIsPreservedNotJustTheFirst() { + TupleDescriptor tuple = tupleOf("k1", "k2", "payload"); + PluginDrivenScanNode node = nodeAnswering(tuple, ImmutableSet.of("k1", "k2")); + Set required = Sets.newHashSet(slotIdOf(tuple, "payload")); + + PhysicalPlanTranslator.preserveConnectorMustReadSlots(node, required); + + // WHY: composite keys are the normal case for the readers this exists for; keeping only one column + // of a two-column key compares the wrong thing. MUTATION: `break` after the first match -> red. + Assertions.assertEquals( + ImmutableSet.of(slotIdOf(tuple, "payload"), slotIdOf(tuple, "k1"), slotIdOf(tuple, "k2")), + required); + } + + @Test + public void preservedSlotSurvivesThePruneItself() { + TupleDescriptor tuple = tupleOf("id", "name", "amount"); + PluginDrivenScanNode node = nodeAnswering(tuple, ImmutableSet.of("id")); + Set required = Sets.newHashSet(slotIdOf(tuple, "name")); + + // The real prune step, driven end to end: it is what decides which slots the scan reads, and the + // branch under test sits inside it. + Deencapsulation.invoke(new PhysicalPlanTranslator(), "updateScanSlotsMaterialization", + node, required, Sets.newHashSet(), new PlanTranslatorContext()); + + // WHY: this is the only assertion that also pins the DISPATCH — that a plugin-driven scan reaches + // the branch at all. MUTATION: deleting the `else if (scanNode instanceof PluginDrivenScanNode)` + // arm -> 'id' pruned away -> red. MUTATION: preserving AFTER the removeIf -> also red. + Assertions.assertEquals(ImmutableList.of("id", "name"), + tuple.getSlots().stream().map(s -> s.getColumn().getName()).collect(Collectors.toList()), + "the connector's column must be read; the unreferenced one must still be pruned"); + } + + @Test + public void columnTheScanDoesNotHaveFailsLoud() { + TupleDescriptor tuple = tupleOf("id", "name"); + PluginDrivenScanNode node = nodeAnswering(tuple, ImmutableSet.of("id", "ghost")); + Set required = Sets.newHashSet(slotIdOf(tuple, "name")); + + AnalysisException thrown = Assertions.assertThrows(AnalysisException.class, + () -> PhysicalPlanTranslator.preserveConnectorMustReadSlots(node, required)); + + // WHY: a name matching no slot means the connector and the engine disagree about the table. Reading + // on would hand the reader a scan without a column it said it needs — wrong rows, silently. The + // message must name the column, because that is the only clue to which side is stale. MUTATION: + // skipping unknown names instead of throwing -> no exception -> red. + Assertions.assertTrue(thrown.getMessage().contains("ghost"), + "the failure must name the column the scan does not have: " + thrown.getMessage()); + } +}