Micronaut Reactor

Integration between Micronaut and Reactor

Version: 4.3.1-SNAPSHOT

1 Introduction

Micronaut Reactor adds support for Project Reactor to a Micronaut 2.x application. If you are using Micronaut 1.x you do not need this module.

  • Converters for Reactor types so that Reactor types can be used in controllers and clients

  • Instrumentation for Reactor types so that tracing works with Reactor types

  • A version of the Http client that supports Reactor

2 Release History

For this project, you can find a list of releases (with release notes) here:

3 Breaking Changes

5.0.0

RxJava 2 Support Removed

Micronaut Reactor no longer supports RxJava 2. The RxJava 2 dependency has been removed and RxJava 3 are supported because Micronaut 5 drops RxJava 2 support.

If your application uses RxJava with this module, update imports from io.reactivex. to io.reactivex.rxjava3. and depend on io.reactivex.rxjava3:rxjava instead of io.reactivex.rxjava2:rxjava.

3.0.0

Reactor Http Client

The Reactor Http client has changed to return Mono<HttpResponse<O>> instead of Flux<HttpResponse<O>>. This is semantically more correct as there can only be a single response per request.

4 Quick Start

Add the following dependency to your Micronaut application:

implementation("io.micronaut.reactor:micronaut-reactor")
<dependency>
    <groupId>io.micronaut.reactor</groupId>
    <artifactId>micronaut-reactor</artifactId>
</dependency>
[tool.pyronaut.dependencies]
runtime = [
    "io.micronaut.reactor:micronaut-reactor",
]

To use the HTTP client add the following dependency:

implementation("io.micronaut.reactor:micronaut-reactor-http-client")
<dependency>
    <groupId>io.micronaut.reactor</groupId>
    <artifactId>micronaut-reactor-http-client</artifactId>
</dependency>
[tool.pyronaut.dependencies]
runtime = [
    "io.micronaut.reactor:micronaut-reactor-http-client",
]

To use the Reactor variation of the Micronaut HTTP client inject the ReactorHttpClient interface (or one the other variants). For example:

import io.micronaut.http.HttpRequest;
import io.micronaut.http.sse.Event;
import io.micronaut.reactor.http.client.ReactorHttpClient;
import io.micronaut.reactor.http.client.ReactorSseClient;
import io.micronaut.reactor.http.client.ReactorStreamingHttpClient;
import jakarta.inject.Inject;
import jakarta.inject.Singleton;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;

import java.net.URI;
import java.util.Map;

@Singleton
public class HeadlineClient {

    @Inject ReactorHttpClient httpClient; // regular client
    @Inject ReactorSseClient sseClient; // server sent events
    @Inject ReactorStreamingHttpClient streamingClient; // streaming

    public Mono<String> latest(URI server) {
        return httpClient.retrieve(HttpRequest.GET(server.resolve("/headlines/latest")));
    }

    public Flux<String> events(URI server) {
        return sseClient.eventStream(HttpRequest.GET(server.resolve("/headlines/events")), String.class)
                .map(Event::getData);
    }

    public Flux<Map<String, Object>> stream(URI server) {
        return streamingClient.jsonStream(HttpRequest.GET(server.resolve("/headlines/stream")));
    }
}
from typing import Annotated

from jakarta.inject import Inject, Singleton
from java.lang import String
from java.net import URI
from micronaut.http import HttpRequest
from micronaut.reactor.http.client import ReactorHttpClient, ReactorSseClient, ReactorStreamingHttpClient
from reactor.core.publisher import Flux, Mono

@Singleton
class HeadlineClient:

    http_client: Annotated[ReactorHttpClient, Inject]  # regular client
    sse_client: Annotated[ReactorSseClient, Inject]  # server sent events
    streaming_client: Annotated[ReactorStreamingHttpClient, Inject]  # streaming

    def latest(self, server: URI) -> Mono[str]:
        return self.http_client.retrieve(HttpRequest.GET(server.resolve("/headlines/latest")))

    def events(self, server: URI) -> Flux[str]:
        return (self.sse_client.eventStream(HttpRequest.GET(server.resolve("/headlines/events")), String)
                .map(lambda event: event.getData()))

    def stream(self, server: URI) -> Flux[dict[str, object]]:
        return self.streaming_client.jsonStream(HttpRequest.GET(server.resolve("/headlines/stream")))
import io.micronaut.http.HttpRequest
import io.micronaut.http.sse.Event
import io.micronaut.reactor.http.client.ReactorHttpClient
import io.micronaut.reactor.http.client.ReactorSseClient
import io.micronaut.reactor.http.client.ReactorStreamingHttpClient
import jakarta.inject.Inject
import jakarta.inject.Singleton
import reactor.core.publisher.Flux
import reactor.core.publisher.Mono
import java.net.URI

@Singleton
class HeadlineClient {

    @Inject lateinit var httpClient: ReactorHttpClient // regular client
    @Inject lateinit var sseClient: ReactorSseClient // server sent events
    @Inject lateinit var streamingClient: ReactorStreamingHttpClient // streaming

    fun latest(server: URI): Mono<String> =
        httpClient.retrieve(HttpRequest.GET<Any>(server.resolve("/headlines/latest")))

    fun events(server: URI): Flux<String> =
        sseClient.eventStream(HttpRequest.GET<Any>(server.resolve("/headlines/events")), String::class.java)
            .map { it.data }

    fun stream(server: URI): Flux<Map<String, Any>> =
        streamingClient.jsonStream(HttpRequest.GET<Any>(server.resolve("/headlines/stream")))
}
import io.micronaut.http.HttpRequest
import io.micronaut.http.sse.Event
import io.micronaut.reactor.http.client.ReactorHttpClient
import io.micronaut.reactor.http.client.ReactorSseClient
import io.micronaut.reactor.http.client.ReactorStreamingHttpClient
import jakarta.inject.Inject
import jakarta.inject.Singleton
import reactor.core.publisher.Flux
import reactor.core.publisher.Mono

@Singleton
class HeadlineClient {

    @Inject ReactorHttpClient httpClient // regular client
    @Inject ReactorSseClient sseClient // server sent events
    @Inject ReactorStreamingHttpClient streamingClient // streaming

    Mono<String> latest(URI server) {
        httpClient.retrieve(HttpRequest.GET(server.resolve("/headlines/latest")))
    }

    Flux<String> events(URI server) {
        sseClient.eventStream(HttpRequest.GET(server.resolve("/headlines/events")), String)
                .map(Event::getData)
    }

    Flux<Map<String, Object>> stream(URI server) {
        streamingClient.jsonStream(HttpRequest.GET(server.resolve("/headlines/stream")))
    }
}

5 Reactor context propagation

When using Micronaut Reactor within a reactive chain, add the following dependency to handle the Reactor Context Propagation.

implementation("io.micrometer:context-propagation")
<dependency>
    <groupId>io.micrometer</groupId>
    <artifactId>context-propagation</artifactId>
</dependency>
[tool.pyronaut.dependencies]
runtime = [
    "io.micrometer:context-propagation",
]

6 Repository

You can find the source code of this project in this repository: