fix: preflight required kafka sink
This commit is contained in:
@@ -67,10 +67,19 @@ public final class DisruptorEventBus implements AutoCloseable {
|
|||||||
public CompletableFuture<Void> publishAndAwait(VehicleEvent event, String requiredSinkName) {
|
public CompletableFuture<Void> publishAndAwait(VehicleEvent event, String requiredSinkName) {
|
||||||
String requiredSink = Objects.requireNonNull(requiredSinkName, "requiredSinkName");
|
String requiredSink = Objects.requireNonNull(requiredSinkName, "requiredSinkName");
|
||||||
published.incrementAndGet();
|
published.incrementAndGet();
|
||||||
List<CompletableFuture<Void>> requiredFutures = new ArrayList<>();
|
List<EventSink> acceptingSinks = sinks.stream()
|
||||||
for (EventSink sink : sinks) {
|
.filter(sink -> sink.accepts(event))
|
||||||
if (!sink.accepts(event)) continue;
|
.toList();
|
||||||
|
List<EventSink> 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<CompletableFuture<Void>> requiredFutures = new ArrayList<>();
|
||||||
|
for (EventSink sink : acceptingSinks) {
|
||||||
CompletableFuture<Void> future = CompletableFuture
|
CompletableFuture<Void> future = CompletableFuture
|
||||||
.supplyAsync(() -> publishToSinkAndTrack(sink, event), awaitExecutor)
|
.supplyAsync(() -> publishToSinkAndTrack(sink, event), awaitExecutor)
|
||||||
.thenCompose(f -> f);
|
.thenCompose(f -> f);
|
||||||
@@ -78,10 +87,6 @@ public final class DisruptorEventBus implements AutoCloseable {
|
|||||||
requiredFutures.add(future);
|
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));
|
return CompletableFuture.allOf(requiredFutures.toArray(CompletableFuture[]::new));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -6,12 +6,16 @@ import com.lingniu.ingest.api.sink.EventSink;
|
|||||||
import org.junit.jupiter.api.Test;
|
import org.junit.jupiter.api.Test;
|
||||||
|
|
||||||
import java.time.Instant;
|
import java.time.Instant;
|
||||||
|
import java.util.List;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.concurrent.CompletableFuture;
|
import java.util.concurrent.CompletableFuture;
|
||||||
|
import java.util.concurrent.CompletionException;
|
||||||
|
import java.util.concurrent.CountDownLatch;
|
||||||
import java.util.concurrent.TimeUnit;
|
import java.util.concurrent.TimeUnit;
|
||||||
import java.util.concurrent.atomic.AtomicLong;
|
import java.util.concurrent.atomic.AtomicLong;
|
||||||
|
|
||||||
import static org.assertj.core.api.Assertions.assertThat;
|
import static org.assertj.core.api.Assertions.assertThat;
|
||||||
|
import static org.assertj.core.api.Assertions.assertThatThrownBy;
|
||||||
|
|
||||||
class DisruptorEventBusAwaitTest {
|
class DisruptorEventBusAwaitTest {
|
||||||
|
|
||||||
@@ -27,6 +31,20 @@ class DisruptorEventBusAwaitTest {
|
|||||||
assertThat(sink.publishThreadId.get()).isNotEqualTo(callerThreadId);
|
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() {
|
private static VehicleEvent.Heartbeat event() {
|
||||||
Instant now = Instant.parse("2026-06-23T10:00:00Z");
|
Instant now = Instant.parse("2026-06-23T10:00:00Z");
|
||||||
return new VehicleEvent.Heartbeat(
|
return new VehicleEvent.Heartbeat(
|
||||||
@@ -42,6 +60,7 @@ class DisruptorEventBusAwaitTest {
|
|||||||
private static final class CapturingSink implements EventSink {
|
private static final class CapturingSink implements EventSink {
|
||||||
private final String name;
|
private final String name;
|
||||||
private final AtomicLong publishThreadId = new AtomicLong(-1);
|
private final AtomicLong publishThreadId = new AtomicLong(-1);
|
||||||
|
private final CountDownLatch published = new CountDownLatch(1);
|
||||||
|
|
||||||
private CapturingSink(String name) {
|
private CapturingSink(String name) {
|
||||||
this.name = name;
|
this.name = name;
|
||||||
@@ -55,6 +74,7 @@ class DisruptorEventBusAwaitTest {
|
|||||||
@Override
|
@Override
|
||||||
public CompletableFuture<Void> publish(VehicleEvent event) {
|
public CompletableFuture<Void> publish(VehicleEvent event) {
|
||||||
publishThreadId.set(Thread.currentThread().threadId());
|
publishThreadId.set(Thread.currentThread().threadId());
|
||||||
|
published.countDown();
|
||||||
return CompletableFuture.completedFuture(null);
|
return CompletableFuture.completedFuture(null);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user