hudi-agent commented on code in PR #19811:
URL: https://github.com/apache/hudi/pull/19811#discussion_r3961973667
##########
hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/HoodieMultiTableStreamer.java:
##########
@@ -389,6 +406,12 @@ public static class Config implements Serializable {
+ " source-fetch -> Transform -> Hudi Write in loop")
public Boolean continuousMode = false;
+ @Parameter(names = {"--fail-fast-on-continuous"},
+ description = "Only applies in continuous mode. When enabled, the
failure of any single table sync fails the "
+ + "whole job. The remaining table syncs are shut down and the
process exits with a non-zero status. When "
+ + "disabled (default), each table is synced independently and a
single failure does not stop the others.")
+ public Boolean failFastOnContinuousMode = false;
+
@Parameter(names = {"--min-sync-interval-seconds"},
Review Comment:
🤖 nit: the CLI flag is `--fail-fast-on-continuous` but the field is
`failFastOnContinuousMode`; could you align them (e.g.
`--fail-fast-on-continuous-mode`) so it's easier to grep between the two?
<sub><i>⚠️ AI-generated; verify before applying. React 👍/👎 to flag
quality.</i></sub>
##########
hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/HoodieMultiTableStreamer.java:
##########
@@ -461,28 +484,194 @@ private static String resetTarget(Config configuration,
String database, String
/**
* Creates actual HoodieDeltaStreamer objects for every table/topic and does
incremental sync.
+ *
+ * <p>In continuous mode each table's sync blocks until it is shut down, so
the tables are synced concurrently.
+ * Otherwise the tables are synced sequentially, one after another.
*/
public void sync() {
+ try {
+ if (continuousMode) {
+ syncContinuously();
+ } else {
+ syncSequentially();
+ }
+ } finally {
+ log.info("Ingestion was successful for topics: {}", successTables);
+ if (!failedTables.isEmpty()) {
+ log.info("Ingestion failed for topics: {}", failedTables);
+ }
+ }
+ }
+
+ private void syncSequentially() {
for (TableExecutionContext context : tableExecutionContexts) {
HoodieStreamer streamer = null;
try {
streamer = new HoodieStreamer(context.getConfig(), jssc,
Option.ofNullable(context.getProperties()));
streamer.sync();
successTables.add(Helpers.getTableWithDatabase(context));
- streamer.shutdownGracefully();
} catch (Exception e) {
log.error("error while running MultiTableDeltaStreamer for table: {}",
context.getTableName(), e);
failedTables.add(Helpers.getTableWithDatabase(context));
} finally {
if (streamer != null) {
- streamer.shutdownGracefully();
+ shutdownQuietly(streamer, context);
}
}
}
+ }
+
+ /**
+ * Syncs all tables concurrently, one thread per table. Used for continuous
mode where each table's sync blocks
+ * indefinitely.
+ *
+ * <p>When {@code --fail-fast-on-continuous} is enabled, the first table
failure fails the whole job. The sibling
+ * streamers are shut down and a {@link HoodieException} is thrown so the
caller can exit with a non-zero status.
+ * Otherwise, every table is synced independently and a single failure does
not affect the others.
+ */
+ private void syncContinuously() {
+ if (tableExecutionContexts.isEmpty()) {
+ return;
+ }
+ // Streamer instances are registered from worker threads, so a thread-safe
list is required.
+ final List<HoodieStreamer> streamerInstances = new
CopyOnWriteArrayList<>();
+ // Set once fail fast trips, so tasks that register their streamer
afterwards stop before starting the sync.
+ final AtomicBoolean shutdownRequested = new AtomicBoolean(false);
+ final ExecutorService executor =
Executors.newFixedThreadPool(tableExecutionContexts.size(),
+ new CustomizedThreadFactory("multi-table-streamer", true));
+ boolean terminated = false;
+ try {
+ final List<CompletableFuture<Void>> tableFutures =
tableExecutionContexts.stream()
+ .map(context -> CompletableFuture.runAsync(
+ () -> runTableSync(context, streamerInstances,
shutdownRequested), executor))
+ .collect(Collectors.toList());
+
+ if (failFastOnContinuousMode) {
+ log.info("Fail fast enabled in continuous mode. The whole job fails on
any single table failure");
+ awaitFailFast(tableFutures, streamerInstances, shutdownRequested);
+ } else {
+ CompletableFuture.allOf(tableFutures.toArray(new
CompletableFuture[0])).join();
+ }
+ } finally {
+ // Wait for every worker thread to finish (including its finally
cleanup) before returning, so sync() does not
+ // return while a table is still writing and main() then stops the
shared Spark context under it.
+ terminated = shutdownExecutor(executor);
+ }
+ // If the workers never terminated, ingestion may still be running. Fail
loudly instead of returning as if the
+ // cleanup succeeded, so the caller does not silently proceed to Spark
teardown with live writers.
+ if (!terminated) {
+ throw new HoodieException("Timed out shutting down table ingestion
workers in continuous mode");
+ }
+ }
+
+ /**
+ * Syncs one table on the calling worker thread. Rethrows only under fail
fast; otherwise the failure is recorded
+ * in {@link #failedTables} and the sibling tables carry on.
+ */
+ private void runTableSync(TableExecutionContext context,
List<HoodieStreamer> streamerInstances, AtomicBoolean shutdownRequested) {
+ HoodieStreamer streamer = null;
+ try {
+ streamer = new HoodieStreamer(context.getConfig(), jssc,
Option.ofNullable(context.getProperties()));
+ streamerInstances.add(streamer);
+ // Register before checking the flag so a concurrent shutdownStreamers()
always sees this streamer.
+ if (shutdownRequested.get()) {
+ return;
+ }
+ streamer.sync();
+ // A streamer registered just before fail fast tripped can reach here
without ever ingesting.
+ // shutdown() call will be a no-op because its ingestion service hadn't
started yet.
+ // Don't count that as a success.
+ if (!shutdownRequested.get()) {
+ successTables.add(Helpers.getTableWithDatabase(context));
+ }
+ } catch (Exception e) {
+ String table = Helpers.getTableWithDatabase(context);
+ log.error("error while running MultiTableDeltaStreamer for table: {}",
table, e);
+ failedTables.add(table);
+ if (failFastOnContinuousMode) {
+ // Name the table so the thrown exception identifies the culprit, not
the siblings torn down after it.
+ throw new HoodieException("Table sync failed in continuous mode for
table: " + table, e);
+ }
+ } finally {
+ if (streamer != null) {
+ shutdownQuietly(streamer, context);
+ }
+ }
+ }
+
+ /**
+ * Waits until either every table sync finishes successfully or the first
one fails. On the first failure, the
+ * remaining streamers are shut down and a {@link HoodieException} is
thrown. {@link FutureUtils#allOf} only trips
+ * on an <em>exceptional</em> completion, so a table that terminates
normally (e.g. via a
+ * {@link PostWriteTerminationStrategy}) does not abort its siblings.
+ */
+ private void awaitFailFast(List<CompletableFuture<Void>> tableFutures,
List<HoodieStreamer> streamerInstances, AtomicBoolean shutdownRequested) {
+ try {
+ FutureUtils.allOf(tableFutures).join();
+ } catch (CompletionException e) {
+ Throwable cause = unwrapFailFastFailure(e);
+ log.error("error while running MultiTableDeltaStreamer, shutting down
remaining tables as fail fast is enabled", cause);
+ shutdownRequested.set(true);
+ // shutdownStreamers only interrupts; the executor teardown in
syncContinuously() waits for the siblings to stop.
+ shutdownStreamers(streamerInstances);
+ throw new HoodieException("Fail fast is enabled and a table sync failed
in continuous mode.", cause);
+ }
+ }
+
+ /**
+ * Releases a streamer's resources, logging rather than propagating a
failure to do so. Closing can throw, and this
+ * runs in a {@code finally} on the failure path where escaping would mask
the table failure and, in
+ * {@link #syncSequentially()}, abort the tables not synced yet.
+ */
+ private static void shutdownQuietly(HoodieStreamer streamer,
TableExecutionContext context) {
+ try {
+ streamer.shutdownGracefully();
+ } catch (Exception e) {
+ log.warn("error while shutting down the streamer for table: {}",
context.getTableName(), e);
+ }
+ }
+
+ // FutureUtils.allOf can wrap the real failure in more than one layer of
CompletionException.
+ private static Throwable unwrapFailFastFailure(CompletionException e) {
+ Throwable cause = e;
Review Comment:
🤖 nit: nothing in this method is fail-fast specific; would you consider
naming it `unwrapCompletionException` so it reads as a general helper rather
than something tied to the fail-fast path?
<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]