nsivabalan commented on code in PR #19033: URL: https://github.com/apache/hudi/pull/19033#discussion_r3710076654
########## hudi-sync/hudi-hive-sync/src/main/java/org/apache/hudi/hive/util/ParallelDispatch.java: ########## @@ -0,0 +1,240 @@ +/* + * 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.hive.util; + +import org.apache.hudi.common.util.VisibleForTesting; + +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.Callable; +import java.util.concurrent.CancellationException; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.Future; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; + +/** + * Handle to one fan-out batch of partition work: the submitted futures plus the shared + * abort flag the tasks consult before running. + * + * <p>The abort flag exists because waiting on futures in <i>submission</i> order is not + * enough to stop queued work. If a later task fails quickly while an earlier one is slow, + * the awaiting thread is still parked on the earlier {@code Future.get()}, and the + * executor happily keeps starting every queued task in the meantime. By the time the + * failure is observed, most of the "not-yet-started" work has already run. + * + * <p>Two mechanisms fix that, and both are needed: + * <ul> + * <li>a {@link CountDownLatch} tripped by the <i>first</i> abort, so the awaiting + * thread wakes on failure rather than on its turn in submission order;</li> + * <li>a task-side {@link #aborted()} check on entry, because cancelling from the + * awaiting thread is inherently late — a worker can pull its next task off the + * queue at any moment.</li> + * </ul> + * + * <p>Shared by {@link HiveDriverPool} (Hive {@code Driver} statements) and + * {@link IMetaStoreClientPool} (Thrift {@code dropPartition} batches), which fan out over + * different execution models but need identical abort-on-first-error semantics. + */ +public final class ParallelDispatch { + + private final List<Future<?>> futures; + private final int total; + private final AtomicInteger settled = new AtomicInteger(0); + private final AtomicBoolean aborted = new AtomicBoolean(false); + private final CountDownLatch done = new CountDownLatch(1); + private volatile boolean sealed; + + ParallelDispatch(int total) { + this.total = total; + this.futures = new ArrayList<>(total); + } + + void add(Future<?> future) { + futures.add(future); + } + + // Called once submission finishes. A task that settles before the last submit + // would otherwise see settled < total and never trip the latch, so re-check here. + void sealed() { + sealed = true; + signalIfComplete(); + } + + boolean aborted() { + return aborted.get(); + } + + void abort() { + aborted.set(true); + done.countDown(); + } + + void taskSettled() { + settled.incrementAndGet(); + signalIfComplete(); + } + + private void signalIfComplete() { + if (sealed && settled.get() >= total) { + done.countDown(); + } + } + + void awaitSettledOrAborted() { + if (total == 0) { + return; + } + try { + done.await(); + } catch (InterruptedException ie) { + Thread.currentThread().interrupt(); + aborted.set(true); + } + } + + // mayInterruptIfRunning=false: a worker may be mid-statement against a Hive Driver or + // a Thrift client, and we don't want to tear that down partway. Cancel only tasks that + // haven't started; in-flight work runs to completion. + int cancelPending() { + int cancelled = 0; + for (Future<?> f : futures) { + if (f.cancel(false)) { + cancelled++; + } + } + return cancelled; + } + + List<Future<?>> futures() { + return futures; + } + + /** + * Wraps {@code body} so it observes this batch's abort flag: it skips itself if a + * sibling has already failed, trips the flag if it fails, and always records that it + * settled so the awaiting thread can be released. + */ + Callable<Void> guard(Callable<Void> body, String skipMessage) { + return () -> { + if (aborted()) { + throw new CancellationException(skipMessage); + } + try { + return body.call(); + } catch (Throwable t) { + abort(); + throw t; + } finally { + taskSettled(); + } + }; + } + + /** + * Waits for the batch to settle (or abort), cancels whatever had not started, and + * returns the outcome. Errors are observed in <i>completion</i> order, not submission + * order, so a failure on a fast worker stops the other queues even while a slow worker + * is still mid-statement. + */ + Outcome awaitOutcome() { + awaitSettledOrAborted(); + int cancelled = cancelPending(); + + Exception firstError = null; + int completed = 0; + List<Exception> suppressed = new ArrayList<>(); + for (Future<?> f : futures) { + try { + f.get(); + completed++; + } catch (CancellationException ce) { + // Either we cancelled it before it started, or the task itself observed the + // abort flag and bailed. Not a new failure; just note it for the summary. + cancelled++; Review Comment: The double-count is real and reproduces exactly as you describe (`cancelled == 4` for one failure plus two never-started). I tried to fix it and am deliberately **not** shipping the fix in this PR, because every version I wrote traded a wrong log number for a swallowed exception. Details, since I think this is worth your eyes: Counting in only one place looks safe — a future cancelled in `cancelPending()` also surfaces as `CancellationException` from `get()`, so dropping the first count should be lossless. But it removes the thing that was masking a pre-existing race in the wake protocol. In `guard`: ```java } catch (Throwable t) { abort(t); // trips the latch -> awaitOutcome wakes here throw t; } finally { taskSettled(); // has not run yet; the future is not yet complete } ``` `abort()` counts down `done`, so `awaitSettledOrAborted()` returns while the failing task is still between its `catch` and its `finally`. `awaitOutcome` then calls `cancelPending()`, which cancels **the failing task's own future** before its result was ever collected. Measured on this branch, one failure plus two never-started: ``` current FAIL=EXEC:IllegalStateException cancelled=4 (wrong) failure preserved naive fix FAIL=BARE_CANCEL cancelled=2 (right) failure destroyed ``` So today's over-count is exactly what keeps the real cause alive — `cancelPending()`'s tally papers over the future it just corrupted. With the count deduplicated, a genuine metastore error on a DROP batch is reported to the operator as a bare `CancellationException`. I tried four variants — gating `signalIfComplete` on `sealed`, moving the wake into `taskSettled`, a separate task-side `skipped` counter, and an `isDone()` guard before `cancel(false)`. The best got the failure rate to ~1/120 in a 120-iteration stress run, but `isDone()` can only narrow the window, never close it: the task is mid-unwind, so `isDone()` is false right up until it is not. A correct fix means the failing task must fully settle before the awaiting thread wakes — i.e. reworking the `abort()` / `taskSettled()` / `signalIfComplete()` protocol. That is the same machinery implementing the abort-on-first-error guarantee from your earlier review, so I would rather not rework it as a side effect of a log-counter fix in this PR. As it stands the cost is one wrong number in an INFO/ERROR line that nothing reads programmatically. Happy to do it here if you would prefer, or as a focused follow-up — your call. Either way the race is worth recording, since it is latent in the current code and the obvious cleanup is what exposes it. -- 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]
