Spring WebFlux 反应式编程与 Reactor 实战

深入探讨 Spring WebFlux 与 Project Reactor 反应式编程模型,涵盖 Mono/Flux 核心抽象、 常用操作符、背压策略、Netty 架构、WebClient、R2DBC 及 WebFlux 与 MVC 的 8 维度对比。 附带 12+ 完整 Java 代码示例与中文注释。

反应式编程已成为构建高并发、低延迟服务的核心范式。Spring WebFlux 是 Spring 5 引入的完全非阻塞 Web 框架,依托 Project Reactor 实现,基于响应式流规范(Reactive Streams Specification)定义了 PublisherSubscriberSubscriptionProcessor 四个接口。本文从理论到实战,系统讲解 WebFlux 核心原理与工程实践。

1. 反应式系统四要素

反应式宣言定义了四大特性:

  • 响应性:系统以一致、及时的方式响应,维持可预测的响应时间。
  • 弹性:系统在故障时保持响应,通过隔离、复制、降级实现自愈。
  • 弹性(Elastic):系统根据负载自动扩缩容,利用水平扩展应对流量激增。
  • 消息驱动:组件通过异步消息传递实现松耦合。

Spring WebFlux 正是上述理念的 Java 落地实现。响应式流规范确保跨库兼容性。

import org.reactivestreams.*;

/** 自定义 Publisher:模拟数据发布 */
public class SimplePublisher implements Publisher<Integer> {
    private final Integer[] data;
    public SimplePublisher(Integer[] data) { this.data = data; }

    @Override
    public void subscribe(Subscriber<? super Integer> subscriber) {
        subscriber.onSubscribe(new Subscription() {
            private int index = 0;
            private boolean cancelled = false;
            @Override
            public void request(long n) {
                for (int i = 0; i < n && index < data.length; i++) {
                    if (cancelled) return;
                    subscriber.onNext(data[index++]);
                }
                if (index >= data.length) subscriber.onComplete();
            }
            @Override
            public void cancel() { cancelled = true; }
        });
    }
}

设计洞察:反应式系统将数据流与时间解耦,订阅者决定消费时机,天然适合微服务异步编排。

2. Reactor 基础:Mono vs Flux

Project Reactor 提供两个核心发布者实现。

2.1 Mono — 零个或一个元素的异步序列

Mono<T> 表示异步计算的最终结果,可能是 0 个(empty)、1 个(just)或异常(error)。

import reactor.core.publisher.Mono;
import java.time.Duration;
import java.util.Optional;

public class MonoDemo {
    /** Mono.just:直接包装已存在的值 */
    public Mono<String> getUserName() { return Mono.just("Alice"); }

    /** Mono.fromCallable:延迟执行阻塞操作 */
    public Mono<String> fetchFromDB() {
        return Mono.fromCallable(() -> { Thread.sleep(100); return "User:42"; });
    }

    /** Mono.empty 与 Mono.error */
    public Mono<Optional<String>> findById(String id) {
        if (id == null || id.isBlank()) return Mono.error(new IllegalArgumentException("ID 不能为空"));
        if ("notfound".equals(id)) return Mono.empty();
        return Mono.just(Optional.of("User[" + id + "]"));
    }

    /** Mono.delay:延迟定时 */
    public Mono<Long> delayedSignal() {
        return Mono.delay(Duration.ofSeconds(2))
                .doOnNext(t -> System.out.println("延迟任务触发,tick=" + t));
    }
}

2.2 Flux — 零到多个元素的异步序列

Flux<T> 表示包含 0 到 N 个元素的响应式流,最常用。

import reactor.core.publisher.Flux;
import reactor.core.scheduler.Schedulers;
import java.time.Duration;
import java.util.List;
import java.util.stream.IntStream;

public class FluxDemo {
    /** Flux.just:有限流 */
    public Flux<String> fruitStream() { return Flux.just("Apple", "Banana", "Cherry"); }

    /** Flux.fromIterable:从集合创建 */
    public Flux<Integer> fromList(List<Integer> numbers) { return Flux.fromIterable(numbers).log(); }

    /** Flux.range:连续整数序列 */
    public Flux<Integer> rangeStream() { return Flux.range(1, 10); }

    /** Flux.interval:定时发射 */
    public Flux<Long> ticker() { return Flux.interval(Duration.ofMillis(500)).take(10); }

