akarnokd / akarnokd/RxJavaFiberInterop

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

Aperta
#98 2 commenti 0 reazioni 0 assegnatari Vedi su GitHub
enhancement
Lingua principale
Java
Stelle
40
Fork
7
Metriche di merge delle PR
Nessuna PR unita negli ultimi 30g

Descrizione

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.

Guida per i contributori

Nessuna guida per i contributori indicizzata per questo repository

Valutazione

Questa issue non è ancora stata valutata.

Ricevi le nuove issue nella tua casella

Un breve riepilogo di issue GitHub adatte ai principianti.