构建响应式实时服务
开发端到端的响应式实时服务,充分发挥 Project Reactor 的能力。
构建响应式实时服务 是 CoddyKit 上的免费 WebSockets & Real-Time Systems with Spring 课时。 这是第 3 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 WebSockets & Real-Time Systems with Spring 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 WebSockets & Real-Time Systems with Spring 课程共包含 4 节课。
本课时的部分内容尚未翻译,以英文显示。
Reactive Real-Time Services
Welcome! In this lesson, we'll build end-to-end reactive real-time services using Spring WebFlux and Project Reactor.
Reactive services are excellent for handling many concurrent connections efficiently. They offer better scalability and responsiveness compared to traditional blocking approaches.
Project Reactor: Flux & Mono
At the heart of reactive programming in Spring is Project Reactor. It provides two key types for handling data streams:
Flux: Represents a stream of 0 to N items. Think of it as a publisher that can emit multiple values over time.Mono: Represents a stream of 0 to 1 item. Useful for operations that return a single result or no result (likevoid).
These types allow us to compose asynchronous operations in a clear and non-blocking way.
WebFlux WebSocket Handlers
Spring WebFlux uses the WebSocketHandler interface to manage WebSocket connections. Its main method, handle(), takes a WebSocketSession and returns a Mono.
This Mono signifies that the handling process is complete once the reactive stream it represents finishes. We can use Flux inside to send continuous messages.
Designing a Reactive Data Source
To build a real-time service, we need a source of data. Let's create a simple Flux that emits a message periodically. This simulates a real-time data feed, like a stock ticker or a sensor reading.
We'll use Flux.interval() to generate events and map() to transform them into useful messages.
Implementing a Ticker Service
Here's a basic WebSocketHandler that sends a 'tick' message every second. It uses the Flux.interval() we discussed.
The session.send() method takes a Flux to push data to the client.
import org.springframework.web.reactive.socket.WebSocketHandler;
import org.springframework.web.reactive.socket.WebSocketMessage;
import org.springframework.web.reactive.socket.WebSocketSession;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import java.time.Duration;
public class TimeTickerHandler implements WebSocketHandler {
@Override
public Mono<Void> handle(WebSocketSession session) {
// Send messages to the client
Flux<WebSocketMessage> output = Flux.interval(Duration.ofSeconds(1))
.map(i -> session.textMessage("Tick #" + i));
// Receive messages from the client (and ignore them for now)
// We use .then() to ensure the Mono<Void> completes only when the session closes.
Mono<Void> input = session.receive().then();
return session.send(output).and(input);
}
}
Full Runnable Ticker Service
To make our TimeTickerHandler runnable, we need a Spring Boot application. This example sets up the WebFlux server and registers our handler.
Access this via ws://localhost:8080/ticker in a WebSocket client (like Postman or a browser's DevTools console) to see it in action.
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.context.annotation.Bean;
import org.springframework.web.reactive.handler.SimpleUrlHandlerMapping;
import org.springframework.web.reactive.socket.WebSocketHandler;
import org.springframework.web.reactive.socket.WebSocketMessage;
import org.springframework.web.reactive.socket.WebSocketSession;
import org.springframework.web.reactive.socket.server.support.WebSocketHandlerAdapter;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import java.time.Duration;
import java.util.HashMap;
import java.util.Map;
@SpringBootApplication
public class ReactiveTickerApplication {
public static void main(String[] args) {
SpringApplication.run(ReactiveTickerApplication.class, args);
}
@Bean
public SimpleUrlHandlerMapping webSocketHandlerMapping(WebSocketHandler webSocketHandler) {
Map<String, WebSocketHandler> map = new HashMap<>();
map.put("/ticker", webSocketHandler);
return new SimpleUrlHandlerMapping(map, 1);
}
@Bean
public WebSocketHandler webSocketHandler() {
return new WebSocketHandler() {
@Override
public Mono<Void> handle(WebSocketSession session) {
// Send a 'tick' message every second
Flux<WebSocketMessage> output = Flux.interval(Duration.ofSeconds(1))
.map(i -> session.textMessage("Tick #" + i + " at " + System.currentTimeMillis()));
// Handle incoming messages (e.g., echo them back, or process commands)
// For this example, we'll just log and then complete the input stream
Mono<Void> input = session.receive()
.doOnNext(msg -> System.out.println("Received: " + msg.getPayloadAsText()))
.then(); // ensures the Mono completes after processing all incoming
return session.send(output).and(input);
}
};
}
@Bean
public WebSocketHandlerAdapter handlerAdapter() {
return new WebSocketHandlerAdapter();
}
}
Handling Client Input
Our previous ticker only sent data. To make it truly interactive, we can also process messages coming from the client.
The session.receive() method returns a Flux that represents incoming messages. You can subscribe to this Flux to react to client input, for example, by filtering, transforming, or using the data to control the output stream.
Error Handling in Reactive Streams
Errors can occur in any part of a reactive pipeline. Project Reactor provides operators to handle these gracefully, preventing your application from crashing:
onErrorResume(): Recovers from an error by switching to an alternative publisher.doOnError(): Performs a side-effect (like logging) when an error occurs, then re-throws it or completes.retry(): Retries the sequence if an error occurs.
Using these helps build robust real-time services that can recover from transient issues.
Backpressure Management
Backpressure is crucial for reactive systems. It's a mechanism where a consumer can signal to a producer that it's receiving data too quickly and needs the producer to slow down.
Project Reactor handles backpressure automatically. When a client can't keep up, the WebSocket connection might buffer messages or eventually close, but the server-side Flux won't overwhelm itself or the network.
Reactive Service Concepts
Which of the following are key characteristics of building reactive real-time services with Spring WebFlux and Project Reactor?
Recap: Reactive Real-Time
We've explored how to build reactive real-time services using Spring WebFlux and Project Reactor.
- We saw how
Fluxcan generate continuous data streams. - We implemented a
WebSocketHandlerto push these streams to clients. - We configured a basic Spring Boot application to host our reactive WebSocket endpoint.
- We touched upon error handling and backpressure, vital for robust systems.
These principles enable highly scalable and responsive real-time applications.
常见问题解答
「构建响应式实时服务」课时是免费的吗?
是的 — 「构建响应式实时服务」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 WebSockets & Real-Time Systems with Spring 课程的其余内容,请升级到 CoddyKit PRO。 WebSockets & Real-Time Systems with Spring 课程共包含 4 节课。
「构建响应式实时服务」这节课中我会学到什么?
开发端到端的响应式实时服务,充分发挥 Project Reactor 的能力。 你通过在浏览器中直接运行的动手代码来练习 WebSockets & Real-Time Systems with Spring,全天候 AI 导师会在你学习这节课的过程中回答你的问题。
学习 WebSockets & Real-Time Systems with Spring 需要有经验吗?
无需任何先前经验。CoddyKit 上的 WebSockets & Real-Time Systems with Spring 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 3 节课,共 4 节。
「构建响应式实时服务」课时需要多长时间?
大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。
我能在这节 WebSockets & Real-Time Systems with Spring 课中编写并运行代码吗?
能。每节 WebSockets & Real-Time Systems with Spring 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。