Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -20,10 +20,12 @@

import static org.apache.pulsar.broker.loadbalance.extensions.models.SplitDecision.Label.Failure;
import static org.apache.pulsar.broker.loadbalance.extensions.models.SplitDecision.Reason.Unknown;
import com.google.common.annotations.VisibleForTesting;
import java.util.Map;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import lombok.CustomLog;
import org.apache.pulsar.broker.loadbalance.extensions.channel.ServiceUnitState;
import org.apache.pulsar.broker.loadbalance.extensions.channel.ServiceUnitStateData;
Expand All @@ -47,39 +49,38 @@ public SplitManager(SplitCounter splitCounter) {
}

private void complete(String serviceUnit, Throwable ex) {
inFlightSplitRequests.computeIfPresent(serviceUnit, (__, future) -> {
if (!future.isDone()) {
if (ex != null) {
future.completeExceptionally(ex);
} else {
future.complete(null);
}
}
return null;
});
CompletableFuture<Void> future = inFlightSplitRequests.remove(serviceUnit);
if (future == null || future.isDone()) {
return;
}
if (ex != null) {
future.completeExceptionally(ex);
} else {
future.complete(null);
}
}

public CompletableFuture<Void> waitAsync(CompletableFuture<Void> eventPubFuture,
String bundle,
SplitDecision decision,
long timeout,
TimeUnit timeoutUnit) {
return eventPubFuture
.thenCompose(__ -> inFlightSplitRequests.computeIfAbsent(bundle, ignore -> {
log.info().attr("bundle", bundle).attr("timeout", timeout)
.attr("timeoutUnit", timeoutUnit)
.log("Published the bundle split event for bundle: . "
+ "Waiting the split event to complete. Timeout");
CompletableFuture<Void> future = new CompletableFuture<>();
future.orTimeout(timeout, timeoutUnit).whenComplete((v, ex) -> {
if (ex != null) {
inFlightSplitRequests.remove(bundle);
log.warn().attr("bundle", bundle).exception(ex)
.log("Timed out while waiting for the bundle split event");
}
});
return future;
}))
return eventPubFuture.thenCompose(__ -> {
CompletableFuture<Void> future = inFlightSplitRequests.computeIfAbsent(bundle, ignore -> {
log.info().attr("bundle", bundle).attr("timeout", timeout)
.attr("timeoutUnit", timeoutUnit)
.log("Published the bundle split event for bundle: . "
+ "Waiting the split event to complete. Timeout");
return new CompletableFuture<Void>().orTimeout(timeout, timeoutUnit);
});
// Return the dependent stage so callers cannot observe completion before timeout cleanup finishes.
return future.whenComplete((v, ex) -> {
if (ex instanceof TimeoutException && inFlightSplitRequests.remove(bundle, future)) {
log.warn().attr("bundle", bundle).exception(ex)
.log("Timed out while waiting for the bundle split event");
}
});
})
.whenComplete((__, ex) -> {
if (ex != null) {
log.error().attr("bundle", bundle).exception(ex)
Expand Down Expand Up @@ -109,12 +110,16 @@ public void handleEvent(String serviceUnit, ServiceUnitStateData data, Throwable

public void close() {
inFlightSplitRequests.forEach((bundle, future) -> {
if (!future.isDone()) {
if (inFlightSplitRequests.remove(bundle, future) && !future.isDone()) {
String msg = String.format("Splitting bundle: %s, but the manager already closed.", bundle);
log.warn(msg);
future.completeExceptionally(new IllegalStateException(msg));
}
});
inFlightSplitRequests.clear();
}

@VisibleForTesting
int getInFlightSplitRequestCount() {
return inFlightSplitRequests.size();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import lombok.CustomLog;
import org.apache.commons.lang3.StringUtils;
import org.apache.pulsar.broker.loadbalance.extensions.channel.ServiceUnitState;
Expand Down Expand Up @@ -117,35 +118,35 @@ private void complete(String serviceUnit, Throwable ex) {
LatencyMetric.ASSIGN.endMeasurement(serviceUnit);
}

inFlightUnloadRequest.computeIfPresent(serviceUnit, (__, future) -> {
if (!future.isDone()) {
if (ex != null) {
future.completeExceptionally(ex);
} else {
future.complete(null);
}
}
return null;
});
CompletableFuture<Void> future = inFlightUnloadRequest.remove(serviceUnit);
if (future == null || future.isDone()) {
return;
}
if (ex != null) {
future.completeExceptionally(ex);
} else {
future.complete(null);
}
}

public CompletableFuture<Void> waitAsync(CompletableFuture<Void> eventPubFuture,
String bundle,
UnloadDecision decision,
long timeout,
TimeUnit timeoutUnit) {
return eventPubFuture.thenCompose(__ -> inFlightUnloadRequest.computeIfAbsent(bundle, ignore -> {
log.debug().attr("bundle", bundle).attr("timeout", timeout).attr("timeoutUnit", timeoutUnit)
.log("Handle unload bundle: , timeout");
CompletableFuture<Void> future = new CompletableFuture<>();
future.orTimeout(timeout, timeoutUnit).whenComplete((v, ex) -> {
if (ex != null) {
inFlightUnloadRequest.remove(bundle);
return eventPubFuture.thenCompose(__ -> {
CompletableFuture<Void> future = inFlightUnloadRequest.computeIfAbsent(bundle, ignore -> {
log.debug().attr("bundle", bundle).attr("timeout", timeout).attr("timeoutUnit", timeoutUnit)
.log("Handle unload bundle: , timeout");
return new CompletableFuture<Void>().orTimeout(timeout, timeoutUnit);
});
// Return the dependent stage so callers cannot observe completion before timeout cleanup finishes.
return future.whenComplete((v, ex) -> {
if (ex instanceof TimeoutException && inFlightUnloadRequest.remove(bundle, future)) {
log.warn().attr("bundle", bundle).exception(ex).log("Failed to wait unload for serviceUnit");
}
});
return future;
})).whenComplete((__, ex) -> {
}).whenComplete((__, ex) -> {
if (ex != null) {
counter.update(Failure, Unknown);
log.warn().attr("bundle", bundle).exception(ex).log("Failed to unload bundle");
Expand Down Expand Up @@ -207,12 +208,16 @@ public void handleEvent(String serviceUnit, ServiceUnitStateData data, Throwable

public void close() {
inFlightUnloadRequest.forEach((bundle, future) -> {
if (!future.isDone()) {
if (inFlightUnloadRequest.remove(bundle, future) && !future.isDone()) {
String msg = String.format("Unloading bundle: %s, but the unload manager already closed.", bundle);
log.warn(msg);
future.completeExceptionally(new IllegalStateException(msg));
}
});
inFlightUnloadRequest.clear();
}

@VisibleForTesting
int getInFlightUnloadRequestCount() {
return inFlightUnloadRequest.size();
}
}
Loading