    /** Flux.create:桥接回调式 API */
    public Flux<String> createStream() {
        return Flux.create(emitter -> {
            IntStream.rangeClosed(1, 5).forEach(i -> emitter.next("事件-" + i));
            emitter.complete();
        }, FluxSink.OverflowStrategy.BUFFER);
    }

    /** Flux.generate:有状态同步生成 */
    public Flux<Long> fibonacci() {
        return Flux.generate(() -> new long[]{0, 1},
            (state, sink) -> {
                long next = state[0] + state[1];
                state[0] = state[1]; state[1] = next;
                sink.next(state[0]);
                if (state[0] > 100) sink.complete();
                return state;
            });
    }
}

2.3 订阅与执行模型

Reactor 的核心是惰性执行 —— 只有调用 .subscribe() 或由 WebFlux 框架驱动时,数据流才会运转。

// 冷发布者:每次订阅重新执行
Flux<String> coldFlux = Flux.defer(() -> {
    System.out.println("冷流数据源被调用"); return Flux.just("a", "b", "c");
});
coldFlux.subscribe(System.out::println);
coldFlux.subscribe(System.out::println);  // 再次触发

// 热发布者:共享同一数据源
Flux<Long> hotFlux = Flux.interval(Duration.ofSeconds(1)).take(5).share();
hotFlux.subscribe(v -> System.out.println("订阅者1: " + v));
Thread.sleep(2500);
hotFlux.subscribe(v -> System.out.println("订阅者2: " + v));

生产环境勿用 Thread.sleep() 阻塞 Reactor 线程,应使用 Mono.delay()StepVerifier.withVirtualTime() 测试。

3. 操作符:map / flatMap / filter / switchMap / zip / merge

Reactor 提供 200+ 操作符,以下是六类核心操作符。

3.1 map — 一对一同步转换

/** map:同步转换,元素数量不变 */
public Flux<String> mapExample() {
    return Flux.range(1, 5).map(i -> "编号-" + i).map(String::toUpperCase);
}

3.2 flatMap — 一对多异步展开

import reactor.core.publisher.Mono;
import java.time.Duration;

/** flatMap:将元素映射为 Publisher 后异步合并,并发度高,顺序不保证 */
public Flux<String> flatMapExample() {
    return Flux.just("user:1", "user:2", "user:3")
            .flatMap(id -> Mono.just("详情[" + id + "]")
                    .delayElement(Duration.ofMillis((long)(Math.random() * 100))));
}

/** flatMapSequential:按源顺序合并 */
public Flux<String> flatMapSequentialExample() {
    return Flux.range(1, 5).flatMapSequential(i ->
            Mono.just("结果" + i).delayElement(Duration.ofMillis(50)));
}

3.3 filter — 条件过滤

/** filter:保留满足条件的元素 */
public Flux<Integer> filterExample() {
    return Flux.range(1, 20).filter(i -> i % 2 == 0).filter(i -> i > 5).take(3);
}

3.4 switchMap — 仅保留最新流

/** switchMap:新元素到达时取消旧流,适合搜索框自动补全 */
public Flux<String> searchAutoComplete(Flux<String> inputStream) {
    return inputStream.filter(kw -> kw.length() >= 2)
            .switchMap(kw -> performSearch(kw)
                    .timeout(Duration.ofSeconds(3))
                    .onErrorResume(e -> Flux.just("搜索超时:" + kw)));
}
private Flux<String> performSearch(String kw) {
    return Flux.just(kw + "-结果1", kw + "-结果2").delayElements(Duration.ofMillis(300));
}

3.5 zip — 多流对齐合并

import reactor.util.function.Tuple2;

/** zip:多流按索引一对一合并,最短流决定长度 */
public Mono<UserProfile> getUserProfile(String userId) {
    Mono<User> userMono = userService.findById(userId);
    Mono<List<Order>> ordersMono = orderService.findByUserId(userId);
    Mono<VipStatus> vipMono = vipService.getStatus(userId);
    return Mono.zip(userMono, ordersMono, vipMono)
            .map(t -> new UserProfile(t.getT1(), t.getT2(), t.getT3()));
}

/** zipWith:两个流对齐 */
public Flux<Tuple2<String, Integer>> zipWithExample() {
    Flux<String> names = Flux.just("Alice", "Bob", "Carol");
    Flux<Integer> scores = Flux.just(90, 85, 95, 70);  // 第四个被忽略
    return names.zipWith(scores);
}

