hudi-agent commented on code in PR #19485:
URL: https://github.com/apache/hudi/pull/19485#discussion_r3703323676
##########
hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/HoodieDeltaStreamerTestBase.java:
##########
@@ -763,22 +766,47 @@ static HoodieInstant
assertCommitMetadataForIncrSource(String expected, String t
return lastInstant;
}
+ /**
+ * Polls {@code condition} until it holds, the deltastreamer future
finishes, or the timeout expires.
+ *
+ * <p>On timeout the last error the condition threw is attached to the
failure. Without it the only
+ * output is a bare {@link TimeoutException} pointing at this method,
which says nothing about which
+ * assertion never held - the reason HUDI-6843 stayed open: every report
of it looks identical.
+ */
static void waitTillCondition(Function<Boolean, Boolean> condition, Future
dsFuture, long timeoutInSecs) throws Exception {
- Future<Boolean> res = Executors.newSingleThreadExecutor().submit(() -> {
- boolean ret = false;
- while (!ret && !dsFuture.isDone()) {
- try {
- Thread.sleep(2000);
- ret = condition.apply(true);
- log.info("Condition completed successfully");
- } catch (Throwable error) {
- log.debug("Got error waiting for condition", error);
- ret = false;
+ AtomicReference<Throwable> lastError = new AtomicReference<>();
+ ExecutorService executor = Executors.newSingleThreadExecutor();
+ try {
+ Future<Boolean> res = executor.submit(() -> {
+ boolean ret = false;
+ while (!ret && !dsFuture.isDone()) {
+ try {
+ Thread.sleep(2000);
+ ret = condition.apply(true);
+ if (ret) {
+ log.info("Condition completed successfully");
+ }
+ } catch (Throwable error) {
Review Comment:
🤖 Since this catches `Throwable` (which includes `InterruptedException`) and
just loops again, does the new `executor.shutdownNow()` in the `finally`
actually stop the worker when `dsFuture` never completes — the exact timeout
case this targets? `Thread.sleep` clears the interrupt flag when it throws, so
the loop re-enters `sleep(2000)` and keeps polling. Might be worth breaking out
on interrupt (e.g. `!Thread.currentThread().isInterrupted()` in the loop
condition) so the leak the comment references is really closed.
<sub><i>⚠️ AI-generated; verify before applying. React 👍/👎 to flag
quality.</i></sub>
##########
hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestWaitTillCondition.java:
##########
@@ -0,0 +1,80 @@
+/*
+ * 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.hudi.utilities.deltastreamer;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.Future;
+
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Covers {@code HoodieDeltaStreamerTestBase.TestHelpers#waitTillCondition},
the helper every
+ * continuous-mode deltastreamer test waits on.
+ *
+ * <p>HUDI-6843 is a flaky timeout in that wait whose only output was
+ * {@code java.util.concurrent.TimeoutException} at this method, with no
indication of which assertion in
+ * the condition never held - the condition's error was logged at debug and
discarded. That is why every
+ * report of the flake looks the same and none of them is actionable.
+ */
+class TestWaitTillCondition {
+
+ /** A deltastreamer future that never finishes, as a continuous-mode job
would be. */
+ private static final Future<?> RUNNING = new CompletableFuture<>();
+
+ @Test
+ void timeoutFailureNamesTheLastConditionFailure() {
+ String assertionText = "assertAtleastNDeltaCommits: expected at least 3
delta commits but got 2";
+
+ AssertionError error = assertThrows(AssertionError.class,
+ () -> HoodieDeltaStreamerTestBase.TestHelpers.waitTillCondition(
+ ignored -> {
+ throw new AssertionError(assertionText);
Review Comment:
🤖 The poll interval inside `waitTillCondition` is 2s but this timeout is 3s,
so there's only ~1s of slack before the first evaluation has to record
`lastError`. On a loaded CI, could a delayed worker startup or sleep overrun
push the first evaluation past 3s, leaving `lastError` null so the
`assertionText` assertion below fails? Widening the timeout would give this
de-flaking test more margin.
<sub><i>⚠️ AI-generated; verify before applying. React 👍/👎 to flag
quality.</i></sub>
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]