This is an automated email from the ASF dual-hosted git repository.
nsivabalan pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/hudi.git
The following commit(s) were added to refs/heads/master by this push:
new 463b6159fe7f fix(timeline-service): release the server when close
fails (#20062)
463b6159fe7f is described below
commit 463b6159fe7fae8f9132acda3f473233fd265a5d
Author: Lin Liu <[email protected]>
AuthorDate: Fri Sep 25 12:36:31 2026 -0700
fix(timeline-service): release the server when close fails (#20062)
Guard each shutdown stage so one failure cannot skip the others, always
release the references,
and log a warning when a stage fails.
TimelineService.close() — each of the three stages wrapped; app released in
a finally
EmbeddedTimelineService.stopForBasePath() — server / viewManager released
in a finally
EmbeddedTimelineService.shutdownAllTimelineServers() — the sweep continues
past a failing
server, and the metric is decremented in a finally
---
.../client/embedded/EmbeddedTimelineService.java | 26 +++++--
.../embedded/TestEmbeddedTimelineService.java | 84 ++++++++++++++++++++++
.../hudi/timeline/service/TimelineService.java | 32 +++++++--
3 files changed, 131 insertions(+), 11 deletions(-)
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/embedded/EmbeddedTimelineService.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/embedded/EmbeddedTimelineService.java
index 13a075d09f19..159bc3a13346 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/embedded/EmbeddedTimelineService.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/embedded/EmbeddedTimelineService.java
@@ -123,8 +123,15 @@ public class EmbeddedTimelineService {
public static void shutdownAllTimelineServers() {
RUNNING_SERVICES.entrySet().forEach(entry -> {
log.info("Closing Timeline server");
- entry.getValue().server.close();
- METRICS_REGISTRY.set(NUM_EMBEDDED_TIMELINE_SERVERS,
NUM_SERVERS_RUNNING.decrementAndGet());
+ try {
+ entry.getValue().server.close();
+ } catch (Exception e) {
+ // Keep sweeping: an unguarded throw here would abandon every server
after this one and
+ // skip the clear() below, leaving the registry pointing at servers
nobody can reach.
+ log.warn("Timeline server did not close cleanly during shutdown;
continuing", e);
+ } finally {
+ METRICS_REGISTRY.set(NUM_EMBEDDED_TIMELINE_SERVERS,
NUM_SERVERS_RUNNING.decrementAndGet());
+ }
log.info("Closed Timeline server");
});
RUNNING_SERVICES.clear();
@@ -243,10 +250,17 @@ public class EmbeddedTimelineService {
// continue rest of shutdown outside of the synchronized block to avoid
excess blocking
if (basePaths.isEmpty() && null != server) {
log.info("Closing Timeline server");
- this.server.close();
- METRICS_REGISTRY.set(NUM_EMBEDDED_TIMELINE_SERVERS,
NUM_SERVERS_RUNNING.decrementAndGet());
- this.server = null;
- this.viewManager = null;
+ try {
+ this.server.close();
+ } catch (Exception e) {
+ // Release the references anyway: holding on to a server that failed
to close leaves this
+ // instance permanently unclosable, since every later call re-enters
this same branch.
+ log.warn("Timeline server did not close cleanly; releasing the
reference anyway", e);
+ } finally {
+ METRICS_REGISTRY.set(NUM_EMBEDDED_TIMELINE_SERVERS,
NUM_SERVERS_RUNNING.decrementAndGet());
+ this.server = null;
+ this.viewManager = null;
+ }
log.info("Closed Timeline server");
}
}
diff --git
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/client/embedded/TestEmbeddedTimelineService.java
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/client/embedded/TestEmbeddedTimelineService.java
index 6ea30694f2c5..6b1e23510ef6 100644
---
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/client/embedded/TestEmbeddedTimelineService.java
+++
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/client/embedded/TestEmbeddedTimelineService.java
@@ -31,6 +31,7 @@ import org.junit.jupiter.api.Test;
import org.mockito.Mockito;
import static
org.apache.hudi.common.testutils.HoodieTestUtils.getDefaultStorageConf;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
@@ -38,6 +39,7 @@ import static org.junit.jupiter.api.Assertions.assertNotSame;
import static org.junit.jupiter.api.Assertions.assertSame;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
@@ -219,4 +221,86 @@ public class TestEmbeddedTimelineService extends
HoodieCommonTestHarness {
verify(mockService2,
times(1)).unregisterBasePath(writeConfig2.getBasePath());
verify(mockService2, times(1)).close();
}
+
+ /**
+ * A timeline service whose close() throws must still leave this instance
releasable.
+ *
+ * <p>The discriminating input is a TimelineService that throws from
close(): the rows returned
+ * and the number of services created are identical either way, so only the
state left behind
+ * after a failed close separates the two behaviours. Before the guard, the
throw propagated out
+ * of stopForBasePath and left `server` non-null, so the instance could
never be closed.
+ */
+ @Test
+ public void stopForBasePathReleasesServerWhenCloseThrows() throws Exception {
+ HoodieEngineContext engineContext = new
HoodieLocalEngineContext(getDefaultStorageConf());
+ HoodieWriteConfig writeConfig = HoodieWriteConfig.newBuilder()
+ .withPath(tempDir.resolve("table_close_throws").toString())
+ .withEmbeddedTimelineServerEnabled(true)
+ .build();
+ EmbeddedTimelineService.TimelineServiceCreator mockCreator =
+ Mockito.mock(EmbeddedTimelineService.TimelineServiceCreator.class);
+ TimelineService mockService = Mockito.mock(TimelineService.class);
+ when(mockCreator.create(any(), any(), any())).thenReturn(mockService);
+ when(mockService.startService()).thenReturn(456);
+ doThrow(new RuntimeException("jetty refused to
stop")).when(mockService).close();
+
+ EmbeddedTimelineService service =
EmbeddedTimelineService.getOrStartEmbeddedTimelineService(
+ engineContext, null, writeConfig, mockCreator);
+
+ // The failing close must not escape.
+ assertDoesNotThrow(() ->
service.stopForBasePath(writeConfig.getBasePath()));
+ verify(mockService, times(1)).close();
+
+ // And the reference must have been released, so a second stop does not
re-enter the close
+ // branch. Without the guard `server` is still set here and close() would
be called again.
+ assertDoesNotThrow(() ->
service.stopForBasePath(writeConfig.getBasePath()));
+ verify(mockService, times(1)).close();
+ }
+
+ /**
+ * One server that fails to close must not abandon the others in the
registry.
+ *
+ * <p>Discriminating input: two registered services where the close of one
throws. Before the
+ * guard the forEach aborted on the first throw, so the second server was
never closed and
+ * RUNNING_SERVICES was never cleared.
+ */
+ @Test
+ public void shutdownAllContinuesPastAFailingServer() throws Exception {
+ HoodieEngineContext engineContext = new
HoodieLocalEngineContext(getDefaultStorageConf());
+ HoodieWriteConfig throwingConfig = HoodieWriteConfig.newBuilder()
+ .withPath(tempDir.resolve("table_throwing").toString())
+ .withEmbeddedTimelineServerEnabled(true)
+ .withEmbeddedTimelineServerReuseEnabled(true)
+
.withMetadataConfig(HoodieMetadataConfig.newBuilder().enable(true).build())
+ .build();
+ EmbeddedTimelineService.TimelineServiceCreator throwingCreator =
+ Mockito.mock(EmbeddedTimelineService.TimelineServiceCreator.class);
+ TimelineService throwingService = Mockito.mock(TimelineService.class);
+ when(throwingCreator.create(any(), any(),
any())).thenReturn(throwingService);
+ when(throwingService.startService()).thenReturn(654);
+ doThrow(new RuntimeException("jetty refused to
stop")).when(throwingService).close();
+ EmbeddedTimelineService.getOrStartEmbeddedTimelineService(
+ engineContext, null, throwingConfig, throwingCreator);
+
+ // A different identifier, so this lands as a separate entry in
RUNNING_SERVICES.
+ HoodieWriteConfig healthyConfig = HoodieWriteConfig.newBuilder()
+ .withPath(tempDir.resolve("table_healthy").toString())
+ .withEmbeddedTimelineServerEnabled(true)
+ .withEmbeddedTimelineServerReuseEnabled(true)
+
.withMetadataConfig(HoodieMetadataConfig.newBuilder().enable(false).build())
+ .build();
+ EmbeddedTimelineService.TimelineServiceCreator healthyCreator =
+ Mockito.mock(EmbeddedTimelineService.TimelineServiceCreator.class);
+ TimelineService healthyService = Mockito.mock(TimelineService.class);
+ when(healthyCreator.create(any(), any(),
any())).thenReturn(healthyService);
+ when(healthyService.startService()).thenReturn(655);
+ EmbeddedTimelineService.getOrStartEmbeddedTimelineService(
+ engineContext, null, healthyConfig, healthyCreator);
+
+ assertDoesNotThrow(EmbeddedTimelineService::shutdownAllTimelineServers);
+
+ // Both were attempted regardless of iteration order.
+ verify(throwingService, times(1)).close();
+ verify(healthyService, times(1)).close();
+ }
}
diff --git
a/hudi-timeline-service/src/main/java/org/apache/hudi/timeline/service/TimelineService.java
b/hudi-timeline-service/src/main/java/org/apache/hudi/timeline/service/TimelineService.java
index 95bda8f83a6e..f41ea070d019 100644
---
a/hudi-timeline-service/src/main/java/org/apache/hudi/timeline/service/TimelineService.java
+++
b/hudi-timeline-service/src/main/java/org/apache/hudi/timeline/service/TimelineService.java
@@ -303,16 +303,38 @@ public class TimelineService {
}
}
+ /**
+ * Shuts the service down.
+ *
+ * <p>Each stage is guarded so a failure in one does not skip the rest, and
{@code app} is
+ * released either way. Without that, a throw from any stage leaves {@code
app} set and the
+ * caller's reference non-null, so the instance can never be closed on a
later attempt while
+ * its Jetty threads keep running.
+ */
public void close() {
log.info("Closing Timeline Service with port {}", serverPort);
- if (requestHandler != null) {
- this.requestHandler.stop();
+ try {
+ if (requestHandler != null) {
+ this.requestHandler.stop();
+ }
+ } catch (Exception e) {
+ log.warn("Failed to stop the timeline request handler on port {};
continuing shutdown",
+ serverPort, e);
}
- if (this.app != null) {
- this.app.stop();
+ try {
+ if (this.app != null) {
+ this.app.stop();
+ }
+ } catch (Exception e) {
+ log.warn("Failed to stop the Javalin app on port {}; continuing
shutdown", serverPort, e);
+ } finally {
this.app = null;
}
- this.fsViewsManager.close();
+ try {
+ this.fsViewsManager.close();
+ } catch (Exception e) {
+ log.warn("Failed to close the file system view manager on port {}",
serverPort, e);
+ }
log.info("Closed Timeline Service with port {}", serverPort);
}