3.6 merge & concat — 多流合并策略

/** merge:交错合并,同时订阅 */
public Flux<String> mergeExample() {
    Flux<String> a = Flux.interval(Duration.ofMillis(300)).map(i -> "A" + i).take(3);
    Flux<String> b = Flux.interval(Duration.ofMillis(500)).map(i -> "B" + i).take(3);
    return Flux.merge(a, b);
}

/** concat:顺序合并,先消费完第一个 */
public Flux<String> concatExample() {
    Flux<String> s1 = Flux.just("1-a", "1-b").delayElements(Duration.ofMillis(100));
    Flux<String> s2 = Flux.just("2-a", "2-b");
    return Flux.concat(s1, s2);
}

选择速查:同步转 map,异步转 flatMap,顺序用 concatMap,最新用 switchMap,多流聚合用 zip

4. 背压策略:BUFFER / DROP / LATEST

背压解决生产速度远超消费速度时的内存与线程问题。

import reactor.core.publisher.Flux;
import reactor.core.publisher.FluxSink;
import reactor.core.scheduler.Schedulers;
import java.time.Duration;

public class BackpressureDemo {
    /** BUFFER:缓冲到队列,不丢数据,内存风险高 */
    public Flux<Integer> bufferStrategy() {
        return Flux.range(1, 1_000_000).onBackpressureBuffer(1000);
    }

    /** DROP:丢弃无法处理的元素,适合日志/监控 */
    public Flux<Long> dropStrategy() {
        return Flux.interval(Duration.ofMillis(1))
                .onBackpressureDrop(d -> System.out.println("丢弃:" + d))
                .publishOn(Schedulers.boundedElastic(), 1);
    }

    /** LATEST:只保留最新元素,低内存,适合仪表盘 */
    public Flux<Double> latestStrategy() {
        return Flux.<Double>create(e -> new Thread(() -> {
            for (int i = 0; i < 10000; i++) e.next(Math.random() * 100);
            e.complete();
        }).start(), FluxSink.OverflowStrategy.LATEST)
        .sample(Duration.ofMillis(100));
    }

    /** ERROR:超出处理能力即报错 */
    public Flux<Integer> errorStrategy() { return Flux.range(1, 100).onBackpressureError(); }
}
背压策略数据丢失内存风险适用场景
BUFFER高(无限缓冲区)金融交易、订单处理等不能丢数据场景
DROP日志流、监控指标、采样统计
LATEST是(中间值)极低实时仪表盘、传感器读数、股价推送
ERROR严格速率匹配场景,不匹配即失败

5. WebFlux 架构:Netty vs Tomcat、RouterFunction vs @Controller

5.1 Netty 运行时配置

WebFlux 默认使用 Netty(NIO/EventLoop),也支持 Servlet 3.1+ 容器。

import org.springframework.boot.web.embedded.netty.NettyReactiveWebServerFactory;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import io.netty.channel.ChannelOption;

/** 自定义 Netty:调整 EventLoop 与连接参数 */
@Configuration
public class NettyConfig {
    @Bean
    public NettyReactiveWebServerFactory nettyFactory() {
        NettyReactiveWebServerFactory f = new NettyReactiveWebServerFactory();
        f.addServerCustomizers(s -> s
                .option(ChannelOption.SO_BACKLOG, 1024)
                .option(ChannelOption.CONNECT_TIMEOUT_MILLIS, 5000)
                .childOption(ChannelOption.TCP_NODELAY, true)
                .childOption(ChannelOption.SO_KEEPALIVE, true));
        return f;
    }
}

5.2 RouterFunction 函数式路由

import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.web.reactive.function.server.*;
import reactor.core.publisher.Mono;
import static org.springframework.web.reactive.function.server.RequestPredicates.*;
import static org.springframework.web.reactive.function.server.RouterFunctions.route;

/** 函数式路由:路由与处理器分离,易于测试,无反射开销 */
@Configuration
public class UserRouter {
    @Bean
    public RouterFunction<ServerResponse> userRoutes(UserHandler h) {
        return route().GET("/users", accept(APPLICATION_JSON), h::listUsers)
                .GET("/users/{id}", h::getUser)
                .POST("/users", contentType(APPLICATION_JSON), h::createUser)
                .PUT("/users/{id}", h::updateUser)
                .DELETE("/users/{id}", h::deleteUser)
                .build();
    }
}

