diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/committer/FlussTableLakeSnapshotCommitter.java b/fluss-client/src/main/java/org/apache/fluss/client/tiering/FlussTableLakeSnapshotCommitter.java similarity index 89% rename from fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/committer/FlussTableLakeSnapshotCommitter.java rename to fluss-client/src/main/java/org/apache/fluss/client/tiering/FlussTableLakeSnapshotCommitter.java index 76474a5b265..0fb01130404 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/committer/FlussTableLakeSnapshotCommitter.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/tiering/FlussTableLakeSnapshotCommitter.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.fluss.flink.tiering.committer; +package org.apache.fluss.client.tiering; import org.apache.fluss.annotation.VisibleForTesting; import org.apache.fluss.client.metadata.MetadataUpdater; @@ -197,7 +197,7 @@ public void commit( } } - void commit( + public void commit( long tableId, long snapshotId, String lakeBucketTieredOffsetsPath, @@ -233,14 +233,6 @@ void commit( } } - /** - * Converts the prepare commit parameters to a {@link PrepareLakeTableSnapshotRequest}. - * - * @param tableId the table ID - * @param tablePath the table path - * @param logEndOffsets the log end offsets for each bucket - * @return the prepared commit request - */ private PrepareLakeTableSnapshotRequest toPrepareLakeTableSnapshotRequest( long tableId, TablePath tablePath, Map logEndOffsets) { PrepareLakeTableSnapshotRequest prepareLakeTableSnapshotRequest = @@ -264,29 +256,6 @@ private PrepareLakeTableSnapshotRequest toPrepareLakeTableSnapshotRequest( return prepareLakeTableSnapshotRequest; } - /** - * Converts the commit parameters to a {@link CommitLakeTableSnapshotRequest}. - * - *

This method creates a request that includes: - * - *

- * - * @param tableId the table ID - * @param snapshotId the lake snapshot ID - * @param tieredBucketOffsetsPath the file path where the tiered bucket offsets is stored - * @param readableBucketTieredOffsetsPath the file path where the readable bucket offsets is - * stored - * @param logEndOffsets the log end offsets for each bucket - * @param logMaxTieredTimestamps the max tiered timestamps for each bucket - * @param earliestSnapshotIDToKeep the earliest snapshot ID to keep. Null means keep only the - * latest (discard all previous). -1 ({@link LakeCommitResult#KEEP_ALL_PREVIOUS}) means keep - * all previous snapshots (infinite retention). - * @return the commit request - */ private CommitLakeTableSnapshotRequest toCommitLakeTableSnapshotRequest( long tableId, long snapshotId, @@ -298,12 +267,10 @@ private CommitLakeTableSnapshotRequest toCommitLakeTableSnapshotRequest( CommitLakeTableSnapshotRequest commitLakeTableSnapshotRequest = new CommitLakeTableSnapshotRequest(); - // Add lake table snapshot metadata PbLakeTableSnapshotMetadata pbLakeTableSnapshotMetadata = commitLakeTableSnapshotRequest.addLakeTableSnapshotMetadata(); pbLakeTableSnapshotMetadata.setSnapshotId(snapshotId); pbLakeTableSnapshotMetadata.setTableId(tableId); - // tiered snapshot file path is equal to readable snapshot currently pbLakeTableSnapshotMetadata.setTieredBucketOffsetsFilePath(tieredBucketOffsetsPath); if (readableBucketTieredOffsetsPath != null) { pbLakeTableSnapshotMetadata.setReadableBucketOffsetsFilePath( @@ -313,8 +280,6 @@ private CommitLakeTableSnapshotRequest toCommitLakeTableSnapshotRequest( pbLakeTableSnapshotMetadata.setEarliestSnapshotIdToKeep(earliestSnapshotIDToKeep); } - // Add PbLakeTableSnapshotInfo for metrics reporting (to notify tablet servers about - // synchronized log end offsets and max timestamps) if (!logEndOffsets.isEmpty()) { commitLakeTableSnapshotRequest = addLogEndOffsets( @@ -328,7 +293,7 @@ private CommitLakeTableSnapshotRequest toCommitLakeTableSnapshotRequest( } @VisibleForTesting - protected CommitLakeTableSnapshotRequest addLogEndOffsets( + public CommitLakeTableSnapshotRequest addLogEndOffsets( CommitLakeTableSnapshotRequest commitLakeTableSnapshotRequest, long tableId, long snapshotId, @@ -358,7 +323,7 @@ protected CommitLakeTableSnapshotRequest addLogEndOffsets( } @VisibleForTesting - CoordinatorGateway getCoordinatorGateway() { + public CoordinatorGateway getCoordinatorGateway() { return coordinatorGateway; } diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TableBucketWriteResult.java b/fluss-client/src/main/java/org/apache/fluss/client/tiering/TableBucketWriteResult.java similarity index 98% rename from fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TableBucketWriteResult.java rename to fluss-client/src/main/java/org/apache/fluss/client/tiering/TableBucketWriteResult.java index abec3c6c21f..de87f009481 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TableBucketWriteResult.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/tiering/TableBucketWriteResult.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.fluss.flink.tiering.source; +package org.apache.fluss.client.tiering; import org.apache.fluss.lake.writer.LakeWriter; import org.apache.fluss.metadata.TableBucket; diff --git a/fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringCommitResult.java b/fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringCommitResult.java new file mode 100644 index 00000000000..72e466fb69b --- /dev/null +++ b/fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringCommitResult.java @@ -0,0 +1,52 @@ +/* + * 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.fluss.client.tiering; + +import org.apache.fluss.lake.committer.TieringStats; + +import javax.annotation.Nullable; + +/** + * The result of one table's commit round, holding the lake committable (nullable for empty commits + * where no data was written) and the associated tiering statistics. + * + * @param the type of the lake committable + */ +public class TieringCommitResult { + + /** The lake committable, or {@code null} if nothing was written in this round. */ + @Nullable private final Committable committable; + + /** Per-table tiering statistics collected during this round. */ + @Nullable private final TieringStats stats; + + public TieringCommitResult(@Nullable Committable committable, @Nullable TieringStats stats) { + this.committable = committable; + this.stats = stats; + } + + @Nullable + public Committable getCommittable() { + return committable; + } + + @Nullable + public TieringStats getStats() { + return stats; + } +} diff --git a/fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringCommitter.java b/fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringCommitter.java new file mode 100644 index 00000000000..a7888751478 --- /dev/null +++ b/fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringCommitter.java @@ -0,0 +1,265 @@ +/* + * 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.fluss.client.tiering; + +import org.apache.fluss.client.admin.Admin; +import org.apache.fluss.client.metadata.LakeSnapshot; +import org.apache.fluss.config.Configuration; +import org.apache.fluss.exception.LakeTableSnapshotNotExistException; +import org.apache.fluss.lake.committer.CommittedLakeSnapshot; +import org.apache.fluss.lake.committer.LakeCommitResult; +import org.apache.fluss.lake.committer.LakeCommitter; +import org.apache.fluss.lake.writer.LakeTieringFactory; +import org.apache.fluss.metadata.TableBucket; +import org.apache.fluss.metadata.TableInfo; +import org.apache.fluss.metadata.TablePath; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import javax.annotation.Nullable; + +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; + +import static org.apache.fluss.lake.committer.LakeCommitter.FLUSS_LAKE_SNAP_BUCKET_OFFSET_PROPERTY; + +/** + * Engine-agnostic tiering committer that commits write results to a lake and the Fluss cluster. + * + *

This class extracts the core commit logic from engine-specific operators so that any compute + * engine (Flink, Spark, etc.) can reuse the same commit protocol. + * + * @param the type of individual write results produced by lake writers + * @param the type of the aggregated committable for the lake + */ +public class TieringCommitter { + + private static final Logger LOG = LoggerFactory.getLogger(TieringCommitter.class); + + /** + * Commits the collected write results for one table to the lake and Fluss. When buckets + * advanced their tiered offsets without writing any data (e.g. splits only covering empty WAL + * batches), an empty lake snapshot is committed to persist the tiering progress. + * + *

Always returns a non-null {@link TieringCommitResult}. When nothing needs to be committed, + * {@link TieringCommitResult#getCommittable()} is {@code null} and stats are {@code null}. + */ + public TieringCommitResult commitWriteResults( + Admin admin, + long tableId, + TablePath tablePath, + Configuration flussConf, + Configuration lakeTieringConfig, + LakeTieringFactory lakeTieringFactory, + FlussTableLakeSnapshotCommitter flussTableLakeSnapshotCommitter, + List> committableWriteResults) + throws Exception { + // filter down to buckets that actually produced data + List> nonEmptyResults = + committableWriteResults.stream() + .filter(r -> r.writeResult() != null) + .collect(Collectors.toList()); + + // collect tiered offsets from all buckets, including those finished without writing + // data, otherwise their splits would be regenerated forever; buckets with unknown + // progress (e.g. splits skipped in this round) are excluded + Map logEndOffsets = new HashMap<>(); + Map logMaxTieredTimestamps = new HashMap<>(); + for (TableBucketWriteResult writeResult : committableWriteResults) { + if (writeResult.logEndOffset() < 0) { + continue; + } + TableBucket tableBucket = writeResult.tableBucket(); + logEndOffsets.put(tableBucket, writeResult.logEndOffset()); + if (writeResult.maxTimestamp() >= 0) { + logMaxTieredTimestamps.put(tableBucket, writeResult.maxTimestamp()); + } + } + + // nothing was written and no tiered offset advanced — nothing to commit + if (nonEmptyResults.isEmpty() && logEndOffsets.isEmpty()) { + LOG.info( + "Commit tiering write results is empty for table {}, table path {}", + tableId, + tablePath); + return new TieringCommitResult<>(null, null); + } + + if (nonEmptyResults.isEmpty()) { + LOG.info( + "No data was written for table {} (table path {}) but buckets {} advanced " + + "their tiered offsets, committing an empty lake snapshot to " + + "persist the tiering progress.", + tableId, + tablePath, + logEndOffsets.keySet()); + } + + // Check if the table was dropped and recreated during tiering. + // If the current table id differs from the committable's table id, fail this commit + // to avoid dirty commit to a newly created table. + TableInfo currentTableInfo = admin.getTableInfo(tablePath).get(); + if (currentTableInfo.getTableId() != tableId) { + throw new IllegalStateException( + String.format( + "The current table id %s for table path %s is different from the table id %s in the committable. " + + "This usually happens when a table was dropped and recreated during tiering. " + + "Aborting commit to prevent dirty commit.", + currentTableInfo.getTableId(), tablePath, tableId)); + } + + try (LakeCommitter lakeCommitter = + lakeTieringFactory.createLakeCommitter( + new TieringCommitterInitContext( + tablePath, currentTableInfo, lakeTieringConfig, flussConf))) { + List writeResults = + nonEmptyResults.stream() + .map(TableBucketWriteResult::writeResult) + .collect(Collectors.toList()); + + // to committable + Committable committable = lakeCommitter.toCommittable(writeResults); + // before commit to lake, check fluss not missing any lake snapshot committed by fluss + LakeSnapshot flussCurrentLakeSnapshot = getLatestLakeSnapshot(admin, tablePath); + checkFlussNotMissingLakeSnapshot( + flussTableLakeSnapshotCommitter, + tablePath, + tableId, + lakeCommitter, + committable, + flussCurrentLakeSnapshot == null + ? null + : flussCurrentLakeSnapshot.getSnapshotId()); + + // get the lake bucket offsets file storing the log end offsets + String lakeBucketTieredOffsetsFile = + flussTableLakeSnapshotCommitter.prepareLakeSnapshot( + tableId, tablePath, logEndOffsets); + + // record the lake snapshot bucket offsets file to snapshot property + Map snapshotProperties = + Collections.singletonMap( + FLUSS_LAKE_SNAP_BUCKET_OFFSET_PROPERTY, lakeBucketTieredOffsetsFile); + LakeCommitResult lakeCommitResult = + lakeCommitter.commit(committable, snapshotProperties); + // commit to fluss + flussTableLakeSnapshotCommitter.commit( + tableId, + tablePath, + lakeCommitResult, + lakeBucketTieredOffsetsFile, + logEndOffsets, + logMaxTieredTimestamps); + return new TieringCommitResult<>(committable, lakeCommitResult.getTieringStats()); + } + } + + @Nullable + private LakeSnapshot getLatestLakeSnapshot(Admin admin, TablePath tablePath) throws Exception { + LakeSnapshot flussCurrentLakeSnapshot; + try { + flussCurrentLakeSnapshot = admin.getLatestLakeSnapshot(tablePath).get(); + } catch (Exception e) { + Throwable throwable = e.getCause(); + if (throwable instanceof LakeTableSnapshotNotExistException) { + flussCurrentLakeSnapshot = null; + } else { + throw e; + } + } + return flussCurrentLakeSnapshot; + } + + private void checkFlussNotMissingLakeSnapshot( + FlussTableLakeSnapshotCommitter flussTableLakeSnapshotCommitter, + TablePath tablePath, + long tableId, + LakeCommitter lakeCommitter, + Committable committable, + Long flussCurrentLakeSnapshot) + throws Exception { + // get Fluss missing lake snapshot in Lake + CommittedLakeSnapshot missingCommittedSnapshot = + lakeCommitter.getMissingLakeSnapshot(flussCurrentLakeSnapshot); + + // fluss's known snapshot is less than lake snapshot committed by fluss + // fail this commit since the data is read from the log end-offset of a invalid fluss + // known lake snapshot, which means the data already has been committed to lake, + // not to commit to lake to avoid data duplicated + if (missingCommittedSnapshot != null) { + String lakeSnapshotOffsetPath = + missingCommittedSnapshot + .getSnapshotProperties() + .get(FLUSS_LAKE_SNAP_BUCKET_OFFSET_PROPERTY); + + // should only will happen in v0.7 which won't put offsets info + // to properties + if (lakeSnapshotOffsetPath == null) { + throw new IllegalStateException( + String.format( + "Can't find %s field from snapshot property.", + FLUSS_LAKE_SNAP_BUCKET_OFFSET_PROPERTY)); + } + + // the fluss-offsets will be a json string if it's tiered by v0.8, + // since this code path should be rare, we do not consider backward compatibility + // and throw IllegalStateException directly + String trimmedPath = lakeSnapshotOffsetPath.trim(); + if (trimmedPath.contains("{")) { + throw new IllegalStateException( + String.format( + "The %s field in snapshot property is a JSON string (tiered by v0.8), " + + "which is not supported to restore. Snapshot ID: %d, Table: {tablePath=%s, tableId=%d}.", + FLUSS_LAKE_SNAP_BUCKET_OFFSET_PROPERTY, + missingCommittedSnapshot.getLakeSnapshotId(), + tablePath, + tableId)); + } + + // commit this missing snapshot to fluss + flussTableLakeSnapshotCommitter.commit( + tableId, + missingCommittedSnapshot.getLakeSnapshotId(), + lakeSnapshotOffsetPath, + // don't care readable snapshot and offsets, + null, + // use empty log offsets, log max timestamp, since we can't know that + // in last tiering, it doesn't matter for they are just used to + // report metrics + Collections.emptyMap(), + Collections.emptyMap(), + LakeCommitResult.KEEP_ALL_PREVIOUS); + // abort this committable to delete the written files + lakeCommitter.abort(committable); + throw new IllegalStateException( + String.format( + "The current Fluss's lake snapshot %d is less than" + + " lake actual snapshot %d committed by Fluss for table: {tablePath=%s, tableId=%d}," + + " missing snapshot: %s.", + flussCurrentLakeSnapshot, + missingCommittedSnapshot.getLakeSnapshotId(), + tablePath, + tableId, + missingCommittedSnapshot)); + } + } +} diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/committer/TieringCommitterInitContext.java b/fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringCommitterInitContext.java similarity index 97% rename from fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/committer/TieringCommitterInitContext.java rename to fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringCommitterInitContext.java index 79b7aaeb4f6..a4f5934d655 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/committer/TieringCommitterInitContext.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringCommitterInitContext.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.fluss.flink.tiering.committer; +package org.apache.fluss.client.tiering; import org.apache.fluss.config.Configuration; import org.apache.fluss.lake.committer.CommitterInitContext; diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split/TieringLogSplit.java b/fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringLogSplit.java similarity index 99% rename from fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split/TieringLogSplit.java rename to fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringLogSplit.java index a1e983b5030..194851e2bcf 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split/TieringLogSplit.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringLogSplit.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.fluss.flink.tiering.source.split; +package org.apache.fluss.client.tiering; import org.apache.fluss.metadata.TableBucket; import org.apache.fluss.metadata.TablePath; diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split/TieringSnapshotSplit.java b/fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringSnapshotSplit.java similarity index 99% rename from fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split/TieringSnapshotSplit.java rename to fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringSnapshotSplit.java index 8c3bad36452..f4e7901ec71 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split/TieringSnapshotSplit.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringSnapshotSplit.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.fluss.flink.tiering.source.split; +package org.apache.fluss.client.tiering; import org.apache.fluss.metadata.TableBucket; import org.apache.fluss.metadata.TablePath; diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split/TieringSplit.java b/fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringSplit.java similarity index 96% rename from fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split/TieringSplit.java rename to fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringSplit.java index b41aef3f707..dae0d3dad9b 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split/TieringSplit.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringSplit.java @@ -15,19 +15,18 @@ * limitations under the License. */ -package org.apache.fluss.flink.tiering.source.split; +package org.apache.fluss.client.tiering; import org.apache.fluss.metadata.TableBucket; import org.apache.fluss.metadata.TablePath; -import org.apache.flink.api.connector.source.SourceSplit; - import javax.annotation.Nullable; +import java.io.Serializable; import java.util.Objects; /** The base table split for tiering service. */ -public abstract class TieringSplit implements SourceSplit { +public abstract class TieringSplit implements Serializable { public static final byte TIERING_SNAPSHOT_SPLIT_FLAG = 1; public static final byte TIERING_LOG_SPLIT_FLAG = 2; @@ -91,6 +90,9 @@ public TieringSplit( this.tieringRoundTimestamp = tieringRoundTimestamp; } + /** Returns the unique identifier for this split. */ + public abstract String splitId(); + /** Checks whether this split is a primary key table split to tier. */ public final boolean isTieringSnapshotSplit() { return getClass() == TieringSnapshotSplit.class; @@ -128,7 +130,7 @@ public TieringLogSplit asTieringLogSplit() { return (TieringLogSplit) this; } - protected byte splitKind() { + public byte splitKind() { if (isTieringSnapshotSplit()) { return TIERING_SNAPSHOT_SPLIT_FLAG; } else if (isTieringLogSplit()) { diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split/TieringSplitGenerator.java b/fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringSplitGenerator.java similarity index 98% rename from fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split/TieringSplitGenerator.java rename to fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringSplitGenerator.java index a4b2638f309..9e7f06b98d5 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split/TieringSplitGenerator.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringSplitGenerator.java @@ -15,13 +15,14 @@ * limitations under the License. */ -package org.apache.fluss.flink.tiering.source.split; +package org.apache.fluss.client.tiering; import org.apache.fluss.client.admin.Admin; import org.apache.fluss.client.initializer.BucketOffsetsRetrieverImpl; import org.apache.fluss.client.initializer.OffsetsInitializer.BucketOffsetsRetriever; import org.apache.fluss.client.metadata.KvSnapshots; import org.apache.fluss.client.metadata.LakeSnapshot; +import org.apache.fluss.exception.FlussRuntimeException; import org.apache.fluss.exception.LakeTableSnapshotNotExistException; import org.apache.fluss.metadata.PartitionInfo; import org.apache.fluss.metadata.TableBucket; @@ -29,7 +30,6 @@ import org.apache.fluss.metadata.TablePath; import org.apache.fluss.utils.ExceptionUtils; -import org.apache.flink.util.FlinkRuntimeException; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -75,7 +75,7 @@ public List generateTableSplits(TableInfo tableInfo) throws Except if (t instanceof LakeTableSnapshotNotExistException) { lakeSnapshotInfo = null; } else { - throw new FlinkRuntimeException( + throw new FlussRuntimeException( String.format( "Failed to get table snapshot for table %s", tableInfo.getTablePath()), @@ -127,7 +127,7 @@ private List generatePartitionTableSplit( .getLatestKvSnapshots(tableInfo.getTablePath(), partitionName) .get(); } catch (Exception e) { - throw new FlinkRuntimeException( + throw new FlussRuntimeException( String.format( "Failed to get table snapshot for table %s and partition %s", tableInfo.getTablePath(), partitionName), @@ -163,7 +163,7 @@ private List generateNonPartitionedTableSplit( try { latestKvSnapshots = flussAdmin.getLatestKvSnapshots(tableInfo.getTablePath()).get(); } catch (Exception e) { - throw new FlinkRuntimeException( + throw new FlussRuntimeException( String.format( "Failed to get table snapshot for table %s", tableInfo.getTablePath()), diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringWriterInitContext.java b/fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringWriterInitContext.java similarity index 98% rename from fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringWriterInitContext.java rename to fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringWriterInitContext.java index f67b44176be..c8fdd1fc124 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringWriterInitContext.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringWriterInitContext.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.fluss.flink.tiering.source; +package org.apache.fluss.client.tiering; import org.apache.fluss.lake.writer.WriterInitContext; import org.apache.fluss.metadata.TableBucket; diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/committer/CommittableMessageTypeInfo.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/committer/CommittableMessageTypeInfo.java index d541721d19f..0fc76e9b39e 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/committer/CommittableMessageTypeInfo.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/committer/CommittableMessageTypeInfo.java @@ -17,8 +17,8 @@ package org.apache.fluss.flink.tiering.committer; +import org.apache.fluss.client.tiering.TableBucketWriteResult; import org.apache.fluss.flink.adapter.TypeInformationAdapter; -import org.apache.fluss.flink.tiering.source.TableBucketWriteResult; import org.apache.fluss.lake.serializer.SimpleVersionedSerializer; import org.apache.flink.api.common.typeinfo.TypeInformation; diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/committer/TieringCommitOperator.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/committer/TieringCommitOperator.java index 477beea1134..6e1a2a07e6e 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/committer/TieringCommitOperator.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/committer/TieringCommitOperator.java @@ -20,22 +20,17 @@ import org.apache.fluss.client.Connection; import org.apache.fluss.client.ConnectionFactory; import org.apache.fluss.client.admin.Admin; -import org.apache.fluss.client.metadata.LakeSnapshot; +import org.apache.fluss.client.tiering.FlussTableLakeSnapshotCommitter; +import org.apache.fluss.client.tiering.TableBucketWriteResult; +import org.apache.fluss.client.tiering.TieringCommitResult; +import org.apache.fluss.client.tiering.TieringCommitter; import org.apache.fluss.config.Configuration; -import org.apache.fluss.exception.LakeTableSnapshotNotExistException; import org.apache.fluss.flink.tiering.event.FailedTieringEvent; import org.apache.fluss.flink.tiering.event.FinishedTieringEvent; -import org.apache.fluss.flink.tiering.source.TableBucketWriteResult; import org.apache.fluss.flink.tiering.source.TieringSource; -import org.apache.fluss.lake.committer.CommittedLakeSnapshot; -import org.apache.fluss.lake.committer.LakeCommitResult; -import org.apache.fluss.lake.committer.LakeCommitter; -import org.apache.fluss.lake.committer.TieringStats; import org.apache.fluss.lake.writer.LakeTieringFactory; import org.apache.fluss.lake.writer.LakeWriter; import org.apache.fluss.metadata.TableBucket; -import org.apache.fluss.metadata.TableInfo; -import org.apache.fluss.metadata.TablePath; import org.apache.fluss.utils.ExceptionUtils; import org.apache.flink.runtime.operators.coordination.OperatorEventGateway; @@ -48,15 +43,12 @@ import javax.annotation.Nullable; import java.util.ArrayList; -import java.util.Collections; import java.util.HashMap; import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Set; -import java.util.stream.Collectors; -import static org.apache.fluss.lake.committer.LakeCommitter.FLUSS_LAKE_SNAP_BUCKET_OFFSET_PROPERTY; import static org.apache.fluss.utils.Preconditions.checkState; /** @@ -85,6 +77,7 @@ public class TieringCommitOperator private final Configuration lakeTieringConfig; private final LakeTieringFactory lakeTieringFactory; private final FlussTableLakeSnapshotCommitter flussTableLakeSnapshotCommitter; + private final TieringCommitter tieringCommitter; private Connection connection; private Admin admin; @@ -96,21 +89,9 @@ public class TieringCommitOperator collectedTableBucketWriteResults; /** - * The result of one table's commit round, holding the lake committable (nullable for empty - * commits where no data was written) and the associated tiering statistics. + * The result of one table's commit round is now represented by {@link TieringCommitResult} from + * fluss-client. */ - private final class CommitResult { - /** The lake committable, or {@code null} if nothing was written in this round. */ - @Nullable final Committable committable; - /** Per-table tiering statistics collected during this round. */ - @Nullable final TieringStats stats; - - CommitResult(@Nullable Committable committable, @Nullable TieringStats stats) { - this.committable = committable; - this.stats = stats; - } - } - public TieringCommitOperator( StreamOperatorParameters> parameters, Configuration flussConf, @@ -118,6 +99,7 @@ public TieringCommitOperator( LakeTieringFactory lakeTieringFactory) { this.lakeTieringFactory = lakeTieringFactory; this.flussTableLakeSnapshotCommitter = new FlussTableLakeSnapshotCommitter(flussConf); + this.tieringCommitter = new TieringCommitter<>(); this.collectedTableBucketWriteResults = new HashMap<>(); this.flussConfig = flussConf; this.lakeTieringConfig = lakeTieringConfig; @@ -152,20 +134,26 @@ public void processElement(StreamRecord> str if (committableWriteResults != null) { try { - CommitResult commitResult = - commitWriteResults( + TieringCommitResult commitResult = + tieringCommitter.commitWriteResults( + admin, tableId, tableBucketWriteResult.tablePath(), + flussConfig, + lakeTieringConfig, + lakeTieringFactory, + flussTableLakeSnapshotCommitter, committableWriteResults); // only emit downstream when actual data was written - if (commitResult.committable != null) { + if (commitResult.getCommittable() != null) { output.collect( - new StreamRecord<>(new CommittableMessage<>(commitResult.committable))); + new StreamRecord<>( + new CommittableMessage<>(commitResult.getCommittable()))); } // notify that the table id has been finished tier operatorEventGateway.sendEventToCoordinator( new SourceEventWrapper( - new FinishedTieringEvent(tableId, commitResult.stats))); + new FinishedTieringEvent(tableId, commitResult.getStats()))); } catch (Exception e) { // if any exception happens, send to source coordinator to mark it as failed operatorEventGateway.sendEventToCoordinator( @@ -181,205 +169,6 @@ public void processElement(StreamRecord> str } } - /** - * Commits the collected write results for one table to the lake and Fluss. When buckets - * advanced their tiered offsets without writing any data (e.g. splits only covering empty WAL - * batches), an empty lake snapshot is committed to persist the tiering progress. - */ - private CommitResult commitWriteResults( - long tableId, - TablePath tablePath, - List> committableWriteResults) - throws Exception { - // filter down to buckets that actually produced data - List> nonEmptyResults = - committableWriteResults.stream() - .filter(r -> r.writeResult() != null) - .collect(Collectors.toList()); - - // collect tiered offsets from all buckets, including those finished without writing - // data, otherwise their splits would be regenerated forever; buckets with unknown - // progress (e.g. splits skipped in this round) are excluded - Map logEndOffsets = new HashMap<>(); - Map logMaxTieredTimestamps = new HashMap<>(); - for (TableBucketWriteResult writeResult : committableWriteResults) { - if (writeResult.logEndOffset() < 0) { - continue; - } - TableBucket tableBucket = writeResult.tableBucket(); - logEndOffsets.put(tableBucket, writeResult.logEndOffset()); - if (writeResult.maxTimestamp() >= 0) { - logMaxTieredTimestamps.put(tableBucket, writeResult.maxTimestamp()); - } - } - - // nothing was written and no tiered offset advanced — nothing to commit - if (nonEmptyResults.isEmpty() && logEndOffsets.isEmpty()) { - LOG.info( - "Commit tiering write results is empty for table {}, table path {}", - tableId, - tablePath); - return new CommitResult(null, null); - } - - if (nonEmptyResults.isEmpty()) { - LOG.info( - "No data was written for table {} (table path {}) but buckets {} advanced " - + "their tiered offsets, committing an empty lake snapshot to " - + "persist the tiering progress.", - tableId, - tablePath, - logEndOffsets.keySet()); - } - - // Check if the table was dropped and recreated during tiering. - // If the current table id differs from the committable's table id, fail this commit - // to avoid dirty commit to a newly created table. - TableInfo currentTableInfo = admin.getTableInfo(tablePath).get(); - if (currentTableInfo.getTableId() != tableId) { - throw new IllegalStateException( - String.format( - "The current table id %s for table path %s is different from the table id %s in the committable. " - + "This usually happens when a table was dropped and recreated during tiering. " - + "Aborting commit to prevent dirty commit.", - currentTableInfo.getTableId(), tablePath, tableId)); - } - - try (LakeCommitter lakeCommitter = - lakeTieringFactory.createLakeCommitter( - new TieringCommitterInitContext( - tablePath, currentTableInfo, lakeTieringConfig, flussConfig))) { - List writeResults = - nonEmptyResults.stream() - .map(TableBucketWriteResult::writeResult) - .collect(Collectors.toList()); - - // to committable - Committable committable = lakeCommitter.toCommittable(writeResults); - // before commit to lake, check fluss not missing any lake snapshot committed by fluss - LakeSnapshot flussCurrentLakeSnapshot = getLatestLakeSnapshot(tablePath); - checkFlussNotMissingLakeSnapshot( - tablePath, - tableId, - lakeCommitter, - committable, - flussCurrentLakeSnapshot == null - ? null - : flussCurrentLakeSnapshot.getSnapshotId()); - - // get the lake bucket offsets file storing the log end offsets - String lakeBucketTieredOffsetsFile = - flussTableLakeSnapshotCommitter.prepareLakeSnapshot( - tableId, tablePath, logEndOffsets); - - // record the lake snapshot bucket offsets file to snapshot property - Map snapshotProperties = - Collections.singletonMap( - FLUSS_LAKE_SNAP_BUCKET_OFFSET_PROPERTY, lakeBucketTieredOffsetsFile); - LakeCommitResult lakeCommitResult = - lakeCommitter.commit(committable, snapshotProperties); - // commit to fluss - flussTableLakeSnapshotCommitter.commit( - tableId, - tablePath, - lakeCommitResult, - lakeBucketTieredOffsetsFile, - logEndOffsets, - logMaxTieredTimestamps); - return new CommitResult(committable, lakeCommitResult.getTieringStats()); - } - } - - @Nullable - private LakeSnapshot getLatestLakeSnapshot(TablePath tablePath) throws Exception { - LakeSnapshot flussCurrentLakeSnapshot; - try { - flussCurrentLakeSnapshot = admin.getLatestLakeSnapshot(tablePath).get(); - } catch (Exception e) { - Throwable throwable = e.getCause(); - if (throwable instanceof LakeTableSnapshotNotExistException) { - // do-nothing - flussCurrentLakeSnapshot = null; - } else { - throw e; - } - } - return flussCurrentLakeSnapshot; - } - - private void checkFlussNotMissingLakeSnapshot( - TablePath tablePath, - long tableId, - LakeCommitter lakeCommitter, - Committable committable, - Long flussCurrentLakeSnapshot) - throws Exception { - // get Fluss missing lake snapshot in Lake - CommittedLakeSnapshot missingCommittedSnapshot = - lakeCommitter.getMissingLakeSnapshot(flussCurrentLakeSnapshot); - - // fluss's known snapshot is less than lake snapshot committed by fluss - // fail this commit since the data is read from the log end-offset of a invalid fluss - // known lake snapshot, which means the data already has been committed to lake, - // not to commit to lake to avoid data duplicated - if (missingCommittedSnapshot != null) { - String lakeSnapshotOffsetPath = - missingCommittedSnapshot - .getSnapshotProperties() - .get(FLUSS_LAKE_SNAP_BUCKET_OFFSET_PROPERTY); - - // should only will happen in v0.7 which won't put offsets info - // to properties - if (lakeSnapshotOffsetPath == null) { - throw new IllegalStateException( - String.format( - "Can't find %s field from snapshot property.", - FLUSS_LAKE_SNAP_BUCKET_OFFSET_PROPERTY)); - } - - // the fluss-offsets will be a json string if it's tiered by v0.8, - // since this code path should be rare, we do not consider backward compatibility - // and throw IllegalStateException directly - String trimmedPath = lakeSnapshotOffsetPath.trim(); - if (trimmedPath.contains("{")) { - throw new IllegalStateException( - String.format( - "The %s field in snapshot property is a JSON string (tiered by v0.8), " - + "which is not supported to restore. Snapshot ID: %d, Table: {tablePath=%s, tableId=%d}.", - FLUSS_LAKE_SNAP_BUCKET_OFFSET_PROPERTY, - missingCommittedSnapshot.getLakeSnapshotId(), - tablePath, - tableId)); - } - - // commit this missing snapshot to fluss - flussTableLakeSnapshotCommitter.commit( - tableId, - missingCommittedSnapshot.getLakeSnapshotId(), - lakeSnapshotOffsetPath, - // don't care readable snapshot and offsets, - null, - // use empty log offsets, log max timestamp, since we can't know that - // in last tiering, it doesn't matter for they are just used to - // report metrics - Collections.emptyMap(), - Collections.emptyMap(), - LakeCommitResult.KEEP_ALL_PREVIOUS); - // abort this committable to delete the written files - lakeCommitter.abort(committable); - throw new IllegalStateException( - String.format( - "The current Fluss's lake snapshot %d is less than" - + " lake actual snapshot %d committed by Fluss for table: {tablePath=%s, tableId=%d}," - + " missing snapshot: %s.", - flussCurrentLakeSnapshot, - missingCommittedSnapshot.getLakeSnapshotId(), - tablePath, - tableId, - missingCommittedSnapshot)); - } - } - private void registerTableBucketWriteResult( long tableId, TableBucketWriteResult tableBucketWriteResult) { collectedTableBucketWriteResults diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/committer/TieringCommitOperatorFactory.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/committer/TieringCommitOperatorFactory.java index efced7aeabb..4c573ee8457 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/committer/TieringCommitOperatorFactory.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/committer/TieringCommitOperatorFactory.java @@ -17,8 +17,8 @@ package org.apache.fluss.flink.tiering.committer; +import org.apache.fluss.client.tiering.TableBucketWriteResult; import org.apache.fluss.config.Configuration; -import org.apache.fluss.flink.tiering.source.TableBucketWriteResult; import org.apache.fluss.lake.writer.LakeTieringFactory; import org.apache.flink.streaming.api.operators.AbstractStreamOperatorFactory; diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TableBucketWriteResultEmitter.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TableBucketWriteResultEmitter.java index 2b12337f849..8421d487569 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TableBucketWriteResultEmitter.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TableBucketWriteResultEmitter.java @@ -17,6 +17,7 @@ package org.apache.fluss.flink.tiering.source; +import org.apache.fluss.client.tiering.TableBucketWriteResult; import org.apache.fluss.flink.tiering.source.state.TieringSplitState; import org.apache.fluss.lake.committer.LakeCommitter; diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TableBucketWriteResultSerializer.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TableBucketWriteResultSerializer.java index 36517609557..1ce19d463db 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TableBucketWriteResultSerializer.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TableBucketWriteResultSerializer.java @@ -17,6 +17,7 @@ package org.apache.fluss.flink.tiering.source; +import org.apache.fluss.client.tiering.TableBucketWriteResult; import org.apache.fluss.metadata.TableBucket; import org.apache.fluss.metadata.TablePath; diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TableBucketWriteResultTypeInfo.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TableBucketWriteResultTypeInfo.java index 424673c26c1..9322f6ee7b3 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TableBucketWriteResultTypeInfo.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TableBucketWriteResultTypeInfo.java @@ -17,6 +17,7 @@ package org.apache.fluss.flink.tiering.source; +import org.apache.fluss.client.tiering.TableBucketWriteResult; import org.apache.fluss.flink.adapter.TypeInformationAdapter; import org.apache.fluss.lake.serializer.SimpleVersionedSerializer; diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringSource.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringSource.java index 696d4722b8a..d892ad4714c 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringSource.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringSource.java @@ -19,9 +19,10 @@ import org.apache.fluss.client.Connection; import org.apache.fluss.client.ConnectionFactory; +import org.apache.fluss.client.tiering.TableBucketWriteResult; import org.apache.fluss.config.Configuration; import org.apache.fluss.flink.tiering.source.enumerator.TieringSourceEnumerator; -import org.apache.fluss.flink.tiering.source.split.TieringSplit; +import org.apache.fluss.flink.tiering.source.split.FlinkTieringSplit; import org.apache.fluss.flink.tiering.source.split.TieringSplitSerializer; import org.apache.fluss.flink.tiering.source.state.TieringSourceEnumeratorState; import org.apache.fluss.flink.tiering.source.state.TieringSourceEnumeratorStateSerializer; @@ -56,7 +57,9 @@ */ public class TieringSource implements Source< - TableBucketWriteResult, TieringSplit, TieringSourceEnumeratorState> { + TableBucketWriteResult, + FlinkTieringSplit, + TieringSourceEnumeratorState> { public static final String TIERING_SOURCE_TRANSFORMATION_UID = "$$fluss_tiering_source_operator$$"; @@ -85,15 +88,15 @@ public Boundedness getBoundedness() { } @Override - public SplitEnumerator createEnumerator( - SplitEnumeratorContext splitEnumeratorContext) { + public SplitEnumerator createEnumerator( + SplitEnumeratorContext splitEnumeratorContext) { return new TieringSourceEnumerator( flussConf, splitEnumeratorContext, lakeTieringFactory, pollTieringTableIntervalMs); } @Override - public SplitEnumerator restoreEnumerator( - SplitEnumeratorContext splitEnumeratorContext, + public SplitEnumerator restoreEnumerator( + SplitEnumeratorContext splitEnumeratorContext, TieringSourceEnumeratorState tieringSourceEnumeratorState) { // stateless operator return new TieringSourceEnumerator( @@ -101,7 +104,7 @@ public SplitEnumerator restoreEnumer } @Override - public SimpleVersionedSerializer getSplitSerializer() { + public SimpleVersionedSerializer getSplitSerializer() { return TieringSplitSerializer.INSTANCE; } @@ -112,7 +115,7 @@ public SimpleVersionedSerializer getSplitSerializer() { } @Override - public SourceReader, TieringSplit> createReader( + public SourceReader, FlinkTieringSplit> createReader( SourceReaderContext sourceReaderContext) { FutureCompletingBlockingQueue>> elementsQueue = new FutureCompletingBlockingQueue<>(); diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringSourceFetcherManager.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringSourceFetcherManager.java index ac72aad6643..63352b47661 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringSourceFetcherManager.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringSourceFetcherManager.java @@ -18,8 +18,9 @@ package org.apache.fluss.flink.tiering.source; +import org.apache.fluss.client.tiering.TableBucketWriteResult; import org.apache.fluss.flink.adapter.SingleThreadFetcherManagerAdapter; -import org.apache.fluss.flink.tiering.source.split.TieringSplit; +import org.apache.fluss.flink.tiering.source.split.FlinkTieringSplit; import org.apache.flink.configuration.Configuration; import org.apache.flink.connector.base.source.reader.RecordsWithSplitIds; @@ -40,14 +41,14 @@ */ public class TieringSourceFetcherManager extends SingleThreadFetcherManagerAdapter< - TableBucketWriteResult, TieringSplit> { + TableBucketWriteResult, FlinkTieringSplit> { private static final Logger LOG = LoggerFactory.getLogger(TieringSourceFetcherManager.class); public TieringSourceFetcherManager( FutureCompletingBlockingQueue>> elementsQueue, - Supplier, TieringSplit>> + Supplier, FlinkTieringSplit>> splitReaderSupplier, Configuration configuration, Consumer> splitFinishedHook) { @@ -64,7 +65,7 @@ public void markTableReachTieringMaxDuration(long tableId) { enqueueMarkTableReachTieringMaxDurationTask( splitFetcher, tableId)); } else { - SplitFetcher, TieringSplit> splitFetcher = + SplitFetcher, FlinkTieringSplit> splitFetcher = createSplitFetcher(); LOG.info( "fetchers is empty, enqueue marking tiering max duration for table {}", @@ -75,7 +76,7 @@ public void markTableReachTieringMaxDuration(long tableId) { } private void enqueueMarkTableReachTieringMaxDurationTask( - SplitFetcher, TieringSplit> splitFetcher, + SplitFetcher, FlinkTieringSplit> splitFetcher, long reachTieringDeadlineTable) { splitFetcher.enqueueTask( new SplitFetcherTask() { diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringSourceReader.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringSourceReader.java index 663bc0b82e1..f86747cb3e9 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringSourceReader.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringSourceReader.java @@ -20,10 +20,11 @@ import org.apache.fluss.annotation.Internal; import org.apache.fluss.annotation.VisibleForTesting; import org.apache.fluss.client.Connection; +import org.apache.fluss.client.tiering.TableBucketWriteResult; import org.apache.fluss.flink.adapter.SingleThreadMultiplexSourceReaderBaseAdapter; import org.apache.fluss.flink.tiering.event.TieringReachMaxDurationEvent; import org.apache.fluss.flink.tiering.source.metrics.TieringMetrics; -import org.apache.fluss.flink.tiering.source.split.TieringSplit; +import org.apache.fluss.flink.tiering.source.split.FlinkTieringSplit; import org.apache.fluss.flink.tiering.source.state.TieringSplitState; import org.apache.fluss.lake.writer.LakeTieringFactory; @@ -48,7 +49,7 @@ public final class TieringSourceReader extends SingleThreadMultiplexSourceReaderBaseAdapter< TableBucketWriteResult, TableBucketWriteResult, - TieringSplit, + FlinkTieringSplit, TieringSplitState> { private static final Logger LOG = LoggerFactory.getLogger(TieringSourceReader.class); @@ -144,13 +145,13 @@ protected void onSplitFinished(Map finishedSplitIds) } @Override - public List snapshotState(long checkpointId) { + public List snapshotState(long checkpointId) { // we return empty list to make source reader be stateless return Collections.emptyList(); } @Override - protected TieringSplitState initializedState(TieringSplit split) { + protected TieringSplitState initializedState(FlinkTieringSplit split) { if (split.isTieringSnapshotSplit()) { return new TieringSplitState(split); } else if (split.isTieringLogSplit()) { @@ -161,7 +162,7 @@ protected TieringSplitState initializedState(TieringSplit split) { } @Override - protected TieringSplit toSplitType(String splitId, TieringSplitState splitState) { + protected FlinkTieringSplit toSplitType(String splitId, TieringSplitState splitState) { return splitState.toSourceSplit(); } diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringSplitReader.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringSplitReader.java index 60b911b6f74..c0218f7f38c 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringSplitReader.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringSplitReader.java @@ -25,12 +25,15 @@ import org.apache.fluss.client.table.scanner.log.LogScanner; import org.apache.fluss.client.table.scanner.log.LogScannerImpl; import org.apache.fluss.client.table.scanner.log.ScanRecords; +import org.apache.fluss.client.tiering.TableBucketWriteResult; +import org.apache.fluss.client.tiering.TieringLogSplit; +import org.apache.fluss.client.tiering.TieringSnapshotSplit; +import org.apache.fluss.client.tiering.TieringSplit; +import org.apache.fluss.client.tiering.TieringWriterInitContext; import org.apache.fluss.flink.source.reader.BoundedSplitReader; import org.apache.fluss.flink.source.reader.RecordAndPos; import org.apache.fluss.flink.tiering.source.metrics.TieringMetrics; -import org.apache.fluss.flink.tiering.source.split.TieringLogSplit; -import org.apache.fluss.flink.tiering.source.split.TieringSnapshotSplit; -import org.apache.fluss.flink.tiering.source.split.TieringSplit; +import org.apache.fluss.flink.tiering.source.split.FlinkTieringSplit; import org.apache.fluss.lake.batch.ArrowRecordBatch; import org.apache.fluss.lake.writer.LakeTieringFactory; import org.apache.fluss.lake.writer.LakeWriter; @@ -74,7 +77,7 @@ /** The {@link SplitReader} implementation which will read Fluss and write to lake. */ public class TieringSplitReader - implements SplitReader, TieringSplit> { + implements SplitReader, FlinkTieringSplit> { private static final Logger LOG = LoggerFactory.getLogger(TieringSplitReader.class); @@ -231,14 +234,15 @@ public RecordsWithSplitIds> fetch() throws I } @Override - public void handleSplitsChanges(SplitsChange splitsChange) { + public void handleSplitsChanges(SplitsChange splitsChange) { if (!(splitsChange instanceof SplitsAddition)) { throw new UnsupportedOperationException( String.format( "The SplitChange type of %s is not supported.", splitsChange.getClass())); } - for (TieringSplit split : splitsChange.splits()) { + for (FlinkTieringSplit flinkSplit : splitsChange.splits()) { + TieringSplit split = flinkSplit.unwrap(); LOG.info("add split {}", split.splitId()); if (split.shouldSkipCurrentRound()) { // if the split is forced to ignore, diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/enumerator/TieringSourceEnumerator.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/enumerator/TieringSourceEnumerator.java index f1bf98527af..2d8912fa904 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/enumerator/TieringSourceEnumerator.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/enumerator/TieringSourceEnumerator.java @@ -22,13 +22,14 @@ import org.apache.fluss.client.ConnectionFactory; import org.apache.fluss.client.admin.Admin; import org.apache.fluss.client.metadata.MetadataUpdater; +import org.apache.fluss.client.tiering.TieringSplit; +import org.apache.fluss.client.tiering.TieringSplitGenerator; import org.apache.fluss.config.Configuration; import org.apache.fluss.flink.metrics.FlinkMetricRegistry; import org.apache.fluss.flink.tiering.event.FailedTieringEvent; import org.apache.fluss.flink.tiering.event.FinishedTieringEvent; import org.apache.fluss.flink.tiering.event.TieringReachMaxDurationEvent; -import org.apache.fluss.flink.tiering.source.split.TieringSplit; -import org.apache.fluss.flink.tiering.source.split.TieringSplitGenerator; +import org.apache.fluss.flink.tiering.source.split.FlinkTieringSplit; import org.apache.fluss.flink.tiering.source.state.TieringSourceEnumeratorState; import org.apache.fluss.lake.committer.TieringStats; import org.apache.fluss.lake.writer.LakeTieringFactory; @@ -95,17 +96,17 @@ * */ public class TieringSourceEnumerator - implements SplitEnumerator { + implements SplitEnumerator { private static final Logger LOG = LoggerFactory.getLogger(TieringSourceEnumerator.class); private final Configuration flussConf; - private final SplitEnumeratorContext context; + private final SplitEnumeratorContext context; private final LakeTieringFactory lakeTieringFactory; private final ScheduledExecutorService timerService; private final SplitEnumeratorMetricGroup enumeratorMetricGroup; private final long pollTieringTableIntervalMs; - private final List pendingSplits; + private final List pendingSplits; private final Set readersAwaitingSplit; private final Map tieringTableEpochs; @@ -127,7 +128,7 @@ public class TieringSourceEnumerator public TieringSourceEnumerator( Configuration flussConf, - SplitEnumeratorContext context, + SplitEnumeratorContext context, LakeTieringFactory lakeTieringFactory, long pollTieringTableIntervalMs) { this.flussConf = flussConf; @@ -210,7 +211,7 @@ public void handleSplitRequest(int subtaskId, @Nullable String requesterHostname } @Override - public void addSplitsBack(List splits, int subtaskId) { + public void addSplitsBack(List splits, int subtaskId) { readersAwaitingSplit.add(subtaskId); pendingSplits.addAll(splits); assignSplits(); @@ -338,7 +339,7 @@ protected void handleTableTieringReachMaxDuration( LOG.info("Table {}-{} reached max duration. Force completing.", tablePath, tableId); tieringReachMaxDurationsTables.add(tableId); - for (TieringSplit tieringSplit : pendingSplits) { + for (FlinkTieringSplit tieringSplit : pendingSplits) { if (tieringSplit.getTableBucket().getTableId() == tableId) { // mark this tiering split to skip the current round since the tiering for // this table has timed out, so the tiering source reader can skip them directly @@ -383,7 +384,7 @@ private void assignSplits() { continue; } if (!pendingSplits.isEmpty()) { - TieringSplit tieringSplit = pendingSplits.remove(0); + FlinkTieringSplit tieringSplit = pendingSplits.remove(0); context.assignSplit(tieringSplit, nextAwaitingReader); LOG.info("Assigning split {} to readers {}", tieringSplit, nextAwaitingReader); readersAwaitingSplit.remove(nextAwaitingReader); @@ -462,20 +463,20 @@ private void generateTieringSplits(Tuple3 tieringTable) // shuffle tiering split to avoid splits tiering skew // after introduce tiering max duration Collections.shuffle(tieringSplits); - tieringSplits = populateTieringRoundMetadata(tieringSplits); + List flinkSplits = populateTieringRoundMetadata(tieringSplits); LOG.info( "Generate Tiering {} splits for table {} with cost {}ms.", - tieringSplits.size(), + flinkSplits.size(), tieringTable.f2, System.currentTimeMillis() - start); - if (tieringSplits.isEmpty()) { + if (flinkSplits.isEmpty()) { LOG.info( "Generate Tiering splits for table {} is empty, no need to tier data.", tieringTable.f2.getTableName()); tieringTableEpochs.remove(tieringTable.f0); finishedTables.put(tieringTable.f0, TieringFinishInfo.from(tieringTable.f1)); } else { - pendingSplits.addAll(tieringSplits); + pendingSplits.addAll(flinkSplits); timerService.schedule( () -> @@ -497,18 +498,19 @@ private void generateTieringSplits(Tuple3 tieringTable) } } - private List populateTieringRoundMetadata(List tieringSplits) { + private List populateTieringRoundMetadata(List tieringSplits) { int numberOfSplits = tieringSplits.size(); if (numberOfSplits == 0) { return Collections.emptyList(); } long tieringRoundTimestamp = System.currentTimeMillis(); - List splitsWithMetadata = new ArrayList<>(numberOfSplits); + List splitsWithMetadata = new ArrayList<>(numberOfSplits); for (int splitIndex = 0; splitIndex < numberOfSplits; splitIndex++) { splitsWithMetadata.add( - tieringSplits - .get(splitIndex) - .copy(numberOfSplits, splitIndex, tieringRoundTimestamp)); + new FlinkTieringSplit( + tieringSplits + .get(splitIndex) + .copy(numberOfSplits, splitIndex, tieringRoundTimestamp))); } return splitsWithMetadata; } diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split/FlinkTieringSplit.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split/FlinkTieringSplit.java new file mode 100644 index 00000000000..ddf7d63f29d --- /dev/null +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split/FlinkTieringSplit.java @@ -0,0 +1,152 @@ +/* + * 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.fluss.flink.tiering.source.split; + +import org.apache.fluss.client.tiering.TieringLogSplit; +import org.apache.fluss.client.tiering.TieringSnapshotSplit; +import org.apache.fluss.client.tiering.TieringSplit; +import org.apache.fluss.metadata.TableBucket; +import org.apache.fluss.metadata.TablePath; + +import org.apache.flink.api.connector.source.SourceSplit; + +import javax.annotation.Nullable; + +import java.util.Objects; + +/** + * A Flink-specific wrapper for {@link TieringSplit} that implements Flink's {@link SourceSplit} + * interface. This adapter allows the engine-agnostic {@link TieringSplit} to be used within Flink's + * source framework. + */ +public class FlinkTieringSplit implements SourceSplit { + + private final TieringSplit tieringSplit; + + public FlinkTieringSplit(TieringSplit tieringSplit) { + this.tieringSplit = tieringSplit; + } + + @Override + public String splitId() { + return tieringSplit.splitId(); + } + + /** Returns the underlying engine-agnostic {@link TieringSplit}. */ + public TieringSplit unwrap() { + return tieringSplit; + } + + /** Checks whether this split is a primary key table split to tier. */ + public boolean isTieringSnapshotSplit() { + return tieringSplit.isTieringSnapshotSplit(); + } + + /** Checks whether this split is a log split to tier. */ + public boolean isTieringLogSplit() { + return tieringSplit.isTieringLogSplit(); + } + + /** Casts the underlying split into a {@link TieringSnapshotSplit}. */ + public TieringSnapshotSplit asTieringSnapshotSplit() { + return tieringSplit.asTieringSnapshotSplit(); + } + + /** Casts the underlying split into a {@link TieringLogSplit}. */ + public TieringLogSplit asTieringLogSplit() { + return tieringSplit.asTieringLogSplit(); + } + + /** + * Marks this split to skip reading data in the current round. Once called, the split will not + * be processed and data reading will be skipped. + */ + public void skipCurrentRound() { + tieringSplit.skipCurrentRound(); + } + + /** + * Returns whether this split should skip tiering data in the current round of tiering. + * + * @return true if the split should skip tiering data, false otherwise + */ + public boolean shouldSkipCurrentRound() { + return tieringSplit.shouldSkipCurrentRound(); + } + + public byte splitKind() { + return tieringSplit.splitKind(); + } + + public int getNumberOfSplits() { + return tieringSplit.getNumberOfSplits(); + } + + public int getSplitIndex() { + return tieringSplit.getSplitIndex(); + } + + public boolean isFirstSplit() { + return tieringSplit.isFirstSplit(); + } + + public long getTieringRoundTimestamp() { + return tieringSplit.getTieringRoundTimestamp(); + } + + public TablePath getTablePath() { + return tieringSplit.getTablePath(); + } + + public TableBucket getTableBucket() { + return tieringSplit.getTableBucket(); + } + + @Nullable + public String getPartitionName() { + return tieringSplit.getPartitionName(); + } + + public FlinkTieringSplit copy(int numberOfSplits) { + return new FlinkTieringSplit(tieringSplit.copy(numberOfSplits)); + } + + public FlinkTieringSplit copy(int numberOfSplits, int splitIndex, long tieringRoundTimestamp) { + return new FlinkTieringSplit( + tieringSplit.copy(numberOfSplits, splitIndex, tieringRoundTimestamp)); + } + + @Override + public boolean equals(Object object) { + if (!(object instanceof FlinkTieringSplit)) { + return false; + } + FlinkTieringSplit that = (FlinkTieringSplit) object; + return Objects.equals(tieringSplit, that.tieringSplit); + } + + @Override + public int hashCode() { + return Objects.hash(tieringSplit); + } + + @Override + public String toString() { + return "FlinkTieringSplit{" + tieringSplit + '}'; + } +} diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split/TieringSplitSerializer.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split/TieringSplitSerializer.java index e336ee4670a..493e4fc0327 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split/TieringSplitSerializer.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split/TieringSplitSerializer.java @@ -17,6 +17,9 @@ package org.apache.fluss.flink.tiering.source.split; +import org.apache.fluss.client.tiering.TieringLogSplit; +import org.apache.fluss.client.tiering.TieringSnapshotSplit; +import org.apache.fluss.client.tiering.TieringSplit; import org.apache.fluss.flink.tiering.source.TieringSource; import org.apache.fluss.flink.tiering.source.enumerator.TieringSourceEnumerator; import org.apache.fluss.metadata.TableBucket; @@ -28,14 +31,17 @@ import java.io.IOException; +import static org.apache.fluss.client.tiering.TieringSplit.TIERING_LOG_SPLIT_FLAG; +import static org.apache.fluss.client.tiering.TieringSplit.TIERING_SNAPSHOT_SPLIT_FLAG; + /** - * A serializer for the {@link TieringSplit}. + * A serializer for the {@link FlinkTieringSplit}. * *

This serializer is only used to serialize and deserialize splits sent from {@link * TieringSourceEnumerator} to {@link TieringSource} for network transmission. Therefore, it does * not need to consider compatibility. */ -public class TieringSplitSerializer implements SimpleVersionedSerializer { +public class TieringSplitSerializer implements SimpleVersionedSerializer { public static final TieringSplitSerializer INSTANCE = new TieringSplitSerializer(); @@ -44,9 +50,6 @@ public class TieringSplitSerializer implements SimpleVersionedSerializer SERIALIZER_CACHE = ThreadLocal.withInitial(() -> new DataOutputSerializer(64)); - private static final byte TIERING_SNAPSHOT_SPLIT_FLAG = 1; - private static final byte TIERING_LOG_SPLIT_FLAG = 2; - private static final int CURRENT_VERSION = VERSION_0; @Override @@ -55,8 +58,9 @@ public int getVersion() { } @Override - public byte[] serialize(TieringSplit split) throws IOException { + public byte[] serialize(FlinkTieringSplit flinkSplit) throws IOException { final DataOutputSerializer out = SERIALIZER_CACHE.get(); + TieringSplit split = flinkSplit.unwrap(); byte splitKind = split.splitKind(); out.writeByte(splitKind); @@ -111,7 +115,7 @@ public byte[] serialize(TieringSplit split) throws IOException { } @Override - public TieringSplit deserialize(int version, byte[] serialized) throws IOException { + public FlinkTieringSplit deserialize(int version, byte[] serialized) throws IOException { if (version != VERSION_0) { throw new IOException("Unknown version " + version); } @@ -145,36 +149,42 @@ public TieringSplit deserialize(int version, byte[] serialized) throws IOExcepti int splitIndex = in.readInt(); long tieringRoundTimestamp = in.readLong(); + TieringSplit tieringSplit; if (splitKind == TIERING_SNAPSHOT_SPLIT_FLAG) { // deserialize snapshot id long snapshotId = in.readLong(); // deserialize log offset of snapshot long logOffsetOfSnapshot = in.readLong(); - return new TieringSnapshotSplit( - tablePath, - tableBucket, - partitionName, - snapshotId, - logOffsetOfSnapshot, - numberOfSplits, - skipCurrentRound, - splitIndex, - tieringRoundTimestamp); - } else { + tieringSplit = + new TieringSnapshotSplit( + tablePath, + tableBucket, + partitionName, + snapshotId, + logOffsetOfSnapshot, + numberOfSplits, + skipCurrentRound, + splitIndex, + tieringRoundTimestamp); + } else if (splitKind == TIERING_LOG_SPLIT_FLAG) { // deserialize starting offset long startingOffset = in.readLong(); // deserialize starting offset long stoppingOffset = in.readLong(); - return new TieringLogSplit( - tablePath, - tableBucket, - partitionName, - startingOffset, - stoppingOffset, - numberOfSplits, - skipCurrentRound, - splitIndex, - tieringRoundTimestamp); + tieringSplit = + new TieringLogSplit( + tablePath, + tableBucket, + partitionName, + startingOffset, + stoppingOffset, + numberOfSplits, + skipCurrentRound, + splitIndex, + tieringRoundTimestamp); + } else { + throw new IOException("Unknown split kind " + splitKind); } + return new FlinkTieringSplit(tieringSplit); } } diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/state/TieringSplitState.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/state/TieringSplitState.java index 0690da62866..dec70219ed9 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/state/TieringSplitState.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/state/TieringSplitState.java @@ -17,12 +17,13 @@ package org.apache.fluss.flink.tiering.source.state; -import org.apache.fluss.flink.tiering.source.split.TieringLogSplit; -import org.apache.fluss.flink.tiering.source.split.TieringSnapshotSplit; -import org.apache.fluss.flink.tiering.source.split.TieringSplit; +import org.apache.fluss.client.tiering.TieringLogSplit; +import org.apache.fluss.client.tiering.TieringSnapshotSplit; +import org.apache.fluss.client.tiering.TieringSplit; +import org.apache.fluss.flink.tiering.source.split.FlinkTieringSplit; /** - * The state of a {@link TieringSplit}. + * The state of a {@link FlinkTieringSplit}. * *

Note: The tiering service adopts a stateless design and does not store any progress * information in state during checkpoints. All splits are re-requested from the Fluss cluster in @@ -30,37 +31,40 @@ */ public class TieringSplitState { - protected final TieringSplit tieringSplit; + protected final FlinkTieringSplit flinkTieringSplit; - public TieringSplitState(TieringSplit tieringSplit) { - this.tieringSplit = tieringSplit; + public TieringSplitState(FlinkTieringSplit flinkTieringSplit) { + this.flinkTieringSplit = flinkTieringSplit; } - public TieringSplit toSourceSplit() { + public FlinkTieringSplit toSourceSplit() { + TieringSplit tieringSplit = flinkTieringSplit.unwrap(); if (tieringSplit.isTieringSnapshotSplit()) { - final TieringSnapshotSplit split = (TieringSnapshotSplit) this.tieringSplit; - return new TieringSnapshotSplit( - split.getTablePath(), - split.getTableBucket(), - split.getPartitionName(), - split.getSnapshotId(), - split.getLogOffsetOfSnapshot(), - split.getNumberOfSplits(), - split.shouldSkipCurrentRound(), - split.getSplitIndex(), - split.getTieringRoundTimestamp()); + final TieringSnapshotSplit split = tieringSplit.asTieringSnapshotSplit(); + return new FlinkTieringSplit( + new TieringSnapshotSplit( + split.getTablePath(), + split.getTableBucket(), + split.getPartitionName(), + split.getSnapshotId(), + split.getLogOffsetOfSnapshot(), + split.getNumberOfSplits(), + split.shouldSkipCurrentRound(), + split.getSplitIndex(), + split.getTieringRoundTimestamp())); } else { - final TieringLogSplit split = (TieringLogSplit) tieringSplit; - return new TieringLogSplit( - split.getTablePath(), - split.getTableBucket(), - split.getPartitionName(), - split.getStartingOffset(), - split.getStoppingOffset(), - split.getNumberOfSplits(), - split.shouldSkipCurrentRound(), - split.getSplitIndex(), - split.getTieringRoundTimestamp()); + final TieringLogSplit split = tieringSplit.asTieringLogSplit(); + return new FlinkTieringSplit( + new TieringLogSplit( + split.getTablePath(), + split.getTableBucket(), + split.getPartitionName(), + split.getStartingOffset(), + split.getStoppingOffset(), + split.getNumberOfSplits(), + split.shouldSkipCurrentRound(), + split.getSplitIndex(), + split.getTieringRoundTimestamp())); } } } diff --git a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/committer/FlussTableLakeSnapshotCommitterTest.java b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/committer/FlussTableLakeSnapshotCommitterTest.java index 875190f4d74..0ad8b440d83 100644 --- a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/committer/FlussTableLakeSnapshotCommitterTest.java +++ b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/committer/FlussTableLakeSnapshotCommitterTest.java @@ -18,6 +18,7 @@ package org.apache.fluss.flink.tiering.committer; import org.apache.fluss.client.metadata.LakeSnapshot; +import org.apache.fluss.client.tiering.FlussTableLakeSnapshotCommitter; import org.apache.fluss.exception.ApiException; import org.apache.fluss.exception.LakeTableSnapshotNotExistException; import org.apache.fluss.flink.utils.FlinkTestBase; diff --git a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/committer/TieringCommitOperatorTest.java b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/committer/TieringCommitOperatorTest.java index 2e4d85d188e..1d95d22cf61 100644 --- a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/committer/TieringCommitOperatorTest.java +++ b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/committer/TieringCommitOperatorTest.java @@ -18,13 +18,14 @@ package org.apache.fluss.flink.tiering.committer; import org.apache.fluss.client.metadata.LakeSnapshot; +import org.apache.fluss.client.tiering.FlussTableLakeSnapshotCommitter; +import org.apache.fluss.client.tiering.TableBucketWriteResult; import org.apache.fluss.exception.LakeTableSnapshotNotExistException; import org.apache.fluss.flink.adapter.StreamOperatorParametersAdapter; import org.apache.fluss.flink.tiering.TestingLakeTieringFactory; import org.apache.fluss.flink.tiering.TestingWriteResult; import org.apache.fluss.flink.tiering.event.FailedTieringEvent; import org.apache.fluss.flink.tiering.event.FinishedTieringEvent; -import org.apache.fluss.flink.tiering.source.TableBucketWriteResult; import org.apache.fluss.flink.utils.FlinkTestBase; import org.apache.fluss.lake.committer.CommittedLakeSnapshot; import org.apache.fluss.metadata.TableBucket; diff --git a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/TableBucketWriteResultSerializerTest.java b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/TableBucketWriteResultSerializerTest.java index dbb40eae17e..588a6d59b9b 100644 --- a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/TableBucketWriteResultSerializerTest.java +++ b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/TableBucketWriteResultSerializerTest.java @@ -17,6 +17,7 @@ package org.apache.fluss.flink.tiering.source; +import org.apache.fluss.client.tiering.TableBucketWriteResult; import org.apache.fluss.flink.tiering.TestingWriteResult; import org.apache.fluss.metadata.TableBucket; import org.apache.fluss.metadata.TablePath; diff --git a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/TieringSourceReaderTest.java b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/TieringSourceReaderTest.java index 9e9de2c7929..30e85014ba6 100644 --- a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/TieringSourceReaderTest.java +++ b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/TieringSourceReaderTest.java @@ -20,12 +20,14 @@ import org.apache.fluss.client.Connection; import org.apache.fluss.client.ConnectionFactory; +import org.apache.fluss.client.tiering.TableBucketWriteResult; +import org.apache.fluss.client.tiering.TieringLogSplit; import org.apache.fluss.config.ConfigOptions; import org.apache.fluss.config.Configuration; import org.apache.fluss.flink.tiering.TestingLakeTieringFactory; import org.apache.fluss.flink.tiering.TestingWriteResult; import org.apache.fluss.flink.tiering.event.TieringReachMaxDurationEvent; -import org.apache.fluss.flink.tiering.source.split.TieringLogSplit; +import org.apache.fluss.flink.tiering.source.split.FlinkTieringSplit; import org.apache.fluss.flink.utils.FlinkTestBase; import org.apache.fluss.metadata.TableBucket; import org.apache.fluss.metadata.TablePath; @@ -78,7 +80,7 @@ void testHandleTieringReachMaxDurationEvent() throws Exception { TieringLogSplit split = new TieringLogSplit( tablePath, new TableBucket(tableId, 0), null, EARLIEST_OFFSET, 100); - reader.addSplits(Collections.singletonList(split)); + reader.addSplits(Collections.singletonList(new FlinkTieringSplit(split))); // send TieringReachMaxDurationEvent TieringReachMaxDurationEvent event = new TieringReachMaxDurationEvent(tableId); @@ -113,7 +115,7 @@ void testHandleTieringReachMaxDurationEvent() throws Exception { // tiering won't be finished if no tiering reach max duration logic 100L); - reader.addSplits(Collections.singletonList(split)); + reader.addSplits(Collections.singletonList(new FlinkTieringSplit(split))); // wait to run one round of tiering to do some tiering FutureCompletingBlockingQueue< @@ -157,7 +159,7 @@ void testHandleTieringReachMaxDurationEvent() throws Exception { EARLIEST_OFFSET, 100L); split.skipCurrentRound(); - reader.addSplits(Collections.singletonList(split)); + reader.addSplits(Collections.singletonList(new FlinkTieringSplit(split))); // should skip tiering for this split retry( Duration.ofMinutes(1), diff --git a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/TieringSplitReaderTest.java b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/TieringSplitReaderTest.java index 5e40640511c..3beefbc1987 100644 --- a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/TieringSplitReaderTest.java +++ b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/TieringSplitReaderTest.java @@ -23,14 +23,15 @@ import org.apache.fluss.client.table.writer.AppendWriter; import org.apache.fluss.client.table.writer.TableWriter; import org.apache.fluss.client.table.writer.UpsertWriter; +import org.apache.fluss.client.tiering.TableBucketWriteResult; +import org.apache.fluss.client.tiering.TieringLogSplit; +import org.apache.fluss.client.tiering.TieringSnapshotSplit; import org.apache.fluss.client.write.HashBucketAssigner; import org.apache.fluss.config.ConfigOptions; import org.apache.fluss.flink.tiering.TestingLakeTieringFactory; import org.apache.fluss.flink.tiering.TestingWriteResult; import org.apache.fluss.flink.tiering.source.metrics.TieringMetrics; -import org.apache.fluss.flink.tiering.source.split.TieringLogSplit; -import org.apache.fluss.flink.tiering.source.split.TieringSnapshotSplit; -import org.apache.fluss.flink.tiering.source.split.TieringSplit; +import org.apache.fluss.flink.tiering.source.split.FlinkTieringSplit; import org.apache.fluss.flink.utils.FlinkTestBase; import org.apache.fluss.lake.writer.LakeWriter; import org.apache.fluss.lake.writer.WriterInitContext; @@ -81,7 +82,7 @@ void testTieringTable() throws Exception { TieringSplitReader tieringSplitReader = createTieringReader(connection, lakeTieringFactory)) { // test empty splits - SplitsAddition splitsAddition = + SplitsAddition splitsAddition = new SplitsAddition<>( Arrays.asList( createLogSplit(tablePath, tableId, 0, EARLIEST_OFFSET, 0), @@ -148,14 +149,14 @@ void testTieringTable() throws Exception { Map> secondRows = putRows(tableId, tablePath, 10); Map expectedRowCount = new HashMap<>(); Set expectFinishTieringSplits = new HashSet<>(); - List logSplits = new ArrayList<>(); + List logSplits = new ArrayList<>(); for (int bucket = 0; bucket < 3; bucket++) { TableBucket tableBucket = new TableBucket(tableId, bucket); long startingOffset = firstRows.get(tableBucket).size(); // -U, +U long stoppingOffset = startingOffset + secondRows.get(tableBucket).size() * 2L; expectedRowCount.put(tableBucket, secondRows.get(tableBucket).size() * 2); - TieringLogSplit tieringLogSplit = + FlinkTieringSplit tieringLogSplit = createLogSplit(tablePath, tableId, bucket, startingOffset, stoppingOffset); logSplits.add(tieringLogSplit); expectFinishTieringSplits.add(tieringLogSplit.splitId()); @@ -189,14 +190,14 @@ void testTieringMixTables() throws Exception { FLUSS_CLUSTER_EXTENSION.triggerAndWaitSnapshot(tablePath1); // first add snapshot split of bucket 0, bucket 1 of table id 0 - SplitsAddition splitsAddition = + SplitsAddition splitsAddition = new SplitsAddition<>( Arrays.asList( createSnapshotSplit(tablePath0, tableId0, 0, 0), createSnapshotSplit(tablePath0, tableId0, 1, 0))); Set table0Splits = splitsAddition.splits().stream() - .map(TieringSplit::splitId) + .map(FlinkTieringSplit::splitId) .collect(Collectors.toSet()); tieringSplitReader.handleSplitsChanges(splitsAddition); @@ -220,7 +221,7 @@ void testTieringMixTables() throws Exception { tieringSplitReader.handleSplitsChanges(splitsAddition); Set table1Splits = splitsAddition.splits().stream() - .map(TieringSplit::splitId) + .map(FlinkTieringSplit::splitId) .collect(Collectors.toSet()); // add bucket2 of table id 0 @@ -235,7 +236,7 @@ void testTieringMixTables() throws Exception { table0Rows.get(new TableBucket(tableId0, 2)).size()))); table0Splits.addAll( splitsAddition.splits().stream() - .map(TieringSplit::splitId) + .map(FlinkTieringSplit::splitId) .collect(Collectors.toSet())); tieringSplitReader.handleSplitsChanges(splitsAddition); @@ -283,7 +284,7 @@ void testTieringMixTables() throws Exception { table2Rows.get(new TableBucket(tableId2, 2)).size()))); Set table2Splits = splitsAddition.splits().stream() - .map(TieringSplit::splitId) + .map(FlinkTieringSplit::splitId) .collect(Collectors.toSet()); tieringSplitReader.handleSplitsChanges(splitsAddition); Map expectedRowCount = @@ -332,7 +333,8 @@ connection, new ThrowOnEmptyCompleteLakeTieringFactory())) { // The custom factory fails if complete() is called on a writer that never received any // record, which captures the regression this test covers. tieringSplitReader.handleSplitsChanges( - new SplitsAddition(Collections.singletonList(tieringLogSplit))); + new SplitsAddition<>( + Collections.singletonList(new FlinkTieringSplit(tieringLogSplit)))); RecordsWithSplitIds> result = tieringSplitReader.fetch(); @@ -380,8 +382,14 @@ void testLakeWriterClosedWhenCompleteFails() throws Exception { tieringSplitReader.handleSplitsChanges( new SplitsAddition<>( Collections.singletonList( - new TieringLogSplit( - tablePath, tableBucket, null, EARLIEST_OFFSET, 2, 1)))); + new FlinkTieringSplit( + new TieringLogSplit( + tablePath, + tableBucket, + null, + EARLIEST_OFFSET, + 2, + 1))))); // the injected complete() failure should propagate out of fetch() assertThatThrownBy( @@ -416,7 +424,7 @@ void testCloseClosesAllInFlightLakeWriters() throws Exception { // add log splits with a stopping offset beyond the log end offset, so the splits // never finish and the lake writers stay in-flight - List logSplits = new ArrayList<>(); + List logSplits = new ArrayList<>(); for (Map.Entry> entry : rows.entrySet()) { logSplits.add( createLogSplit( @@ -474,7 +482,7 @@ void testTieringFirstRowMergeEngineFinishes() throws Exception { } // Build log splits whose stoppingOffset equals the leader's current logEndOffset. - List logSplits = new ArrayList<>(); + List logSplits = new ArrayList<>(); Set splitIds = new HashSet<>(); long totalLogEndOffset = 0L; for (int bucket = 0; bucket < DEFAULT_BUCKET_NUM; bucket++) { @@ -485,7 +493,7 @@ void testTieringFirstRowMergeEngineFinishes() throws Exception { if (stoppingOffset <= 0) { continue; } - TieringLogSplit split = + FlinkTieringSplit split = createLogSplit(tablePath, tableId, bucket, EARLIEST_OFFSET, stoppingOffset); logSplits.add(split); splitIds.add(split.splitId()); @@ -603,20 +611,23 @@ private void verifyTieringRows( } } - private TieringLogSplit createLogSplit( + private FlinkTieringSplit createLogSplit( TablePath tablePath, long tableId, int bucket, long startingOffset, long stoppingOffset) { TableBucket tableBucket = new TableBucket(tableId, bucket); - return new TieringLogSplit(tablePath, tableBucket, null, startingOffset, stoppingOffset, 3); + return new FlinkTieringSplit( + new TieringLogSplit( + tablePath, tableBucket, null, startingOffset, stoppingOffset, 3)); } - private TieringSnapshotSplit createSnapshotSplit( + private FlinkTieringSplit createSnapshotSplit( TablePath tablePath, long tableId, int bucket, long snapshotId) { TableBucket tableBucket = new TableBucket(tableId, bucket); - return new TieringSnapshotSplit(tablePath, tableBucket, null, snapshotId, 10, 3); + return new FlinkTieringSplit( + new TieringSnapshotSplit(tablePath, tableBucket, null, snapshotId, 10, 3)); } private Map> putRows(long tableId, TablePath tablePath, int rows) diff --git a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/TieringWriterInitContextTest.java b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/TieringWriterInitContextTest.java index 8002f476c22..7dcdbc20f51 100644 --- a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/TieringWriterInitContextTest.java +++ b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/TieringWriterInitContextTest.java @@ -17,6 +17,8 @@ package org.apache.fluss.flink.tiering.source; +import org.apache.fluss.client.tiering.TieringWriterInitContext; + import org.junit.jupiter.api.Test; import static org.assertj.core.api.Assertions.assertThat; diff --git a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/enumerator/TieringSourceEnumeratorTest.java b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/enumerator/TieringSourceEnumeratorTest.java index 8d1f62eb348..0358c15cc7a 100644 --- a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/enumerator/TieringSourceEnumeratorTest.java +++ b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/enumerator/TieringSourceEnumeratorTest.java @@ -17,6 +17,10 @@ package org.apache.fluss.flink.tiering.source.enumerator; +import org.apache.fluss.client.tiering.TieringLogSplit; +import org.apache.fluss.client.tiering.TieringSnapshotSplit; +import org.apache.fluss.client.tiering.TieringSplit; +import org.apache.fluss.client.tiering.TieringSplitGenerator; import org.apache.fluss.config.ConfigOptions; import org.apache.fluss.config.Configuration; import org.apache.fluss.exception.NetworkException; @@ -25,10 +29,7 @@ import org.apache.fluss.flink.tiering.event.FinishedTieringEvent; import org.apache.fluss.flink.tiering.event.TieringReachMaxDurationEvent; import org.apache.fluss.flink.tiering.source.TieringTestBase; -import org.apache.fluss.flink.tiering.source.split.TieringLogSplit; -import org.apache.fluss.flink.tiering.source.split.TieringSnapshotSplit; -import org.apache.fluss.flink.tiering.source.split.TieringSplit; -import org.apache.fluss.flink.tiering.source.split.TieringSplitGenerator; +import org.apache.fluss.flink.tiering.source.split.FlinkTieringSplit; import org.apache.fluss.lake.writer.LakeTieringFactory; import org.apache.fluss.metadata.TableBucket; import org.apache.fluss.metadata.TableChange; @@ -88,7 +89,7 @@ void testPrimaryKeyTableWithNoSnapshotSplits() throws Throwable { int numSubtasks = 4; int expectNumberOfSplits = 3; // test get snapshot split & log split and the assignment - try (FlussMockSplitEnumeratorContext context = + try (FlussMockSplitEnumeratorContext context = new FlussMockSplitEnumeratorContext<>(numSubtasks)) { TieringSourceEnumerator enumerator = createTieringSourceEnumerator(flussConf, context); @@ -101,7 +102,7 @@ void testPrimaryKeyTableWithNoSnapshotSplits() throws Throwable { // try to assign splits context.runPeriodicCallable(0); - List actualAssignment = new ArrayList<>(); + List actualAssignment = new ArrayList<>(); context.getSplitsAssignmentSequence() .forEach(a -> a.assignment().values().forEach(actualAssignment::addAll)); @@ -145,7 +146,7 @@ void testPrimaryKeyTableWithNoSnapshotSplits() throws Throwable { + bucketOffsetOfSecondWrite.get(tableBucket), expectNumberOfSplits)); } - List actualLogAssignment = new ArrayList<>(); + List actualLogAssignment = new ArrayList<>(); context.getSplitsAssignmentSequence() .forEach(a -> a.assignment().values().forEach(actualLogAssignment::addAll)); assertTieringSplitsMatch(actualLogAssignment, expectedLogAssignment); @@ -165,7 +166,7 @@ void testPrimaryKeyTableWithSnapshotSplits() throws Throwable { int expectNumberOfSplits = 3; // test get snapshot split assignment - try (FlussMockSplitEnumeratorContext context = + try (FlussMockSplitEnumeratorContext context = new FlussMockSplitEnumeratorContext<>(numSubtasks)) { TieringSourceEnumerator enumerator = createTieringSourceEnumerator(flussConf, context); @@ -187,7 +188,7 @@ void testPrimaryKeyTableWithSnapshotSplits() throws Throwable { bucketOffsetOfInitialWrite.get(tableBucket), expectNumberOfSplits)); } - List actualAssignment = new ArrayList<>(); + List actualAssignment = new ArrayList<>(); context.getSplitsAssignmentSequence() .forEach(a -> a.assignment().values().forEach(actualAssignment::addAll)); assertTieringSplitsMatch(actualAssignment, expectedSnapshotAssignment); @@ -231,7 +232,7 @@ void testPrimaryKeyTableWithSnapshotSplits() throws Throwable { + bucketOffsetOfSecondWrite.get(tableBucket), expectNumberOfSplits)); } - List actualLogAssignment = new ArrayList<>(); + List actualLogAssignment = new ArrayList<>(); context.getSplitsAssignmentSequence() .forEach(a -> a.assignment().values().forEach(actualLogAssignment::addAll)); assertTieringSplitsMatch(actualLogAssignment, expectedLogAssignment); @@ -245,7 +246,7 @@ void testLogTableSplits() throws Throwable { int numSubtasks = 4; int expectNumberOfSplits = 3; // test get log split and the assignment - try (FlussMockSplitEnumeratorContext context = + try (FlussMockSplitEnumeratorContext context = new FlussMockSplitEnumeratorContext<>(numSubtasks)) { TieringSourceEnumerator enumerator = createTieringSourceEnumerator(flussConf, context); @@ -272,7 +273,7 @@ void testLogTableSplits() throws Throwable { bucketOffsetOfFirstWrite.get(bucketId), bucketOffsetOfFirstWrite.size())); } - List actualAssignment = new ArrayList<>(); + List actualAssignment = new ArrayList<>(); context.getSplitsAssignmentSequence() .forEach(a -> a.assignment().values().forEach(actualAssignment::addAll)); @@ -316,7 +317,7 @@ void testLogTableSplits() throws Throwable { + bucketOffsetOfSecondWrite.get(tableBucket), expectNumberOfSplits)); } - List actualLogAssignment = new ArrayList<>(); + List actualLogAssignment = new ArrayList<>(); context.getSplitsAssignmentSequence() .forEach(a -> a.assignment().values().forEach(actualLogAssignment::addAll)); assertTieringSplitsMatch(actualLogAssignment, expectedLogAssignment); @@ -335,7 +336,7 @@ void testPartitionedPrimaryKeyTable() throws Throwable { int numSubtasks = 6; int expectNumberOfSplits = 6; // test get snapshot split assignment - try (FlussMockSplitEnumeratorContext context = + try (FlussMockSplitEnumeratorContext context = new FlussMockSplitEnumeratorContext<>(numSubtasks)) { TieringSourceEnumerator enumerator = createTieringSourceEnumerator(flussConf, context); @@ -348,8 +349,8 @@ void testPartitionedPrimaryKeyTable() throws Throwable { // try to assign splits context.runPeriodicCallable(0); - List actualSnapshotAssignment = new ArrayList<>(); - for (SplitsAssignment splitsAssignment : + List actualSnapshotAssignment = new ArrayList<>(); + for (SplitsAssignment splitsAssignment : context.getSplitsAssignmentSequence()) { splitsAssignment.assignment().values().forEach(actualSnapshotAssignment::addAll); } @@ -413,8 +414,8 @@ void testPartitionedPrimaryKeyTable() throws Throwable { expectNumberOfSplits)); } } - List actualLogAssignment = new ArrayList<>(); - for (SplitsAssignment splitsAssignment : + List actualLogAssignment = new ArrayList<>(); + for (SplitsAssignment splitsAssignment : context.getSplitsAssignmentSequence()) { splitsAssignment.assignment().values().forEach(actualLogAssignment::addAll); } @@ -434,7 +435,7 @@ void testPartitionedLogTableSplits() throws Throwable { int numSubtasks = 6; int expectNumberOfSplits = 6; // test get log split assignment - try (FlussMockSplitEnumeratorContext context = + try (FlussMockSplitEnumeratorContext context = new FlussMockSplitEnumeratorContext<>(numSubtasks)) { TieringSourceEnumerator enumerator = createTieringSourceEnumerator(flussConf, context); @@ -468,8 +469,8 @@ void testPartitionedLogTableSplits() throws Throwable { bucketOffsetOfFirstWrite.size())); } } - List actualAssignment = new ArrayList<>(); - for (SplitsAssignment splitsAssignment : + List actualAssignment = new ArrayList<>(); + for (SplitsAssignment splitsAssignment : context.getSplitsAssignmentSequence()) { splitsAssignment.assignment().values().forEach(actualAssignment::addAll); } @@ -535,8 +536,8 @@ void testPartitionedLogTableSplits() throws Throwable { expectNumberOfSplits)); } } - List actualLogAssignment = new ArrayList<>(); - for (SplitsAssignment splitsAssignment : + List actualLogAssignment = new ArrayList<>(); + for (SplitsAssignment splitsAssignment : context.getSplitsAssignmentSequence()) { splitsAssignment.assignment().values().forEach(actualLogAssignment::addAll); } @@ -553,7 +554,7 @@ void testHandleFailedTieringTableEvent() throws Throwable { Map bucketOffsetOfWrite = appendRow(tablePath, DEFAULT_LOG_TABLE_DESCRIPTOR, 0, 10); // test get log split and the assignment - try (FlussMockSplitEnumeratorContext context = + try (FlussMockSplitEnumeratorContext context = new FlussMockSplitEnumeratorContext<>(numSubtasks)) { TieringSourceEnumerator enumerator = createTieringSourceEnumerator(flussConf, context); @@ -575,7 +576,7 @@ void testHandleFailedTieringTableEvent() throws Throwable { bucketOffsetOfWrite.get(bucketId), expectNumberOfSplits)); } - List actualAssignment = new ArrayList<>(); + List actualAssignment = new ArrayList<>(); context.getSplitsAssignmentSequence() .forEach(a -> a.assignment().values().forEach(actualAssignment::addAll)); @@ -590,7 +591,7 @@ void testHandleFailedTieringTableEvent() throws Throwable { enumerator.handleSplitRequest(subtaskId, "localhost-" + subtaskId); } waitUntilTieringTableSplitAssignmentReady(context, DEFAULT_BUCKET_NUM, 500L); - List actualAssignment1 = new ArrayList<>(); + List actualAssignment1 = new ArrayList<>(); context.getSplitsAssignmentSequence() .forEach(a -> a.assignment().values().forEach(actualAssignment1::addAll)); assertTieringSplitsMatch(actualAssignment1, expectedAssignment); @@ -607,7 +608,7 @@ void testHandleReaderFailOver() throws Throwable { createTable(tablePath2, DEFAULT_LOG_TABLE_DESCRIPTOR); appendRow(tablePath2, DEFAULT_LOG_TABLE_DESCRIPTOR, 0, 10); - try (FlussMockSplitEnumeratorContext context = + try (FlussMockSplitEnumeratorContext context = new FlussMockSplitEnumeratorContext<>(3)) { TieringSourceEnumerator enumerator = createTieringSourceEnumerator(flussConf, context); @@ -682,16 +683,18 @@ private static CommitLakeTableSnapshotRequest genCommitLakeTableSnapshotRequest( * the regular {@code equals}. */ private static void assertTieringSplitsMatch( - List actualSplits, List expectedSplits) { + List actualSplits, List expectedSplits) { assertValidTieringRound(actualSplits); List normalizedActualSplits = actualSplits.stream() .map( - split -> - split.copy( - split.getNumberOfSplits(), - TieringSplit.UNKNOWN_SPLIT_INDEX, - TieringSplit.UNKNOWN_TIERING_ROUND_TIMESTAMP)) + flinkSplit -> { + TieringSplit split = flinkSplit.unwrap(); + return split.copy( + split.getNumberOfSplits(), + TieringSplit.UNKNOWN_SPLIT_INDEX, + TieringSplit.UNKNOWN_TIERING_ROUND_TIMESTAMP); + }) .collect(Collectors.toList()); assertThat(normalizedActualSplits).containsExactlyInAnyOrderElementsOf(expectedSplits); } @@ -701,11 +704,11 @@ private static void assertTieringSplitsMatch( * 0..size-1} with exactly one first split, every split reports the round size as its number of * splits, and all splits share the same positive tiering round timestamp. */ - private static void assertValidTieringRound(List tieringSplits) { + private static void assertValidTieringRound(List tieringSplits) { assertThat(tieringSplits).isNotEmpty(); - assertThat(tieringSplits).filteredOn(TieringSplit::isFirstSplit).hasSize(1); + assertThat(tieringSplits).filteredOn(FlinkTieringSplit::isFirstSplit).hasSize(1); assertThat(tieringSplits) - .extracting(TieringSplit::getSplitIndex) + .extracting(FlinkTieringSplit::getSplitIndex) .containsExactlyInAnyOrderElementsOf( IntStream.range(0, tieringSplits.size()) .boxed() @@ -722,7 +725,7 @@ private static void assertValidTieringRound(List tieringSplits) { } private void registerReaderAndHandleSplitRequests( - FlussMockSplitEnumeratorContext context, + FlussMockSplitEnumeratorContext context, TieringSourceEnumerator enumerator, int numSubtasks, int attemptNumber) { @@ -733,7 +736,7 @@ private void registerReaderAndHandleSplitRequests( } private void registerSingleReaderAndHandleSplitRequests( - FlussMockSplitEnumeratorContext context, + FlussMockSplitEnumeratorContext context, TieringSourceEnumerator enumerator, int subtaskId, int attemptNumber) { @@ -743,7 +746,7 @@ private void registerSingleReaderAndHandleSplitRequests( } private void waitUntilTieringTableSplitAssignmentReady( - FlussMockSplitEnumeratorContext context, + FlussMockSplitEnumeratorContext context, int expectedSplitsNum, long sleepMs) throws Throwable { @@ -758,15 +761,15 @@ private void waitUntilTieringTableSplitAssignmentReady( } private void verifyTieringSplitAssignment( - FlussMockSplitEnumeratorContext context, + FlussMockSplitEnumeratorContext context, int expectedSplitSize, TablePath expectedTablePath) throws Throwable { waitUntilTieringTableSplitAssignmentReady(context, expectedSplitSize, 200); - List> actualAssignment = + List> actualAssignment = context.getSplitsAssignmentSequence(); - List allTieringSplits = + List allTieringSplits = actualAssignment.stream() .flatMap(assignments -> assignments.assignment().values().stream()) .flatMap(List::stream) @@ -777,13 +780,13 @@ private void verifyTieringSplitAssignment( } private TieringSourceEnumerator createTieringSourceEnumerator( - Configuration flussConf, MockSplitEnumeratorContext context) { + Configuration flussConf, MockSplitEnumeratorContext context) { return createTieringSourceEnumerator(flussConf, context, new TestingLakeTieringFactory()); } private TieringSourceEnumerator createTieringSourceEnumerator( Configuration flussConf, - MockSplitEnumeratorContext context, + MockSplitEnumeratorContext context, LakeTieringFactory lakeTieringFactory) { return new TieringSourceEnumerator(flussConf, context, lakeTieringFactory, 500); } @@ -802,7 +805,7 @@ public void validateTable(TableInfo tableInfo) throws IOException { } }; - try (FlussMockSplitEnumeratorContext context = + try (FlussMockSplitEnumeratorContext context = new FlussMockSplitEnumeratorContext<>(1); TieringSourceEnumerator enumerator = createTieringSourceEnumerator(flussConf, context, lakeTieringFactory)) { @@ -817,7 +820,7 @@ public void validateTable(TableInfo tableInfo) throws IOException { @Test void testNetworkErrorInHeartbeatTriggersFailover() throws Exception { - try (FlussMockSplitEnumeratorContext context = + try (FlussMockSplitEnumeratorContext context = new FlussMockSplitEnumeratorContext<>(1)) { TieringSourceEnumerator enumerator = createTieringSourceEnumerator(flussConf, context); FlinkRuntimeException networkError = @@ -835,7 +838,7 @@ void testTableReachMaxTieringDuration() throws Throwable { long tableId = createTable(tablePath, DEFAULT_LOG_TABLE_DESCRIPTOR); int numSubtasks = 2; - try (FlussMockSplitEnumeratorContext context = + try (FlussMockSplitEnumeratorContext context = new FlussMockSplitEnumeratorContext<>(numSubtasks); TieringSourceEnumerator enumerator = createTieringSourceEnumerator(flussConf, context)) { @@ -879,7 +882,7 @@ void testTableReachMaxTieringDuration() throws Throwable { // the split should be marked as skipCurrentRound waitUntilTieringTableSplitAssignmentReady(context, 1, 100L); - List assignedSplits = new ArrayList<>(); + List assignedSplits = new ArrayList<>(); context.getSplitsAssignmentSequence() .forEach(a -> a.assignment().values().forEach(assignedSplits::addAll)); assertThat(assignedSplits).hasSize(1); @@ -925,7 +928,7 @@ void testTableReachMaxTieringDuration() throws Throwable { // Wait for the table to be assigned again waitUntilTieringTableSplitAssignmentReady(context, 2, 500L); - List reassignedSplits = new ArrayList<>(); + List reassignedSplits = new ArrayList<>(); context.getSplitsAssignmentSequence() .forEach(a -> a.assignment().values().forEach(reassignedSplits::addAll)); assertThat(reassignedSplits).hasSize(2); diff --git a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/split/TieringSplitSerializerTest.java b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/split/TieringSplitSerializerTest.java index f9c0913b90c..d87c13dc0c0 100644 --- a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/split/TieringSplitSerializerTest.java +++ b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/split/TieringSplitSerializerTest.java @@ -17,6 +17,9 @@ package org.apache.fluss.flink.tiering.source.split; +import org.apache.fluss.client.tiering.TieringLogSplit; +import org.apache.fluss.client.tiering.TieringSnapshotSplit; +import org.apache.fluss.client.tiering.TieringSplit; import org.apache.fluss.metadata.TableBucket; import org.apache.fluss.metadata.TablePath; @@ -24,7 +27,10 @@ import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.ValueSource; +import java.io.IOException; + import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; /** * Unit tests for serialization and deserialization of {@link TieringSnapshotSplit} and {@link @@ -48,9 +54,10 @@ void testTieringSnapshotSplitSerde(Boolean isPartitionedTable) throws Exception TieringSnapshotSplit tieringSplit = new TieringSnapshotSplit(path, bucket, partitionName, 0L, 200L, 10); - byte[] serialized = serializer.serialize(tieringSplit); - TieringSnapshotSplit deserializedSplit = - (TieringSnapshotSplit) serializer.deserialize(serializer.getVersion(), serialized); + byte[] serialized = serializer.serialize(new FlinkTieringSplit(tieringSplit)); + FlinkTieringSplit deserializedFlinkSplit = + serializer.deserialize(serializer.getVersion(), serialized); + TieringSnapshotSplit deserializedSplit = deserializedFlinkSplit.asTieringSnapshotSplit(); assertThat(deserializedSplit).isEqualTo(tieringSplit); } @@ -84,12 +91,27 @@ void testTieringLogSplitSerde(Boolean isPartitionedTable) throws Exception { TieringLogSplit tieringSplit = new TieringLogSplit(path, bucket, partitionName, 100, 200, 40); - byte[] serialized = serializer.serialize(tieringSplit); - TieringLogSplit deserializedSplit = - (TieringLogSplit) serializer.deserialize(serializer.getVersion(), serialized); + byte[] serialized = serializer.serialize(new FlinkTieringSplit(tieringSplit)); + FlinkTieringSplit deserializedFlinkSplit = + serializer.deserialize(serializer.getVersion(), serialized); + TieringLogSplit deserializedSplit = deserializedFlinkSplit.asTieringLogSplit(); assertThat(deserializedSplit).isEqualTo(tieringSplit); } + @Test + void testDeserializeUnknownSplitKindFailsFast() throws Exception { + TieringLogSplit tieringSplit = + new TieringLogSplit(tablePath, tableBucket, null, 100, 200, 40); + byte[] serialized = serializer.serialize(new FlinkTieringSplit(tieringSplit)); + // the split kind is written as the first byte, tamper it with an unknown kind so that + // deserialization must fail fast instead of silently falling back to a log split + serialized[0] = 99; + + assertThatThrownBy(() -> serializer.deserialize(serializer.getVersion(), serialized)) + .isInstanceOf(IOException.class) + .hasMessageContaining("Unknown split kind 99"); + } + @ParameterizedTest @ValueSource(booleans = {true, false}) void testTieringLogSplitStringExpression(Boolean isPartitionedTable) throws Exception { @@ -114,17 +136,20 @@ void testSkipCurrentRoundSerde() throws Exception { // Test TieringSnapshotSplit with skipCurrentRound set at creation TieringSnapshotSplit snapshotSplitWithSkipCurrentRound = new TieringSnapshotSplit(tablePath, tableBucket, null, 0L, 200L, 10, true); - byte[] serialized = serializer.serialize(snapshotSplitWithSkipCurrentRound); + byte[] serialized = + serializer.serialize(new FlinkTieringSplit(snapshotSplitWithSkipCurrentRound)); TieringSnapshotSplit deserializedSnapshotSplit = - (TieringSnapshotSplit) serializer.deserialize(serializer.getVersion(), serialized); + serializer + .deserialize(serializer.getVersion(), serialized) + .asTieringSnapshotSplit(); assertThat(deserializedSnapshotSplit).isEqualTo(snapshotSplitWithSkipCurrentRound); // Test TieringLogSplit with skipCurrentRound set at creation TieringLogSplit logSplitWithSkipCurrentRound = new TieringLogSplit(tablePath, tableBucket, null, 100, 200, 40, true); - serialized = serializer.serialize(logSplitWithSkipCurrentRound); + serialized = serializer.serialize(new FlinkTieringSplit(logSplitWithSkipCurrentRound)); TieringLogSplit deserializedLogSplit = - (TieringLogSplit) serializer.deserialize(serializer.getVersion(), serialized); + serializer.deserialize(serializer.getVersion(), serialized).asTieringLogSplit(); assertThat(deserializedLogSplit).isEqualTo(logSplitWithSkipCurrentRound); // Test TieringSnapshotSplit with skipCurrentRound set after creation @@ -134,9 +159,11 @@ void testSkipCurrentRoundSerde() throws Exception { snapshotSplit.skipCurrentRound(); assertThat(snapshotSplit.shouldSkipCurrentRound()).isTrue(); - serialized = serializer.serialize(snapshotSplit); + serialized = serializer.serialize(new FlinkTieringSplit(snapshotSplit)); deserializedSnapshotSplit = - (TieringSnapshotSplit) serializer.deserialize(serializer.getVersion(), serialized); + serializer + .deserialize(serializer.getVersion(), serialized) + .asTieringSnapshotSplit(); assertThat(deserializedSnapshotSplit).isEqualTo(snapshotSplit); // Test TieringLogSplit with skipCurrentRound set after creation @@ -146,9 +173,9 @@ void testSkipCurrentRoundSerde() throws Exception { logSplit.skipCurrentRound(); assertThat(logSplit.shouldSkipCurrentRound()).isTrue(); - serialized = serializer.serialize(logSplit); + serialized = serializer.serialize(new FlinkTieringSplit(logSplit)); deserializedLogSplit = - (TieringLogSplit) serializer.deserialize(serializer.getVersion(), serialized); + serializer.deserialize(serializer.getVersion(), serialized).asTieringLogSplit(); assertThat(deserializedLogSplit).isEqualTo(logSplit); } @@ -156,20 +183,96 @@ void testSkipCurrentRoundSerde() throws Exception { void testTieringRoundTimestampSerde() throws Exception { TieringSnapshotSplit snapshotSplit = new TieringSnapshotSplit(tablePath, tableBucket, null, 0L, 200L, 10, 0, 1000L); - byte[] serialized = serializer.serialize(snapshotSplit); + byte[] serialized = serializer.serialize(new FlinkTieringSplit(snapshotSplit)); TieringSnapshotSplit deserializedSnapshotSplit = - (TieringSnapshotSplit) serializer.deserialize(serializer.getVersion(), serialized); + serializer + .deserialize(serializer.getVersion(), serialized) + .asTieringSnapshotSplit(); assertThat(deserializedSnapshotSplit.getSplitIndex()).isZero(); assertThat(deserializedSnapshotSplit.isFirstSplit()).isTrue(); assertThat(deserializedSnapshotSplit.getTieringRoundTimestamp()).isEqualTo(1000L); TieringLogSplit logSplit = new TieringLogSplit(tablePath, tableBucket, null, 100, 200, 40, 2, 2000L); - serialized = serializer.serialize(logSplit); + serialized = serializer.serialize(new FlinkTieringSplit(logSplit)); TieringLogSplit deserializedLogSplit = - (TieringLogSplit) serializer.deserialize(serializer.getVersion(), serialized); + serializer.deserialize(serializer.getVersion(), serialized).asTieringLogSplit(); assertThat(deserializedLogSplit.getSplitIndex()).isEqualTo(2); assertThat(deserializedLogSplit.isFirstSplit()).isFalse(); assertThat(deserializedLogSplit.getTieringRoundTimestamp()).isEqualTo(2000L); } + + @Test + void testFlinkTieringSplitDelegation() { + TieringLogSplit logSplit = new TieringLogSplit(tablePath, tableBucket, null, 10, 200, 5); + FlinkTieringSplit flinkSplit = new FlinkTieringSplit(logSplit); + + assertThat(flinkSplit.splitId()).isNotNull(); + assertThat(flinkSplit.unwrap()).isSameAs(logSplit); + assertThat(flinkSplit.isTieringLogSplit()).isTrue(); + assertThat(flinkSplit.isTieringSnapshotSplit()).isFalse(); + assertThat(flinkSplit.asTieringLogSplit()).isSameAs(logSplit); + assertThat(flinkSplit.getTablePath()).isEqualTo(tablePath); + assertThat(flinkSplit.getTableBucket()).isEqualTo(tableBucket); + assertThat(flinkSplit.getPartitionName()).isNull(); + assertThat(flinkSplit.getNumberOfSplits()).isEqualTo(5); + assertThat(flinkSplit.splitKind()).isEqualTo(TieringSplit.TIERING_LOG_SPLIT_FLAG); + + TieringSnapshotSplit snapshotSplit = + new TieringSnapshotSplit( + tablePath, tableBucket, null, 1L, 100L, 3, false, 0, 5000L); + FlinkTieringSplit flinkSnapshotSplit = new FlinkTieringSplit(snapshotSplit); + assertThat(flinkSnapshotSplit.isTieringSnapshotSplit()).isTrue(); + assertThat(flinkSnapshotSplit.asTieringSnapshotSplit()).isSameAs(snapshotSplit); + assertThat(flinkSnapshotSplit.getSplitIndex()).isZero(); + assertThat(flinkSnapshotSplit.isFirstSplit()).isTrue(); + assertThat(flinkSnapshotSplit.getTieringRoundTimestamp()).isEqualTo(5000L); + } + + @Test + void testFlinkTieringSplitCopyAndSkip() { + TieringLogSplit logSplit = new TieringLogSplit(tablePath, tableBucket, null, 0, 100, 3); + FlinkTieringSplit flinkSplit = new FlinkTieringSplit(logSplit); + + assertThat(flinkSplit.shouldSkipCurrentRound()).isFalse(); + flinkSplit.skipCurrentRound(); + assertThat(flinkSplit.shouldSkipCurrentRound()).isTrue(); + + FlinkTieringSplit copied = flinkSplit.copy(5); + assertThat(copied.getNumberOfSplits()).isEqualTo(5); + assertThat(copied.splitId()).isEqualTo(flinkSplit.splitId()); + + FlinkTieringSplit copiedWithMeta = flinkSplit.copy(10, 2, 9999L); + assertThat(copiedWithMeta.getNumberOfSplits()).isEqualTo(10); + assertThat(copiedWithMeta.getSplitIndex()).isEqualTo(2); + assertThat(copiedWithMeta.isFirstSplit()).isFalse(); + assertThat(copiedWithMeta.getTieringRoundTimestamp()).isEqualTo(9999L); + } + + @Test + void testFlinkTieringSplitEqualsHashCodeToString() { + TieringLogSplit logSplit = new TieringLogSplit(tablePath, tableBucket, null, 0, 100, 3); + FlinkTieringSplit split1 = new FlinkTieringSplit(logSplit); + FlinkTieringSplit split2 = new FlinkTieringSplit(logSplit); + + assertThat(split1).isEqualTo(split2); + assertThat(split1.hashCode()).isEqualTo(split2.hashCode()); + assertThat(split1).isNotEqualTo(null); + assertThat(split1).isNotEqualTo("not a split"); + + String str = split1.toString(); + assertThat(str).startsWith("FlinkTieringSplit{"); + assertThat(str).contains("TieringLogSplit"); + } + + @Test + void testFlinkTieringSplitPartitioned() { + FlinkTieringSplit flinkSplit = + new FlinkTieringSplit( + new TieringLogSplit( + partitionedTablePath, partitionedTableBucket, "p1", 0, 50, 2)); + + assertThat(flinkSplit.getPartitionName()).isEqualTo("p1"); + assertThat(flinkSplit.getTableBucket().getPartitionId()).isEqualTo(100L); + } } diff --git a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/state/TieringSplitStateTest.java b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/state/TieringSplitStateTest.java index ebc08553e6a..68b86206924 100644 --- a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/state/TieringSplitStateTest.java +++ b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/state/TieringSplitStateTest.java @@ -17,8 +17,9 @@ package org.apache.fluss.flink.tiering.source.state; -import org.apache.fluss.flink.tiering.source.split.TieringLogSplit; -import org.apache.fluss.flink.tiering.source.split.TieringSnapshotSplit; +import org.apache.fluss.client.tiering.TieringLogSplit; +import org.apache.fluss.client.tiering.TieringSnapshotSplit; +import org.apache.fluss.flink.tiering.source.split.FlinkTieringSplit; import org.apache.fluss.metadata.TableBucket; import org.apache.fluss.metadata.TablePath; @@ -39,9 +40,10 @@ void testTieringSnapshotSplit() { new TieringSnapshotSplit( tablePath, tableBucket, "partition1", 0L, 200L, 10, 0, 1000L); tieringSnapshotSplit.skipCurrentRound(); - TieringSplitState tieringSnapshotSplitState = new TieringSplitState(tieringSnapshotSplit); - TieringSnapshotSplit restoredSplit = - (TieringSnapshotSplit) tieringSnapshotSplitState.toSourceSplit(); + FlinkTieringSplit flinkSplit = new FlinkTieringSplit(tieringSnapshotSplit); + TieringSplitState tieringSnapshotSplitState = new TieringSplitState(flinkSplit); + FlinkTieringSplit restoredFlinkSplit = tieringSnapshotSplitState.toSourceSplit(); + TieringSnapshotSplit restoredSplit = restoredFlinkSplit.asTieringSnapshotSplit(); assertThat(restoredSplit).isEqualTo(tieringSnapshotSplit); assertThat(restoredSplit.shouldSkipCurrentRound()).isTrue(); assertThat(restoredSplit.getSplitIndex()).isZero(); @@ -57,8 +59,10 @@ void testTieringLogSplit() { TieringLogSplit tieringLogSplit = new TieringLogSplit(tablePath, tableBucket, "partition1", 100L, 200L, 20, 1, 2000L); tieringLogSplit.skipCurrentRound(); - TieringSplitState tieringLogSplitState = new TieringSplitState(tieringLogSplit); - TieringLogSplit restoredSplit = (TieringLogSplit) tieringLogSplitState.toSourceSplit(); + FlinkTieringSplit flinkSplit = new FlinkTieringSplit(tieringLogSplit); + TieringSplitState tieringLogSplitState = new TieringSplitState(flinkSplit); + FlinkTieringSplit restoredFlinkSplit = tieringLogSplitState.toSourceSplit(); + TieringLogSplit restoredSplit = restoredFlinkSplit.asTieringLogSplit(); assertThat(restoredSplit).isEqualTo(tieringLogSplit); assertThat(restoredSplit.shouldSkipCurrentRound()).isTrue(); assertThat(restoredSplit.getSplitIndex()).isEqualTo(1); diff --git a/fluss-lake/fluss-lake-paimon/src/test/java/org/apache/fluss/lake/paimon/utils/DvTableReadableSnapshotRetrieverTest.java b/fluss-lake/fluss-lake-paimon/src/test/java/org/apache/fluss/lake/paimon/utils/DvTableReadableSnapshotRetrieverTest.java index 0ce3d5b8efd..15926a0f34b 100644 --- a/fluss-lake/fluss-lake-paimon/src/test/java/org/apache/fluss/lake/paimon/utils/DvTableReadableSnapshotRetrieverTest.java +++ b/fluss-lake/fluss-lake-paimon/src/test/java/org/apache/fluss/lake/paimon/utils/DvTableReadableSnapshotRetrieverTest.java @@ -21,10 +21,10 @@ import org.apache.fluss.client.Connection; import org.apache.fluss.client.ConnectionFactory; import org.apache.fluss.client.admin.Admin; +import org.apache.fluss.client.tiering.FlussTableLakeSnapshotCommitter; import org.apache.fluss.config.ConfigOptions; import org.apache.fluss.config.Configuration; import org.apache.fluss.exception.FlussRuntimeException; -import org.apache.fluss.flink.tiering.committer.FlussTableLakeSnapshotCommitter; import org.apache.fluss.lake.committer.LakeCommitResult; import org.apache.fluss.metadata.PartitionInfo; import org.apache.fluss.metadata.PartitionSpec;