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/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} ) 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"); 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-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/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/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; + } + }; + } +} 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; + } + }; + } +} 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()); + } +}