@Component
public class UserHandler {
    private final UserRepository repo;
    public UserHandler(UserRepository repo) { this.repo = repo; }

    public Mono<ServerResponse> listUsers(ServerRequest req) {
        return ServerResponse.ok().contentType(APPLICATION_JSON).body(repo.findAll(), User.class);
    }
    public Mono<ServerResponse> getUser(ServerRequest req) {
        return repo.findById(req.pathVariable("id"))
                .flatMap(u -> ServerResponse.ok().bodyValue(u))
                .switchIfEmpty(ServerResponse.notFound().build());
    }
    public Mono<ServerResponse> createUser(ServerRequest req) {
        return req.bodyToMono(User.class).flatMap(repo::save)
                .flatMap(saved -> ServerResponse.created(URI.create("/users/" + saved.id())).bodyValue(saved));
    }
}

5.3 @Controller 注解式

import org.springframework.web.bind.annotation.*;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import java.time.Duration;

/** 注解式 @Controller:学习成本低,与传统 MVC 风格一致 */
@RestController
@RequestMapping("/api/v2/users")
public class UserController {
    private final UserService svc;
    public UserController(UserService svc) { this.svc = svc; }

    @GetMapping public Flux<User> listAll() { return svc.findAll(); }

    @GetMapping("/{id}") public Mono<User> getById(@PathVariable String id) {
        return svc.findById(id);
    }

    @PostMapping public Mono<User> create(@RequestBody Mono<User> userMono) {
        return userMono.flatMap(svc::save);  // 请求体本身就是 Mono
    }

    @GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
    public Flux<User> streamUsers() {
        return svc.streamUpdates().delayElements(Duration.ofSeconds(1));
    }
}

选型建议:新项目优先 RouterFunction;团队有 MVC 经验则选择 @Controller。二者可混用。

6. WebClient 异步 HTTP

WebClient 是 Spring 5 的非阻塞 HTTP 客户端,替代 RestTemplate

import org.springframework.web.reactive.function.client.WebClient;
import org.springframework.web.reactive.function.client.WebClientResponseException;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.netty.http.client.HttpClient;
import reactor.util.retry.Retry;
import io.netty.channel.ChannelOption;
import java.time.Duration;

@Component
public class HttpClientService {
    private final WebClient client;

    public HttpClientService(WebClient.Builder b) {
        this.client = b.baseUrl("https://api.example.com")
                .defaultHeader(HttpHeaders.CONTENT_TYPE, MediaType.APPLICATION_JSON_VALUE)
                .codecs(c -> c.defaultCodecs().maxInMemorySize(2 * 1024 * 1024))
                .clientConnector(new ReactorClientHttpConnector(
                        HttpClient.create()
                                .option(ChannelOption.CONNECT_TIMEOUT_MILLIS, 5000)
                                .responseTimeout(Duration.ofSeconds(10))))
                .build();
    }

    /** GET 单个资源 */
    public Mono<User> getUserById(String id) {
        return client.get().uri("/users/{id}", id)
                .retrieve()
                .onStatus(HttpStatus::is4xxClientError,
                        r -> Mono.error(new UserNotFoundException(id)))
                .bodyToMono(User.class)
                .retryWhen(Retry.backoff(3, Duration.ofMillis(100))
                        .filter(t -> t instanceof IOException))
                .timeout(Duration.ofSeconds(5));
    }

    /** GET 列表资源 */
    public Flux<Order> getUserOrders(String userId) {
        return client.get().uri(u -> u.path("/users/{id}/orders")
                        .queryParam("status", "PAID").queryParam("page", 0).queryParam("size", 20)
                        .build(userId))
                .retrieve().bodyToFlux(Order.class);
    }

    /** POST 提交 */
    public Mono<CreateResult> createOrder(OrderRequest req) {
        return client.post().uri("/orders").bodyValue(req)
                .retrieve().bodyToMono(CreateResult.class);
    }

    /** zip 并发聚合:等待多个 HTTP 结果 */
    public Mono<Dashboard> fetchDashboard(String userId) {
        Mono<User> userMono = getUserById(userId);
        Mono<Wallet> walletMono = client.get().uri("/wallets/{uid}", userId)
                .retrieve().bodyToMono(Wallet.class);
        Mono<NotificationSummary> notifMono = client.get().uri("/notifications/{uid}/summary", userId)
                .retrieve().bodyToMono(NotificationSummary.class);
        return Mono.zip(userMono, walletMono, notifMono)
                .map(t -> new Dashboard(t.getT1(), t.getT2(), t.getT3()));
    }
}

