Skip to content

Commit f789177

Browse files
committed
complete delay fetch
1 parent 91bea7c commit f789177

7 files changed

Lines changed: 207 additions & 4 deletions

File tree

fluss-rpc/src/main/java/org/apache/fluss/rpc/RpcGatewayService.java

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -85,4 +85,10 @@ public String currentListenerName() {
8585

8686
/** Shutdown the gateway service, release any resources. */
8787
public abstract void shutdown();
88+
89+
/**
90+
* Tries to complete all pending delayed actions. Default no-op for services without delayed
91+
* action queue.
92+
*/
93+
public void tryCompleteActions() {}
8894
}

fluss-rpc/src/main/java/org/apache/fluss/rpc/netty/server/FlussRequestHandler.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -96,6 +96,8 @@ public void processRequest(FlussRequest request) {
9696
} catch (Throwable t) {
9797
LOG.debug("Error while executing RPC {}", api, t);
9898
request.fail(stripException(t, InvocationTargetException.class));
99+
} finally {
100+
service.tryCompleteActions();
99101
}
100102
}
101103
}

fluss-server/src/main/java/org/apache/fluss/server/replica/ReplicaManager.java

Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -99,6 +99,8 @@
9999
import org.apache.fluss.server.metrics.group.BucketMetricGroup;
100100
import org.apache.fluss.server.metrics.group.TableMetricGroup;
101101
import org.apache.fluss.server.metrics.group.TabletServerMetricGroup;
102+
import org.apache.fluss.server.replica.delay.ActionQueue;
103+
import org.apache.fluss.server.replica.delay.DelayedActionQueue;
102104
import org.apache.fluss.server.replica.delay.DelayedFetchLog;
103105
import org.apache.fluss.server.replica.delay.DelayedFetchLog.FetchBucketStatus;
104106
import org.apache.fluss.server.replica.delay.DelayedOperationManager;
@@ -187,6 +189,8 @@ public class ReplicaManager implements ServerReconfigurable {
187189
*/
188190
private final DelayedOperationManager<DelayedFetchLog> delayedFetchLogManager;
189191

192+
private final ActionQueue actionQueue;
193+
190194
private final ReplicaFetcherManager replicaFetcherManager;
191195
// The manager used to manager the replica alter, especially the isr expand and shrink.
192196
private final AdjustIsrManager adjustIsrManager;
@@ -309,6 +313,7 @@ public ReplicaManager(
309313
"delay fetch log",
310314
serverId,
311315
conf.getInt(ConfigOptions.LOG_REPLICA_FETCH_OPERATION_PURGE_NUMBER));
316+
this.actionQueue = new DelayedActionQueue();
312317
this.internalListenerName = conf.get(ConfigOptions.INTERNAL_LISTENER_NAME);
313318

314319
this.replicaFetcherManager =
@@ -636,6 +641,10 @@ public void appendRecordsToLog(
636641
appendToLocalLog(entriesPerBucket, requiredAcks, userContext);
637642
LOG.debug("Append records to local log in {} ms", System.currentTimeMillis() - startTime);
638643

644+
// Enqueue delayed fetch completions — not executed here.
645+
// Framework layer invokes tryCompleteActions() after this method returns.
646+
addCompletePurgatoryAction(appendResult);
647+
639648
// maybe do delay write operation.
640649
maybeAddDelayedWrite(
641650
timeoutMs, requiredAcks, entriesPerBucket.size(), appendResult, responseCallback);
@@ -1783,6 +1792,32 @@ private boolean isNonCriticalFetchError(Errors error) {
17831792
|| error == Errors.UNKNOWN_TABLE_OR_BUCKET_EXCEPTION;
17841793
}
17851794

1795+
/**
1796+
* Adds actions to complete delayed fetch log operations for successfully written buckets.
1797+
*
1798+
* <p>Actions are added to the {@link ActionQueue} rather than executed immediately. The
1799+
* framework layer is responsible for invoking {@link #tryCompleteActions()} to execute the
1800+
* queued actions after the write path is fully finished.
1801+
*/
1802+
private void addCompletePurgatoryAction(
1803+
Map<TableBucket, ? extends WriteResultForBucket> writeResults) {
1804+
actionQueue.add(
1805+
() -> {
1806+
for (Map.Entry<TableBucket, ? extends WriteResultForBucket> entry :
1807+
writeResults.entrySet()) {
1808+
if (entry.getValue().succeeded()) {
1809+
delayedFetchLogManager.checkAndComplete(
1810+
new DelayedTableBucketKey(entry.getKey()));
1811+
}
1812+
}
1813+
});
1814+
}
1815+
1816+
/** Tries to complete all pending delayed actions in the action queue. */
1817+
public void tryCompleteActions() {
1818+
actionQueue.tryCompleteActions();
1819+
}
1820+
17861821
private void completeDelayedOperations(TableBucket tableBucket) {
17871822
DelayedTableBucketKey delayedTableBucketKey = new DelayedTableBucketKey(tableBucket);
17881823
delayedWriteManager.checkAndComplete(delayedTableBucketKey);
Lines changed: 39 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,39 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one or more
3+
* contributor license agreements. See the NOTICE file distributed with
4+
* this work for additional information regarding copyright ownership.
5+
* The ASF licenses this file to You under the Apache License, Version 2.0
6+
* (the "License"); you may not use this file except in compliance with
7+
* the License. You may obtain a copy of the License at
8+
*
9+
* http://www.apache.org/licenses/LICENSE-2.0
10+
*
11+
* Unless required by applicable law or agreed to in writing, software
12+
* distributed under the License is distributed on an "AS IS" BASIS,
13+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
14+
* See the License for the specific language governing permissions and
15+
* limitations under the License.
16+
*/
17+
18+
package org.apache.fluss.server.replica.delay;
19+
20+
import org.apache.fluss.annotation.Internal;
21+
22+
/**
23+
* A queue for collecting actions which need to be executed later.
24+
*
25+
* <p>This is used to decouple the enqueuing of delayed operation completions from their execution.
26+
* For example, after appending records, we enqueue actions to complete delayed fetch operations,
27+
* then execute them after the write path is fully finished.
28+
*
29+
* @see DelayedActionQueue
30+
*/
31+
@Internal
32+
public interface ActionQueue {
33+
34+
/** Adds an action to this queue. */
35+
void add(Runnable action);
36+
37+
/** Tries to complete all pending actions in the queue. */
38+
void tryCompleteActions();
39+
}
Lines changed: 62 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,62 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one or more
3+
* contributor license agreements. See the NOTICE file distributed with
4+
* this work for additional information regarding copyright ownership.
5+
* The ASF licenses this file to You under the Apache License, Version 2.0
6+
* (the "License"); you may not use this file except in compliance with
7+
* the License. You may obtain a copy of the License at
8+
*
9+
* http://www.apache.org/licenses/LICENSE-2.0
10+
*
11+
* Unless required by applicable law or agreed to in writing, software
12+
* distributed under the License is distributed on an "AS IS" BASIS,
13+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
14+
* See the License for the specific language governing permissions and
15+
* limitations under the License.
16+
*/
17+
18+
package org.apache.fluss.server.replica.delay;
19+
20+
import org.apache.fluss.annotation.Internal;
21+
22+
import org.slf4j.Logger;
23+
import org.slf4j.LoggerFactory;
24+
25+
import java.util.concurrent.ConcurrentLinkedQueue;
26+
27+
/**
28+
* Default implementation of {@link ActionQueue} that collects actions into a concurrent queue and
29+
* executes them when {@link #tryCompleteActions()} is called.
30+
*
31+
* <p>Uses {@link ConcurrentLinkedQueue} for lock-free enqueue. Actions are executed and removed
32+
* from the queue when {@link #tryCompleteActions()} is called.
33+
*/
34+
@Internal
35+
public class DelayedActionQueue implements ActionQueue {
36+
private static final Logger LOG = LoggerFactory.getLogger(DelayedActionQueue.class);
37+
38+
private final ConcurrentLinkedQueue<Runnable> queue = new ConcurrentLinkedQueue<>();
39+
40+
@Override
41+
public void add(Runnable action) {
42+
queue.add(action);
43+
}
44+
45+
@Override
46+
public void tryCompleteActions() {
47+
int maxToComplete = queue.size();
48+
int count = 0;
49+
while (count < maxToComplete) {
50+
Runnable action = queue.poll();
51+
if (action == null) {
52+
break;
53+
}
54+
try {
55+
action.run();
56+
} catch (Throwable t) {
57+
LOG.error("Failed to complete delayed action.", t);
58+
}
59+
count++;
60+
}
61+
}
62+
}

fluss-server/src/main/java/org/apache/fluss/server/tablet/TabletService.java

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -184,6 +184,11 @@ public String name() {
184184
@Override
185185
public void shutdown() {}
186186

187+
@Override
188+
public void tryCompleteActions() {
189+
replicaManager.tryCompleteActions();
190+
}
191+
187192
@Override
188193
public CompletableFuture<ProduceLogResponse> produceLog(ProduceLogRequest request) {
189194
authorizeTable(WRITE, request.getTableId());

fluss-server/src/test/java/org/apache/fluss/server/replica/delay/DelayedFetchLogTest.java

Lines changed: 58 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -83,7 +83,8 @@ void testCompleteDelayedFetchLog() throws Exception {
8383
assertThat(delayedFetchLogManager.numDelayed()).isEqualTo(1);
8484
assertThat(delayedFetchLogManager.watched()).isEqualTo(1);
8585

86-
// write data.
86+
// Produce data — appendRecordsToLog enqueues a checkAndComplete action,
87+
// which is executed when tryCompleteActions() is called (normally by framework layer).
8788
assertThat(delayedResponse.isDone()).isFalse();
8889
CompletableFuture<List<ProduceLogResultForBucket>> future = new CompletableFuture<>();
8990
replicaManager.appendRecordsToLog(
@@ -92,11 +93,12 @@ void testCompleteDelayedFetchLog() throws Exception {
9293
Collections.singletonMap(tb, genMemoryLogRecordsByObject(DATA1)),
9394
null,
9495
future::complete);
96+
// Simulate framework layer trigger (FlussRequestHandler calls this after invoke).
97+
replicaManager.tryCompleteActions();
9598
assertThat(future.get()).containsOnly(new ProduceLogResultForBucket(tb, 0, 10L));
9699

97-
// check and complete manually
98-
numComplete = delayedFetchLogManager.checkAndComplete(delayedTableBucketKey);
99-
assertThat(numComplete).isEqualTo(1);
100+
// The delayed fetch log should already be completed by tryCompleteActions above.
101+
assertThat(delayedResponse.isDone()).isTrue();
100102
assertThat(delayedFetchLogManager.numDelayed()).isEqualTo(0);
101103
assertThat(delayedFetchLogManager.watched()).isEqualTo(0);
102104

@@ -107,6 +109,58 @@ void testCompleteDelayedFetchLog() throws Exception {
107109
assertLogRecordsEquals(DATA1_ROW_TYPE, resultForBucket.records(), DATA1);
108110
}
109111

112+
@Test
113+
void testProduceAutoCompletesDelayedFetchLog() throws Exception {
114+
TableBucket tb = new TableBucket(DATA1_TABLE_ID, 1);
115+
makeLogTableAsLeader(tb.getBucket());
116+
117+
// Set up a delayed fetch with follower-like params (LOG_END isolation, minBytes=1).
118+
FetchLogResultForBucket preFetchResultForBucket =
119+
new FetchLogResultForBucket(tb, MemoryLogRecords.EMPTY, 0L);
120+
CompletableFuture<Map<TableBucket, FetchLogResultForBucket>> delayedResponse =
121+
new CompletableFuture<>();
122+
DelayedFetchLog delayedFetchLog =
123+
createDelayedFetchLogRequest(
124+
tb,
125+
1, // minFetchBytes = 1, like follower fetch
126+
Duration.ofMinutes(3).toMillis(),
127+
new FetchBucketStatus(
128+
new FetchReqInfo(150001L, 0L, Integer.MAX_VALUE),
129+
new LogOffsetMetadata(0L, 0L, 0),
130+
preFetchResultForBucket),
131+
delayedResponse::complete);
132+
133+
DelayedOperationManager<DelayedFetchLog> delayedFetchLogManager =
134+
replicaManager.getDelayedFetchLogManager();
135+
DelayedTableBucketKey delayedTableBucketKey = new DelayedTableBucketKey(tb);
136+
delayedFetchLogManager.tryCompleteElseWatch(
137+
delayedFetchLog, Collections.singletonList(delayedTableBucketKey));
138+
assertThat(delayedFetchLogManager.numDelayed()).isEqualTo(1);
139+
assertThat(delayedResponse.isDone()).isFalse();
140+
141+
// Produce with acks=1 — response returns immediately, action enqueued.
142+
CompletableFuture<List<ProduceLogResultForBucket>> produceResponse =
143+
new CompletableFuture<>();
144+
replicaManager.appendRecordsToLog(
145+
20000,
146+
1,
147+
Collections.singletonMap(tb, genMemoryLogRecordsByObject(DATA1)),
148+
null,
149+
produceResponse::complete);
150+
// Simulate framework layer trigger (FlussRequestHandler calls this after invoke).
151+
replicaManager.tryCompleteActions();
152+
assertThat(produceResponse.get()).containsOnly(new ProduceLogResultForBucket(tb, 0, 10L));
153+
154+
// The delayed fetch should have been auto-completed by tryCompleteActions.
155+
assertThat(delayedResponse.isDone()).isTrue();
156+
assertThat(delayedFetchLogManager.numDelayed()).isEqualTo(0);
157+
158+
Map<TableBucket, FetchLogResultForBucket> result = delayedResponse.get();
159+
FetchLogResultForBucket resultForBucket = result.get(tb);
160+
assertThat(resultForBucket).isNotNull();
161+
assertLogRecordsEquals(DATA1_ROW_TYPE, resultForBucket.records(), DATA1);
162+
}
163+
110164
@Test
111165
void testDelayFetchLogTimeout() {
112166
TableBucket tb = new TableBucket(DATA1_TABLE_ID, 1);

0 commit comments

Comments
 (0)