From 519d5b0aaced9672c64ebbd82652b7a7a03b9bf7 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=BE=8A=E5=B7=9D?= Date: Wed, 1 Jul 2026 23:22:45 +0800 Subject: [PATCH 1/5] refactor TieringSplit --- .../client/tiering}/TieringLogSplit.java | 2 +- .../client/tiering}/TieringSnapshotSplit.java | 2 +- .../fluss/client/tiering}/TieringSplit.java | 12 +- .../flink/tiering/source/TieringSource.java | 18 ++- .../source/TieringSourceFetcherManager.java | 10 +- .../tiering/source/TieringSourceReader.java | 10 +- .../tiering/source/TieringSplitReader.java | 14 +- .../enumerator/TieringSourceEnumerator.java | 36 +++-- .../source/split/FlinkTieringSplit.java | 152 ++++++++++++++++++ .../source/split/TieringSplitGenerator.java | 3 + .../source/split/TieringSplitSerializer.java | 56 ++++--- .../source/state/TieringSplitState.java | 64 ++++---- .../source/TieringSourceReaderTest.java | 9 +- .../source/TieringSplitReaderTest.java | 40 ++--- .../TieringSourceEnumeratorTest.java | 99 ++++++------ .../split/TieringSplitSerializerTest.java | 47 +++--- .../source/state/TieringSplitStateTest.java | 18 ++- 17 files changed, 395 insertions(+), 197 deletions(-) rename {fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split => fluss-client/src/main/java/org/apache/fluss/client/tiering}/TieringLogSplit.java (99%) rename {fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split => fluss-client/src/main/java/org/apache/fluss/client/tiering}/TieringSnapshotSplit.java (99%) rename {fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split => fluss-client/src/main/java/org/apache/fluss/client/tiering}/TieringSplit.java (96%) create mode 100644 fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split/FlinkTieringSplit.java 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/TieringSource.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringSource.java index 696d4722b8a..fe594d8dc0c 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 @@ -21,7 +21,7 @@ import org.apache.fluss.client.ConnectionFactory; 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 +56,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 +87,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 +103,7 @@ public SplitEnumerator restoreEnumer } @Override - public SimpleVersionedSerializer getSplitSerializer() { + public SimpleVersionedSerializer getSplitSerializer() { return TieringSplitSerializer.INSTANCE; } @@ -112,7 +114,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..fb19afbe24e 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 @@ -19,7 +19,7 @@ package org.apache.fluss.flink.tiering.source; 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 +40,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 +64,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 +75,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..c90da64bc27 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 @@ -23,7 +23,7 @@ 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 +48,7 @@ public final class TieringSourceReader extends SingleThreadMultiplexSourceReaderBaseAdapter< TableBucketWriteResult, TableBucketWriteResult, - TieringSplit, + FlinkTieringSplit, TieringSplitState> { private static final Logger LOG = LoggerFactory.getLogger(TieringSourceReader.class); @@ -144,13 +144,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 +161,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..5b734b67d0f 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,13 @@ 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.TieringLogSplit; +import org.apache.fluss.client.tiering.TieringSnapshotSplit; +import org.apache.fluss.client.tiering.TieringSplit; 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 +75,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 +232,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..067c23d2ee9 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,12 +22,13 @@ 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.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.FlinkTieringSplit; import org.apache.fluss.flink.tiering.source.split.TieringSplitGenerator; import org.apache.fluss.flink.tiering.source.state.TieringSourceEnumeratorState; import org.apache.fluss.lake.committer.TieringStats; @@ -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/TieringSplitGenerator.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split/TieringSplitGenerator.java index a4b2638f309..cd37340614d 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split/TieringSplitGenerator.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split/TieringSplitGenerator.java @@ -22,6 +22,9 @@ 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.client.tiering.TieringLogSplit; +import org.apache.fluss.client.tiering.TieringSnapshotSplit; +import org.apache.fluss.client.tiering.TieringSplit; import org.apache.fluss.exception.LakeTableSnapshotNotExistException; import org.apache.fluss.metadata.PartitionInfo; import org.apache.fluss.metadata.TableBucket; 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..f35a78e5403 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; @@ -29,13 +32,13 @@ import java.io.IOException; /** - * 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(); @@ -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,40 @@ 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); + tieringSplit = + new TieringSnapshotSplit( + tablePath, + tableBucket, + partitionName, + snapshotId, + logOffsetOfSnapshot, + numberOfSplits, + skipCurrentRound, + splitIndex, + tieringRoundTimestamp); } else { // 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); } + 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/source/TieringSourceReaderTest.java b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/TieringSourceReaderTest.java index 9e9de2c7929..369722ab8e6 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,13 @@ import org.apache.fluss.client.Connection; import org.apache.fluss.client.ConnectionFactory; +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 +79,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 +114,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 +158,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..3ad1ddf5ef6 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,14 @@ 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.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 +81,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 +148,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 +189,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 +220,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 +235,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 +283,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 +332,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(); @@ -474,7 +475,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 +486,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 +604,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/enumerator/TieringSourceEnumeratorTest.java b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/source/enumerator/TieringSourceEnumeratorTest.java index 8d1f62eb348..86c01fd74f3 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,9 @@ 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.config.ConfigOptions; import org.apache.fluss.config.Configuration; import org.apache.fluss.exception.NetworkException; @@ -25,9 +28,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.FlinkTieringSplit; import org.apache.fluss.flink.tiering.source.split.TieringSplitGenerator; import org.apache.fluss.lake.writer.LakeTieringFactory; import org.apache.fluss.metadata.TableBucket; @@ -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..8fd43ec51f2 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,8 @@ 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.metadata.TableBucket; import org.apache.fluss.metadata.TablePath; @@ -48,9 +50,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,9 +87,10 @@ 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); } @@ -114,17 +118,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 +141,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 +155,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,18 +165,20 @@ 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); 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); From 8aa6f75ce201c97aa4768b774d77ddcdcbf78cc0 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=BE=8A=E5=B7=9D?= Date: Thu, 2 Jul 2026 14:11:22 +0800 Subject: [PATCH 2/5] add tests --- .../split/TieringSplitSerializerTest.java | 75 +++++++++++++++++++ 1 file changed, 75 insertions(+) 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 8fd43ec51f2..2756209a1d5 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 @@ -19,6 +19,7 @@ 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; @@ -183,4 +184,78 @@ void testTieringRoundTimestampSerde() throws Exception { 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); + } } From 234aa42e0bd56b777fa91b867a2a618eba20c0f3 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=BE=8A=E5=B7=9D?= Date: Fri, 10 Jul 2026 17:22:18 +0800 Subject: [PATCH 3/5] move on --- .../FlussTableLakeSnapshotCommitter.java | 43 +-- .../tiering}/TableBucketWriteResult.java | 2 +- .../client/tiering/TieringCommitResult.java | 52 ++++ .../client/tiering/TieringCommitter.java | 265 ++++++++++++++++++ .../tiering}/TieringCommitterInitContext.java | 2 +- .../tiering}/TieringWriterInitContext.java | 2 +- .../committer/CommittableMessageTypeInfo.java | 2 +- .../committer/TieringCommitOperator.java | 249 ++-------------- .../TieringCommitOperatorFactory.java | 2 +- .../source/TableBucketWriteResultEmitter.java | 1 + .../TableBucketWriteResultSerializer.java | 1 + .../TableBucketWriteResultTypeInfo.java | 1 + .../flink/tiering/source/TieringSource.java | 1 + .../source/TieringSourceFetcherManager.java | 1 + .../tiering/source/TieringSourceReader.java | 1 + .../tiering/source/TieringSplitReader.java | 2 + .../FlussTableLakeSnapshotCommitterTest.java | 1 + .../committer/TieringCommitOperatorTest.java | 3 +- .../TableBucketWriteResultSerializerTest.java | 1 + .../source/TieringSourceReaderTest.java | 1 + .../source/TieringSplitReaderTest.java | 13 +- .../source/TieringWriterInitContextTest.java | 2 + .../DvTableReadableSnapshotRetrieverTest.java | 2 +- 23 files changed, 371 insertions(+), 279 deletions(-) rename {fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/committer => fluss-client/src/main/java/org/apache/fluss/client/tiering}/FlussTableLakeSnapshotCommitter.java (89%) rename {fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source => fluss-client/src/main/java/org/apache/fluss/client/tiering}/TableBucketWriteResult.java (98%) create mode 100644 fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringCommitResult.java create mode 100644 fluss-client/src/main/java/org/apache/fluss/client/tiering/TieringCommitter.java rename {fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/committer => fluss-client/src/main/java/org/apache/fluss/client/tiering}/TieringCommitterInitContext.java (97%) rename {fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source => fluss-client/src/main/java/org/apache/fluss/client/tiering}/TieringWriterInitContext.java (98%) 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: - * - *

    - *
  • Lake table snapshot metadata (snapshot ID, table ID, file paths) - *
  • PbLakeTableSnapshotInfo for metrics reporting (log end offsets and max tiered - * timestamps) - *
- * - * @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/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 fe594d8dc0c..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,6 +19,7 @@ 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.FlinkTieringSplit; 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 fb19afbe24e..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,6 +18,7 @@ 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.FlinkTieringSplit; 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 c90da64bc27..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,6 +20,7 @@ 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; 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 5b734b67d0f..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,9 +25,11 @@ 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; 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 369722ab8e6..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,6 +20,7 @@ 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; 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 3ad1ddf5ef6..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,6 +23,7 @@ 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; @@ -381,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( @@ -417,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( 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-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; From 5762824a645c3a8313f15b815973da59e0f9ebf1 Mon Sep 17 00:00:00 2001 From: "yangchuan.zy" Date: Thu, 30 Jul 2026 19:15:18 +0800 Subject: [PATCH 4/5] [client] Move TieringSplitGenerator to fluss-client tiering package Extract engine-agnostic split generation from fluss-flink so that other engines (e.g. Spark tiering) can reuse it. Replace FlinkRuntimeException with FlussRuntimeException to drop the Flink dependency. --- .../client/tiering}/TieringSplitGenerator.java | 13 +++++-------- .../source/enumerator/TieringSourceEnumerator.java | 2 +- .../enumerator/TieringSourceEnumeratorTest.java | 2 +- 3 files changed, 7 insertions(+), 10 deletions(-) rename {fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/split => fluss-client/src/main/java/org/apache/fluss/client/tiering}/TieringSplitGenerator.java (97%) 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 97% 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 cd37340614d..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,16 +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.client.tiering.TieringLogSplit; -import org.apache.fluss.client.tiering.TieringSnapshotSplit; -import org.apache.fluss.client.tiering.TieringSplit; +import org.apache.fluss.exception.FlussRuntimeException; import org.apache.fluss.exception.LakeTableSnapshotNotExistException; import org.apache.fluss.metadata.PartitionInfo; import org.apache.fluss.metadata.TableBucket; @@ -32,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; @@ -78,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()), @@ -130,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), @@ -166,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/enumerator/TieringSourceEnumerator.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/enumerator/TieringSourceEnumerator.java index 067c23d2ee9..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 @@ -23,13 +23,13 @@ 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.FlinkTieringSplit; -import org.apache.fluss.flink.tiering.source.split.TieringSplitGenerator; import org.apache.fluss.flink.tiering.source.state.TieringSourceEnumeratorState; import org.apache.fluss.lake.committer.TieringStats; import org.apache.fluss.lake.writer.LakeTieringFactory; 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 86c01fd74f3..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 @@ -20,6 +20,7 @@ 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; @@ -29,7 +30,6 @@ import org.apache.fluss.flink.tiering.event.TieringReachMaxDurationEvent; import org.apache.fluss.flink.tiering.source.TieringTestBase; import org.apache.fluss.flink.tiering.source.split.FlinkTieringSplit; -import org.apache.fluss.flink.tiering.source.split.TieringSplitGenerator; import org.apache.fluss.lake.writer.LakeTieringFactory; import org.apache.fluss.metadata.TableBucket; import org.apache.fluss.metadata.TableChange; From 5339daf0f26cba03a9c58b1bd376364405d1c872 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=BE=8A=E5=B7=9D?= Date: Mon, 3 Aug 2026 16:37:52 +0800 Subject: [PATCH 5/5] [flink] Fail fast on unknown split kind in TieringSplitSerializer TieringSplitSerializer#deserialize treated every non-snapshot splitKind as a log split, so an unknown kind was silently deserialized with the wrong field layout instead of failing fast. Match explicitly against TIERING_LOG_SPLIT_FLAG and throw an IOException otherwise. Also drop the serializer's private copies of the split kind flags in favour of the public constants on TieringSplit, so both flags now have exactly one definition and are both actually used. --- .../source/split/TieringSplitSerializer.java | 10 ++++++---- .../split/TieringSplitSerializerTest.java | 17 +++++++++++++++++ 2 files changed, 23 insertions(+), 4 deletions(-) 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 f35a78e5403..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 @@ -31,6 +31,9 @@ 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 FlinkTieringSplit}. * @@ -47,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 @@ -166,7 +166,7 @@ public FlinkTieringSplit deserialize(int version, byte[] serialized) throws IOEx skipCurrentRound, splitIndex, tieringRoundTimestamp); - } else { + } else if (splitKind == TIERING_LOG_SPLIT_FLAG) { // deserialize starting offset long startingOffset = in.readLong(); // deserialize starting offset @@ -182,6 +182,8 @@ public FlinkTieringSplit deserialize(int version, byte[] serialized) throws IOEx skipCurrentRound, splitIndex, tieringRoundTimestamp); + } else { + throw new IOException("Unknown split kind " + splitKind); } return new FlinkTieringSplit(tieringSplit); } 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 2756209a1d5..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 @@ -27,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 @@ -95,6 +98,20 @@ void testTieringLogSplitSerde(Boolean isPartitionedTable) throws Exception { 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 {