优化要点:WebClient 底层复用 Netty 连接池(默认 500 连接),高并发应监控连接池指标。中间链中避免调用 block()subscribe(),否则会阻塞 EventLoop。

7. R2DBC 反应式数据库

R2DBC(Reactive Relational Database Connectivity)替代 JDBC,解决阻塞 I/O 问题。

# application.yml
spring:
  r2dbc:
    url: r2dbc:postgresql://localhost:5432/webflux_db
    username: app_user
    password: secret
    pool:
      initial-size: 10
      max-size: 50
      max-idle-time: 30m
import org.springframework.data.r2dbc.repository.Query;
import org.springframework.data.repository.reactive.ReactiveCrudRepository;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;

/** ReactiveCrudRepository:提供完整的 CRUD 反应式接口 */
public interface UserRepository extends ReactiveCrudRepository<User, Long> {
    Flux<User> findByStatusAndCreatedAtAfter(String status, LocalDateTime after);

    @Query("SELECT * FROM users WHERE email = :email LIMIT 1")
    Mono<User> findByEmail(String email);
}
import org.springframework.r2dbc.core.DatabaseClient;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;

/** DatabaseClient:复杂 SQL 原生操作 */
@Component
public class OrderCustomRepository {
    private final DatabaseClient db;
    public OrderCustomRepository(DatabaseClient db) { this.db = db; }

    public Flux<OrderSummary> findOrderSummariesByUser(String userId) {
        String sql = """
            SELECT o.id, o.total_amount, o.created_at, COUNT(oi.id) as item_count
            FROM orders o LEFT JOIN order_items oi ON o.id = oi.order_id
            WHERE o.user_id = :userId GROUP BY o.id ORDER BY o.created_at DESC""";
        return db.sql(sql).bind("userId", userId)
                .map((row, md) -> new OrderSummary(
                        row.get("id", Long.class), row.get("total_amount", BigDecimal.class),
                        row.get("created_at", LocalDateTime.class), row.get("item_count", Integer.class)))
                .all();
    }

    /** 事务:@Transactional 在反应式中由 ReactiveTransactionManager 驱动 */
    @Transactional
    public Mono<Void> transfer(long fromId, long toId, BigDecimal amount) {
        return db.sql("UPDATE accounts SET balance = balance - :a WHERE id = :f")
                .bind("a", amount).bind("f", fromId).fetch().rowsUpdated()
                .flatMap(u -> u == 0 ? Mono.error(new InsufficientBalanceException()) :
                        db.sql("UPDATE accounts SET balance = balance + :a WHERE id = :t")
                                .bind("a", amount).bind("t", toId).fetch().rowsUpdated().then());
    }
}

迁移注意:R2DBC 无 @OneToMany 懒加载,需显式 JOIN;无二级缓存,依赖外部缓存;连接是非阻塞的,事务方法必须返回 Mono/Flux

8. WebFlux vs MVC 8 维度对比

对比维度Spring WebFluxSpring MVC
运行时模型EventLoop,少量线程处理海量连接每请求一线程,线程池调度
I/O 模型完全非阻塞 NIO同步阻塞 I/O(传统 Servlet)
并发能力单机 10K-100K 并发单机 1K-5K 并发(线程栈受限)
延迟特性长尾延迟稳定,上下文切换极少高并发时线程争抢导致延迟抖动
CPU 利用率高,线程空转少线程阻塞时 CPU 资源浪费
编程模型RouterFunction 或 @Controller注解驱动 @Controller
数据访问R2DBC、MongoDB Reactive、Redis ReactiveJDBC、JPA/Hibernate、MyBatis
生态成熟度较新,部分中间件缺少反应式驱动极其成熟,几乎所有库兼容

选型决策树

  • QPS < 1000 且延迟 > 100ms → MVC(开发效率优先)
  • 需与大量传统 JDBC 集成 → MVC(生态兼容优先)
  • 大量外部 HTTP 聚合调用 → WebFlux(网络并发优势)
  • 高扩展预留 → WebFlux

9. 部署与监控

9.1 生产配置

# application-prod.yml
server:
  netty:
    connection-timeout: 2s
    idle-timeout: 60s
