akarnokd / akarnokd/RxJavaFiberInterop
Cancellation Notification for Long-Running Tasks in the FiberInterop.create
- 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.