akarnokd / akarnokd/RxJavaFiberInterop

Cancellation Notification for Long-Running Tasks in the FiberInterop.create

Offen
#98 2 Kommentare 0 Reaktionen 0 zugewiesene Personen Auf GitHub ansehen
enhancement
Vorherrschende Sprache
Java
Sterne
40
Forks
7
PR-Merge-Kennzahlen
Keine gemergten PRs in 30 T.

Beschreibung

When using RxJava 3 in combination with the FiberInterop library, there is a challenge in handling cancellations for long-running tasks. Specifically, the current implementation does not provide a mechanism for a long-running function to be notified of a cancellation until `emit` is invoked. This causes an issue where a function like `computeValue` cannot be aware of a cancellation event, leading to potential inefficiencies and unnecessary computations.

```java
var executor = Executors.newVirtualThreadPerTaskExecutor();
var scheduler = Schedulers.from(executor);
var cancelled = new AtomicBoolean(false);

var flow = FiberInterop.create(emitter -> {
var value = computeValue(cancelled);
emitter.emit(value);
}, executor).timeout(2, TimeUnit.SECONDS, scheduler);

flow.blockingForEach(value -> {
System.out.println(value);
});
```

In the above code, `computeValue` is a long-running function that takes a `cancelled` flag to potentially halt its execution if the flow is cancelled. However, the `cancelled` flag cannot be updated until the `emit` is called, making it ineffective in stopping the computation early.

To solve this problem, introducing a mechanism like registering `onCancel` callback or setting an interrupt flag on the virtual thread could effectively signal cancellation to the running task.
Here are some examples of how it might look:

```java
var flow = FiberInterop.create(emitter -> {
emitter.onCancel(() -> cancelled.set(true));
var value = computeValue(cancelled);
if (!cancelled.get()) {
emitter.emit(value);
}
}, executor).timeout(2, TimeUnit.SECONDS, scheduler);
```

```java
var flow = FiberInterop.create(emitter -> {
var value = computeValue();
emitter.emit(value);
}, executor, interruptWhenCancelled).timeout(2, TimeUnit.SECONDS, scheduler);
```

This approach ensures that the `computeValue` function can be notified of the cancellation event promptly and can stop its execution accordingly.

Beitragsleitfaden

Für dieses Repository ist kein Beitragsleitfaden indexiert

Bewertung

Dieses Issue wurde noch nicht bewertet.

Neue Issues direkt in Ihr Postfach

Eine kurze Übersicht über anfängerfreundliche GitHub-Issues.