spring:
  webflux:
    base-path: /api
  codec:
    max-in-memory-size: 1MB
logging:
  level:
    reactor.netty: WARN
    org.springframework.web.reactive: WARN

9.2 自定义 Micrometer 指标

import io.micrometer.core.instrument.*;
import org.springframework.stereotype.Component;
import reactor.core.publisher.Mono;
import java.util.concurrent.TimeUnit;

@Component
public class MonitoredOrderService {
    private final Timer timer;
    private final MeterRegistry registry;

    public MonitoredOrderService(MeterRegistry r) {
        this.registry = r;
        this.timer = Timer.builder("order.processing.time")
                .publishPercentiles(0.5, 0.95, 0.99).register(r);
    }

    public Mono<Order> processOrder(OrderRequest req) {
        return Mono.fromCallable(() -> System.nanoTime())
                .flatMap(start -> doProcess(req)
                        .doOnSuccess(o -> {
                            timer.record(System.nanoTime() - start, TimeUnit.NANOSECONDS);
                            registry.counter("order.success").increment();
                        })
                        .doOnError(e -> registry.counter("order.failure",
                                "exception", e.getClass().getSimpleName()).increment()));
    }
    private Mono<Order> doProcess(OrderRequest req) { return Mono.just(new Order()); }
}

9.3 常用监控指标

指标类型Micrometer 指标名说明
HTTP 请求耗时http.server.requests按 URI 和方法聚合的延迟直方图
R2DBC 连接池r2dbc.pool.active.connections活跃连接数
Netty EventLoopreactor.netty.eventloop.pending.tasks积压任务数
JVM 内存jvm.memory.used堆内存与直接内存使用

常见 FAQ

Q1:WebFlux 中调用了阻塞方法(如 JDBC),会怎样?
A:阻塞调用会卡住 EventLoop 线程,导致整个工作线程池停滞。解决方案:迁移到 R2DBC;或使用 Schedulers.boundedElastic() 将阻塞任务隔离到独立线程池,通过 subscribeOn(Schedulers.boundedElastic()) 包装老旧代码。

Q2:Mono/Flux 什么时候调用 subscribe 最合适?
A:在 WebFlux 中,框架层(DispatcherHandler)自动 subscribe,开发者不应手动调用。仅在单元测试、控制台程序或桥接到非反应式代码时显式 subscribe,生产代码中手动 subscribe 会丢失错误处理和生命周期管理。

Q3:WebFlux 支持文件上传下载吗?
A:完全支持。上传使用 Flux<Part> 处理 multipart/form-data;下载使用 DataBufferUtils 流式写入。大文件(>1GB)使用 ResourceRegion 分块传输,避免加载到内存。

Q4:反应式中的异常如何处理?
A:多级机制:onErrorReturn 降级值;onErrorResume 按异常类型切换备用流;retryWhen 指数退避重试;onErrorContinue 跳过错误继续处理。doOnError 仅用于副作用(日志),不会恢复流。

Q5:WebFlux 适合所有 HTTP 服务吗?
A:并非银弹。CRUD 为主、并发 < 1K QPS 的内部服务用 MVC 更稳定。WebFlux 优势在:网关/边缘服务(下游聚合)、实时推送(SSE/WebSocket)、高并发对外 API、资源受限环境(云函数/小容器)。

总结

Spring WebFlux 与 Project Reactor 为 Java 生态系统带来了媲美 Node.js/Go 的异步处理能力,同时保持 JVM 的类型安全与生态深度。掌握以下核心点是构建高并发服务的基础:

  • Mono = 异步的 0/1 值;Flux = 异步的 0/N 流
  • map 同步转;flatMap 异步展开;switchMap 取最新;zip 等全部;merge/concat 合多流
  • 背压四策略:BUFFER 不丢数据、DROP 低内存、LATEST 最新值、ERROR 严格匹配
  • Netty 默认运行时,任何位置避免阻塞
  • R2DBC 替代 JDBC,实现端到端非阻塞

建议团队从边缘服务(通知、网关聚合层)开始试点,逐步积累调优经验,再向核心系统推广。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「java」更多文章

  1. Spring Cloud 微服务全栈实践
  2. Spring Security 6.x 与 OAuth2/JWT 安全认证实战
  3. Spring Data JPA 高级指南:关联映射、N+1 与性能优化