morningman commented on code in PR #66399: URL: https://github.com/apache/doris/pull/66399#discussion_r4120588280
########## fe/fe-connector/fe-connector-fluss/src/main/java/org/apache/doris/connector/fluss/FlussTableHandle.java: ########## @@ -0,0 +1,356 @@ +// 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.fluss; + +import org.apache.doris.connector.spi.handle.ConnectorTableHandle; + +import org.apache.fluss.metadata.DataLakeFormat; +import org.apache.fluss.metadata.TableInfo; +import org.apache.fluss.metadata.TablePath; +import org.apache.fluss.types.DataType; +import org.apache.fluss.types.RowType; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Objects; + +/** + * A fluss table, as everything downstream of {@code getTableHandle} needs to know about it. + * + * <p>The fields are the ones split planning reads, snapshotted from the {@link TableInfo} the handle + * was built from: whether the table has a primary key (which of the two scanner families reads it), + * the bucket count and keys (how many splits and how they line up with the lake's), the partition keys, + * and whether the table is tiered into a lake at all. They are copied rather than looked up again + * because planning must see one coherent view of the table, not a field from before an ALTER and a + * field from after it. + * + * <p>{@link #getProperties()} is the fluss table's own property map, and it carries more than the user + * wrote: for a datalake-enabled table the fluss coordinator merges the cluster's lake-catalog + * connection settings into it under {@code table.datalake.<format>.*}, which is where the lake side of + * a union read gets its catalog configuration. Nothing else in Doris supplies it — the Doris catalog is + * configured with fluss bootstrap servers only. + * + * <p>The column schema is deliberately NOT here — with one exception. It is the one part that a + * statement re-reads through {@link FlussStatementScope}, so the handle stays a small, serializable + * identity object; {@link #getKeyColumnTypes()} carries the types of the primary-key and partition-key + * columns alone, because planning has to reason about those by type (see its javadoc) and re-reading + * the whole schema for them would both cost a round trip and risk answering from a different schema + * version than the rest of this handle describes. + * + * <p>{@link #getReadMode()} is the one field that is NOT a fact about the table: it says WHICH SEGMENT + * of the table this handle stands for. A handle reached through {@code tbl} covers the whole table; one + * reached through {@code tbl$log} covers only the part that is still in fluss's log — the complement of + * what {@code tbl$lake} serves. The two segments partition the table, so the field belongs to identity + * (see {@link #equals}): two handles that name the same table but different segments describe different + * row sets. + */ +public class FlussTableHandle implements ConnectorTableHandle { + + private static final long serialVersionUID = 1L; + + /** + * Which segment of the table a handle stands for. + * + * <p>The names describe the data, not the freshness: {@link #LOG_ONLY} is "what is in the log", which + * stays true however far behind tiering has fallen, whereas a name like "realtime" would be a lie the + * moment tiering stops. Path selection — reading the current Fluss state alone versus combining the + * lake with its log — is a different question and is NOT expressed here; it is the catalog's (or the + * session's) union-read mode. + */ + public enum ReadMode { + /** The whole table. */ + DEFAULT, + /** Only the part not yet tiered into the lake: the log past the lake snapshot's offsets. */ + LOG_ONLY + } + + private final String databaseName; + private final String tableName; + private final String lakeDatabaseName; + private final String lakeTableName; + private final long tableId; + private final int schemaId; + private final boolean hasPrimaryKey; + private final List<String> primaryKeys; + private final List<String> bucketKeys; + private final int bucketCount; + private final List<String> partitionKeys; + private final boolean dataLakeEnabled; + /** The lake format's fluss name ({@code "paimon"}), or {@code null} when the table declares none. */ + private final String dataLakeFormat; + private final Map<String, String> properties; + private final Map<String, DataType> keyColumnTypes; + private final ReadMode readMode; + + public FlussTableHandle(String databaseName, String tableName, long tableId, int schemaId, + boolean hasPrimaryKey, List<String> primaryKeys, List<String> bucketKeys, int bucketCount, + List<String> partitionKeys, boolean dataLakeEnabled, String dataLakeFormat, + Map<String, String> properties, Map<String, DataType> keyColumnTypes) { + this(databaseName, tableName, tableId, schemaId, hasPrimaryKey, primaryKeys, bucketKeys, + bucketCount, partitionKeys, dataLakeEnabled, dataLakeFormat, properties, + keyColumnTypes, databaseName, tableName); + } + + public FlussTableHandle(String databaseName, String tableName, long tableId, int schemaId, + boolean hasPrimaryKey, List<String> primaryKeys, List<String> bucketKeys, int bucketCount, + List<String> partitionKeys, boolean dataLakeEnabled, String dataLakeFormat, + Map<String, String> properties, Map<String, DataType> keyColumnTypes, + String lakeDatabaseName, String lakeTableName) { + this.readMode = ReadMode.DEFAULT; + this.databaseName = Objects.requireNonNull(databaseName, "databaseName"); + this.tableName = Objects.requireNonNull(tableName, "tableName"); + this.lakeDatabaseName = Objects.requireNonNull(lakeDatabaseName, "lakeDatabaseName"); + this.lakeTableName = Objects.requireNonNull(lakeTableName, "lakeTableName"); + this.tableId = tableId; + this.schemaId = schemaId; + this.hasPrimaryKey = hasPrimaryKey; + this.primaryKeys = copyOf(primaryKeys); + this.bucketKeys = copyOf(bucketKeys); + this.bucketCount = bucketCount; + this.partitionKeys = copyOf(partitionKeys); + this.dataLakeEnabled = dataLakeEnabled; + this.dataLakeFormat = dataLakeFormat; + this.properties = properties == null + ? Collections.emptyMap() + : Collections.unmodifiableMap(new LinkedHashMap<>(properties)); + this.keyColumnTypes = keyColumnTypes == null + ? Collections.emptyMap() + : Collections.unmodifiableMap(new LinkedHashMap<>(keyColumnTypes)); + } + + /** + * Re-reads {@code source} at a different segment. The table facts are aliased rather than copied — + * every one of them is already unmodifiable — so the two handles cannot drift apart into describing + * the same table two different ways. + */ + private FlussTableHandle(FlussTableHandle source, ReadMode readMode) { + this.readMode = readMode; + this.databaseName = source.databaseName; + this.tableName = source.tableName; + this.lakeDatabaseName = source.lakeDatabaseName; + this.lakeTableName = source.lakeTableName; + this.tableId = source.tableId; + this.schemaId = source.schemaId; + this.hasPrimaryKey = source.hasPrimaryKey; + this.primaryKeys = source.primaryKeys; + this.bucketKeys = source.bucketKeys; + this.bucketCount = source.bucketCount; + this.partitionKeys = source.partitionKeys; + this.dataLakeEnabled = source.dataLakeEnabled; + this.dataLakeFormat = source.dataLakeFormat; + this.properties = source.properties; + this.keyColumnTypes = source.keyColumnTypes; + } + + /** Snapshots {@code tableInfo} into a handle. */ + public static FlussTableHandle of(TableInfo tableInfo) { + TablePath path = tableInfo.getTablePath(); + TablePath lakePath = tableInfo.getLakeTablePath(); + DataLakeFormat lakeFormat = tableInfo.getTableConfig().getDataLakeFormat().orElse(null); + return new FlussTableHandle( + path.getDatabaseName(), + path.getTableName(), + tableInfo.getTableId(), + tableInfo.getSchemaId(), + tableInfo.hasPrimaryKey(), + tableInfo.getPrimaryKeys(), + tableInfo.getBucketKeys(), + tableInfo.getNumBuckets(), + tableInfo.getPartitionKeys(), + tableInfo.getTableConfig().isDataLakeEnabled(), + lakeFormat == null ? null : lakeFormat.toString(), + tableInfo.getProperties().toMap(), + keyColumnTypes(tableInfo), + lakePath.getDatabaseName(), + lakePath.getTableName()); + } + + /** + * The types of the columns that are part of the primary key or of the partition key, taken from the + * same {@link TableInfo} as every other field here. + */ + private static Map<String, DataType> keyColumnTypes(TableInfo tableInfo) { + RowType rowType = tableInfo.getRowType(); + Map<String, DataType> types = new LinkedHashMap<>(); + List<String> keyColumns = new ArrayList<>(tableInfo.getPrimaryKeys()); + keyColumns.addAll(tableInfo.getPartitionKeys()); + for (String column : keyColumns) { + int index = rowType.getFieldIndex(column); + if (index >= 0) { + types.put(column, rowType.getTypeAt(index)); + } + } + return types; + } + + public TablePath toTablePath() { + return TablePath.of(databaseName, tableName); + } + + public TablePath toLakeTablePath() { + return TablePath.of(lakeDatabaseName, lakeTableName); + } + + public String getDatabaseName() { + return databaseName; + } + + public String getTableName() { + return tableName; + } + + public String getLakeDatabaseName() { + return lakeDatabaseName; + } + + public String getLakeTableName() { + return lakeTableName; + } + + public long getTableId() { + return tableId; + } + + public int getSchemaId() { + return schemaId; + } + + public boolean hasPrimaryKey() { + return hasPrimaryKey; + } + + public List<String> getPrimaryKeys() { + return primaryKeys; + } + + public List<String> getBucketKeys() { + return bucketKeys; + } + + public int getBucketCount() { + return bucketCount; + } + + public List<String> getPartitionKeys() { + return partitionKeys; + } + + public boolean isPartitioned() { + return !partitionKeys.isEmpty(); + } + + public boolean isDataLakeEnabled() { + return dataLakeEnabled; + } + + public String getDataLakeFormat() { + return dataLakeFormat; + } + + public Map<String, String> getProperties() { + return properties; + } + + public ReadMode getReadMode() { + return readMode; + } + + /** Whether this handle covers only the log past the lake snapshot, i.e. it was reached as {@code $log}. */ + public boolean isLogOnly() { + return readMode == ReadMode.LOG_ONLY; + } + + /** The same table, read as its log tail alone. */ + public FlussTableHandle asLogOnly() { + return readMode == ReadMode.LOG_ONLY ? this : new FlussTableHandle(this, ReadMode.LOG_ONLY); + } + + /** + * The fluss types of the primary-key and partition-key columns, by column name. + * + * <p>Split planning needs these two, and only these two, by type. A primary-key table read as the + * union of its lake and its log tail identifies rows across the two halves BY KEY, so a key column + * whose values do not compare exactly the same way on both sides (a float, a timestamp Doris rounds) + * cannot be read that way at all. A partition column is matched the same way one level up: a lake + * split is bound to a fluss partition by comparing the rendered partition values, which is only + * sound for a type both sides render identically. + * + * <p>Everything else about the schema stays out of the handle — the point is not "some of the + * schema", it is the columns whose type decides how the table can be PLANNED. + */ + public Map<String, DataType> getKeyColumnTypes() { + return keyColumnTypes; + } + + /** The primary-key columns that are not partition columns — what a bucket's rows are keyed by. */ + public List<String> getPhysicalPrimaryKeys() { Review Comment: Thanks for tracing the FE/BE failure path. It is accurate for an empty physical key, but `PRIMARY KEY (dt) PARTITIONED BY (dt)` is not a legal Fluss 1.0.0 table. `TableDescriptor.normalizeDistribution()` derives the default bucket key from the primary key minus partition columns, and `defaultBucketKeyOfPrimaryKeyTable()` throws when that list is empty. An explicit bucket key cannot make this descriptor valid either: a partition column cannot be a bucket key, while another column is not part of this table's primary key. See [Fluss 1.0.0 TableDescriptor](https://github.com/apache/fluss/blob/v1.0.0/fluss-common/src/main/java/org/apache/fluss/metadata/TableDescriptor.java#L348-L412). Doris builds the handle from Fluss `TableInfo`, so a valid table cannot reach this empty-key union path. A nonempty-tail regression for the proposed shape would have to fabricate metadata that Fluss refuses to create. We will leave this path unchanged. ########## fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonConnector.java: ########## @@ -389,7 +389,13 @@ public Set<ConnectorCapability> getCapabilities() { // connector-wide: it holds for every paimon DATA table. The narrower question of which // SYSTEM table can honor the clause is answered per table by // PaimonScanPlanProvider.supportsSystemTableOptions. - ConnectorCapability.SUPPORTS_SCAN_PARAM_OPTIONS); + ConnectorCapability.SUPPORTS_SCAN_PARAM_OPTIONS, + // SUPPORTS_NESTED_COLUMN_PRUNE: the paimon JNI scanner mirrors a pruned nested type onto + // paimon's own types and pushes it down (ReadBuilder.withReadType), and the native + // parquet/orc split path resolves the access paths by name. NOT + // SUPPORTS_FIELD_ID_ACCESS_PATH: paimon carries no field id on the Doris column tree, so + // rewriting the paths to ids would make every segment "-1". + ConnectorCapability.SUPPORTS_NESTED_COLUMN_PRUNE); Review Comment: Thanks. The decoder mismatch is real if a new FE executes against an old Paimon JNI scanner. Doris's [upgrade procedure](https://doris.apache.org/docs/4.x/admin-manual/cluster-management/upgrade/) requires all BEs to be upgraded before any FE. For this change, that BE rollout must include the matching `be/plugins/jni/paimon` artifact on every eligible BE: [the build puts the scanner there](https://github.com/apache/doris/blob/221c7d93d207a98c27f48b39f31e05b48b9ba9fc/build.sh#L1511-L1529), outside `bin` and `lib`. The guide's `bin`/`lib` copy and `SHOW BACKENDS Version` alone do not verify the scanner JAR, and `be_exec_version` does not identify its implementation. We will treat the matching Paimon plugin as part of the BE upgrade and upgrade FE only after it is deployed to every eligible BE. Under that required sequence, a new FE cannot run with the old scanner. We are leaving the capability and wire type unchanged in this PR. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
