diff --git a/modules/core/ingest-core/src/main/java/com/lingniu/ingest/core/concurrency/DisruptorEventBus.java b/modules/core/ingest-core/src/main/java/com/lingniu/ingest/core/concurrency/DisruptorEventBus.java index 6558bf33..f3ab22ec 100644 --- a/modules/core/ingest-core/src/main/java/com/lingniu/ingest/core/concurrency/DisruptorEventBus.java +++ b/modules/core/ingest-core/src/main/java/com/lingniu/ingest/core/concurrency/DisruptorEventBus.java @@ -67,10 +67,19 @@ public final class DisruptorEventBus implements AutoCloseable { public CompletableFuture publishAndAwait(VehicleEvent event, String requiredSinkName) { String requiredSink = Objects.requireNonNull(requiredSinkName, "requiredSinkName"); published.incrementAndGet(); - List> requiredFutures = new ArrayList<>(); - for (EventSink sink : sinks) { - if (!sink.accepts(event)) continue; + List acceptingSinks = sinks.stream() + .filter(sink -> sink.accepts(event)) + .toList(); + List requiredSinks = acceptingSinks.stream() + .filter(sink -> requiredSink.equals(sink.name())) + .toList(); + if (requiredSinks.isEmpty()) { + return CompletableFuture.failedFuture(new IllegalStateException( + "required sink '" + requiredSink + "' did not accept event " + event.eventId())); + } + List> requiredFutures = new ArrayList<>(); + for (EventSink sink : acceptingSinks) { CompletableFuture future = CompletableFuture .supplyAsync(() -> publishToSinkAndTrack(sink, event), awaitExecutor) .thenCompose(f -> f); @@ -78,10 +87,6 @@ public final class DisruptorEventBus implements AutoCloseable { requiredFutures.add(future); } } - if (requiredFutures.isEmpty()) { - return CompletableFuture.failedFuture(new IllegalStateException( - "required sink '" + requiredSink + "' did not accept event " + event.eventId())); - } return CompletableFuture.allOf(requiredFutures.toArray(CompletableFuture[]::new)); } diff --git a/modules/core/ingest-core/src/test/java/com/lingniu/ingest/core/concurrency/DisruptorEventBusAwaitTest.java b/modules/core/ingest-core/src/test/java/com/lingniu/ingest/core/concurrency/DisruptorEventBusAwaitTest.java index 2e4fe905..1c39eab6 100644 --- a/modules/core/ingest-core/src/test/java/com/lingniu/ingest/core/concurrency/DisruptorEventBusAwaitTest.java +++ b/modules/core/ingest-core/src/test/java/com/lingniu/ingest/core/concurrency/DisruptorEventBusAwaitTest.java @@ -6,12 +6,16 @@ import com.lingniu.ingest.api.sink.EventSink; import org.junit.jupiter.api.Test; import java.time.Instant; +import java.util.List; import java.util.Map; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionException; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLong; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; class DisruptorEventBusAwaitTest { @@ -27,6 +31,20 @@ class DisruptorEventBusAwaitTest { assertThat(sink.publishThreadId.get()).isNotEqualTo(callerThreadId); } + @Test + void publishAndAwaitDoesNotInvokeOptionalSinksWhenRequiredSinkIsAbsent() throws Exception { + CapturingSink optionalSink = new CapturingSink("event-file-store"); + + try (DisruptorEventBus bus = new DisruptorEventBus(1024, "blocking", List.of(optionalSink))) { + assertThatThrownBy(() -> bus.publishAndAwait(event(), "kafka").join()) + .isInstanceOf(CompletionException.class) + .hasRootCauseInstanceOf(IllegalStateException.class) + .hasMessageContaining("kafka"); + + assertThat(optionalSink.published.await(200, TimeUnit.MILLISECONDS)).isFalse(); + } + } + private static VehicleEvent.Heartbeat event() { Instant now = Instant.parse("2026-06-23T10:00:00Z"); return new VehicleEvent.Heartbeat( @@ -42,6 +60,7 @@ class DisruptorEventBusAwaitTest { private static final class CapturingSink implements EventSink { private final String name; private final AtomicLong publishThreadId = new AtomicLong(-1); + private final CountDownLatch published = new CountDownLatch(1); private CapturingSink(String name) { this.name = name; @@ -55,6 +74,7 @@ class DisruptorEventBusAwaitTest { @Override public CompletableFuture publish(VehicleEvent event) { publishThreadId.set(Thread.currentThread().threadId()); + published.countDown(); return CompletableFuture.completedFuture(null); } }