> Markdown version of [Producing Reactive Streams](https://vaadin.com/docs/next/building-apps/business-logic/background-jobs/interaction/reactive). Section index: [llms.txt](https://vaadin.com/docs/next/building-apps/llms.txt)

# Producing Reactive Streams

You can use `Flux` or `Mono` from [Reactor](https://projectreactor.io/) to allow your background jobs to interact with the user interface. Unlike [callbacks](https://vaadin.com/docs/next/building-apps/business-logic/background-jobs/interaction/callbacks.md) and [futures](https://vaadin.com/docs/next/building-apps/business-logic/background-jobs/interaction/futures.md), reactive streams work with both Flow views and [React views](https://vaadin.com/docs/next/building-apps/react.md). Reactor has an extensive API, which means you can do many things with it. This also means that it can be more difficult to learn than using callbacks or `CompletableFuture`.

This page is about returning the result of a background job to the user who started it. You can also use reactive streams to broadcast updates to all users. That use case is covered in [Consuming Reactive Streams](https://vaadin.com/docs/next/building-apps/server-push/reactive.md) and [Broadcasting to All Users](https://vaadin.com/docs/next/building-apps/server-push/updates.md#broadcasting-to-all-users).

If you’re new to reactive programming, you should read Reactor’s [Introduction to Reactive Programming](https://projectreactor.io/docs/core/release/reference/#intro-reactive) before continuing.

## <a id="returning-a-result"></a>Returning a Result

When you’re using Reactor, you can’t use the `@Async` annotation. Instead, you have to instruct your `Mono` or `Flux` to execute using the Spring `TaskExecutor`. Otherwise, your job executes in the thread that subscribes to the `Mono` or `Flux`.

For example, a background job that returns a string or an exception could be implemented like this:

```java
public Mono<String> startBackgroundJob() {
    return Mono.fromSupplier(this::doSomethingThatTakesALongTime)
               .subscribeOn(Schedulers.fromExecutor(taskExecutor));
}
```

If the `doSomethingThatTakesALongTime()` method throws an exception, the `Mono` terminates with an error.

To update the user interface, you have to subscribe to the `Mono` or `Flux`. For more information about how to do this in a Flow view, see the [Consuming Reactive Streams](https://vaadin.com/docs/next/building-apps/server-push/reactive.md) documentation page. For React views, see [Reactive Services](https://vaadin.com/docs/next/hilla/guides/reactive-services.md).

> **Important:** A React view calls the job through a [browser-callable service](https://vaadin.com/docs/next/building-apps/react/call-services.md), which can return a `Flux` but not a `Mono`. If your job returns a `Mono`, convert it to a `Flux` inside your `@BrowserCallable` service by calling the `Mono.flux()` method.

## <a id="reporting-progress"></a>Reporting Progress

If your background job only needs to report its progress without actually returning a result, you can return a `Flux<Double>`. Your job should then emit progress updates, and complete the stream when done. However, you may often want also to return a result. The simplest way to do this, which also works for React views, is to use the same stream for emitting both progress updates and the end result. The code may be a bit messy, but it works.

You first need to create a data type that can contain both progress updates and the result. For a job that returns a string, it could look like this:

```java
import org.jspecify.annotations.Nullable;

public record BackgroundJobOutput(
        @Nullable Double progressUpdate,
        @Nullable String result
) {
    public static BackgroundJobOutput progressUpdate(double progressUpdate) {
        return new BackgroundJobOutput(progressUpdate, null);
    }

    public static BackgroundJobOutput finished(String result) {
        return new BackgroundJobOutput(null, result);
    }
}
```

The two built-in methods, `progressUpdate()` and `finished()` make the code look better when it’s time to create instances of `BackgroundJobOutput`.

Next, you have to implement the background job like this:

```java
private String doSomethingThatTakesALongTime(Consumer<Double> onProgress) {
    ...
}

public Flux<BackgroundJobOutput> startBackgroundJob() {
    Sinks.Many<Double> progressUpdates = Sinks // (1)
        .many()
        .unicast()
        .onBackpressureError();

    var result = Mono // (2)
        .fromSupplier(() -> doSomethingThatTakesALongTime(
            progressUpdates::tryEmitNext
        ))
        .subscribeOn(Schedulers.fromExecutor(taskExecutor));

    return Flux.merge( // (3)
        progressUpdates.asFlux().map(BackgroundJobOutput::progressUpdate),
        result.map(BackgroundJobOutput::finished)
    );
}
```

1. Create a sink to which you can emit progress updates.

2. Create a `Mono` that emits the result of the background job.

3. Map both streams to `BackgroundJobOutput` and merge them.

When your user interface subscribes to this `Flux`, it needs to check the state of the returned `BackgroundJobOutput` objects. If `progressUpdate` contains a value, it should update the progress indicator. If `result` contains a value, though, the operation is finished.

## <a id="cancelling"></a>Cancelling

You can cancel a subscription to a `Flux` or `Mono` at any time. However, as with `CompletableFuture`, cancelling the subscription doesn’t stop the background job itself. To fix this, you need to tell the background job when it has been cancelled, so that it can stop. Continuing on the earlier example, adding support for cancelling could look like this:

```java
private String doSomethingThatTakesALongTime(
    Consumer<Double> onProgress,
    Supplier<Boolean> isCancelled) {
    ...
}

public Flux<BackgroundJobOutput> startBackgroundJob() {
    var cancelled = new AtomicBoolean(false);
    Sinks.Many<Double> progressUpdates = Sinks
        .many()
        .unicast()
        .onBackpressureError();

    var result = Mono
        .fromSupplier(() -> doSomethingThatTakesALongTime(
            progressUpdates::tryEmitNext, cancelled::get
        ))
        .doOnCancel(() -> cancelled.set(true))
        .subscribeOn(Schedulers.fromExecutor(taskExecutor));

    return Flux.merge(
        progressUpdates.asFlux().map(BackgroundJobOutput::progressUpdate),
        result.map(BackgroundJobOutput::finished)
    );
}
```

If the user interface cancels the subscription, the `cancelled` flag becomes `true`, and the job stops executing at its next iteration.
