You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用Project Reactor计算两个有序冷Flux的集合差并处理背压?

实现支持背压的有序Flux集合差计算

要实现两个冷发布者有序Flux的集合差(B - A,即B中存在但A中不存在的元素),同时避免内存耗尽,核心是利用Reactor的背压机制控制B元素的暂存缓冲区大小,当缓冲区满时暂停B的生产。以下是具体实现方案:

核心思路

  1. 标记元素来源:给A、B Flux的元素打上来源标记,方便后续区分处理。
  2. 有界暂存缓冲区:用有界阻塞队列暂存B中未在A中出现过的元素,队列满时阻塞B的生产,触发背压。
  3. 流式跟踪A元素:用线程安全的集合记录A中已处理的元素,实时过滤B的元素。
  4. 协调流的生命周期:当A流结束时,输出缓冲区中剩余的B元素(即集合差结果),同时处理取消、错误等边界情况。

代码实现

1. 元素标记类

首先定义一个简单的标记类,区分元素来自A还是B:

import reactor.core.publisher.Flux;
import reactor.core.publisher.BaseSubscriber;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.ConcurrentHashMap;
import java.util.Set;

static class TaggedElement<T> {
    enum Source { A, B }
    final T value;
    final Source source;

    TaggedElement(T value, Source source) {
        this.value = value;
        this.source = source;
    }
}

2. 集合差计算方法

实现核心逻辑,通过Flux.push手动控制流的生产,结合BaseSubscriber处理B流的背压:

public static <T> Flux<T> orderedSetDifference(Flux<T> fluxA, Flux<T> fluxB, int bufferCapacity) {
    return Flux.push(sink -> {
        // 有界队列暂存B元素,满时触发背压
        BlockingQueue<T> bBuffer = new LinkedBlockingQueue<>(bufferCapacity);
        // 线程安全集合记录A中已出现的元素
        Set<T> seenInA = ConcurrentHashMap.newKeySet();

        // 订阅A流:记录元素并清理缓冲区中匹配的B元素
        Disposable aDisposable = fluxA.subscribe(
            a -> {
                seenInA.add(a);
                // 移除缓冲区中所有与当前A元素匹配的B元素
                bBuffer.remove(a);
            },
            sink::error,
            () -> {
                // A流结束,输出缓冲区剩余的B元素
                while (!bBuffer.isEmpty()) {
                    sink.next(bBuffer.poll());
                }
                sink.complete();
            }
        );

        // 订阅B流:手动处理背压,仅当缓冲区有空间时请求下一个元素
        BaseSubscriber<T> bSubscriber = new BaseSubscriber<T>() {
            @Override
            protected void hookOnSubscribe(Subscription subscription) {
                // 初始请求1个元素
                request(1);
            }

            @Override
            protected void hookOnNext(T b) {
                try {
                    if (seenInA.contains(b)) {
                        // 元素已在A中出现,直接丢弃,继续请求下一个
                        request(1);
                        return;
                    }
                    // 尝试将元素加入缓冲区,队列满时阻塞,触发背压
                    bBuffer.put(b);
                    // 加入成功,请求下一个元素
                    request(1);
                } catch (InterruptedException e) {
                    Thread.currentThread().interrupt();
                    cancel();
                    sink.error(e);
                }
            }

            @Override
            protected void hookOnError(Throwable throwable) {
                sink.error(throwable);
            }

            @Override
            protected void hookOnComplete() {
                // B流结束,等待A流处理完剩余逻辑
            }
        };
        Disposable bDisposable = fluxB.subscribe(bSubscriber);

        // 当结果流被取消时,及时取消A、B流的订阅,避免资源泄漏
        sink.onCancel(() -> {
            aDisposable.dispose();
            bDisposable.dispose();
        });
    });
}

关键细节说明

  • 背压控制:通过LinkedBlockingQueue的put方法实现阻塞,当缓冲区满时,B流的BaseSubscriber会暂停请求下一个元素,B的生产者会因背压机制停止发送,避免内存耗尽。
  • 冷发布者适配:每次订阅结果流时,A、B流都会重新订阅,符合冷发布者的特性。
  • 线程安全:使用ConcurrentHashMap.newKeySet()和LinkedBlockingQueue保证多线程环境下的安全操作。
  • 结果顺序:输出的差集元素保持B流原有的顺序,仅移除在A中出现过的元素。

使用示例

public static void main(String[] args) {
    Flux<Integer> fluxA = Flux.just(2, 4, 6);
    Flux<Integer> fluxB = Flux.just(1, 2, 3, 4, 5, 6, 7);

    // 计算B - A,缓冲区容量设为3
    orderedSetDifference(fluxA, fluxB, 3)
        .subscribe(System.out::println);
    // 输出结果:1, 3, 5, 7
}

内容的提问来源于stack exchange,提问作者Theozaurus

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.19 13:32:35