github-actions[bot] commented on code in PR #66399:
URL: https://github.com/apache/doris/pull/66399#discussion_r4120289515


##########
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:
   [P1] Handle a primary key made entirely of partition columns before emitting 
suppression. `getPhysicalPrimaryKeys()` returns an empty list for a legal 
`PRIMARY KEY (dt) PARTITIONED BY (dt)` table, but the union planner still emits 
`LAKE_SUPPRESS` and `PK_TAIL` ranges when the lake has a tail. 
`getScanNodeProperties()` serializes an empty `fluss.union.pk_names`; BE then 
leaves `_key_columns` empty and `_prepare_suppression()` fails every lake 
split. Please either include the constant partition key in the suppression 
identity or reject/degrade this union shape before ranges are emitted, and add 
a regression covering a nonempty tail.



-- 
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]

Reply via email to