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

Java Reactor:如何从stdin生成Flux并实现异步读取发布者?

如何用Java Reactor从stdin生成异步Flux?

好问题!要基于Reactor实现从stdin异步读取输入并生成Flux,核心是解决阻塞IO与Reactor非阻塞模型的适配问题——毕竟stdin的readLine()是阻塞操作,不能直接放在Reactor的事件循环线程里执行。下面我给你两种实现方式:一种是快速上手的简便方案,另一种是手动实现Publisher的底层方案,方便你理解原理。

方法一:用Flux.create()快速实现

Reactor提供的Flux.create()是封装阻塞IO最省心的方式,它允许你在单独的线程里执行阻塞操作,然后通过sink向订阅者发送数据、错误或完成信号。

完整代码示例

import reactor.core.publisher.Flux;
import reactor.core.scheduler.Schedulers;
import java.io.BufferedReader;
import java.io.InputStreamReader;

public class StdinFluxDemo {
    public static Flux<String> createStdinFlux() {
        return Flux.create(sink -> {
            try (BufferedReader reader = new BufferedReader(new InputStreamReader(System.in))) {
                // 将阻塞的读取操作放到boundedElastic线程池,避免阻塞Reactor的非阻塞线程
                Schedulers.boundedElastic().schedule(() -> {
                    try {
                        String inputLine;
                        while ((inputLine = reader.readLine()) != null) {
                            // 检查订阅者是否已取消,避免无效发送
                            if (sink.isCancelled()) {
                                break;
                            }
                            sink.next(inputLine);
                        }
                        // 当stdin关闭(比如按下Ctrl+D)时,结束Flux
                        sink.complete();
                    } catch (Exception e) {
                        sink.error(e);
                    }
                });
            } catch (Exception e) {
                sink.error(e);
            }
        });
    }

    public static void main(String[] args) {
        createStdinFlux()
            .subscribe(
                line -> System.out.println("你输入了: " + line),
                error -> System.err.println("读取出错: " + error.getMessage()),
                () -> System.out.println("输入流已关闭")
            );

        // 保持主线程存活,直到输入结束
        try {
            Thread.currentThread().join();
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
}

关键细节解释

  • 线程切换:必须用Schedulers.boundedElastic()来执行阻塞的readLine()——这个线程池专门为阻塞IO任务设计,会自动管理线程的创建和回收,不会影响Reactor的非阻塞性能。
  • 订阅状态检查:每次发送数据前检查sink.isCancelled(),如果订阅者已经取消订阅(比如调用了dispose()),就停止读取,避免浪费资源。
  • 资源清理:用try-with-resources包裹BufferedReader,确保输入流能被正确关闭。

方法二:手动实现Publisher(理解背压原理)

如果你想深入理解Reactor的发布-订阅模型和背压机制,可以手动实现Publisher接口,自己处理订阅关系、背压请求和取消逻辑。

完整代码示例

import org.reactivestreams.Publisher;
import org.reactivestreams.Subscriber;
import org.reactivestreams.Subscription;
import reactor.core.scheduler.Schedulers;
import java.io.BufferedReader;
import java.io.InputStreamReader;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicLong;

public class StdinPublisher implements Publisher<String> {
    @Override
    public void subscribe(Subscriber<? super String> subscriber) {
        StdinSubscription subscription = new StdinSubscription(subscriber);
        subscriber.onSubscribe(subscription);
        // 启动读取线程
        Schedulers.boundedElastic().schedule(subscription::startReading);
    }

    // 自定义Subscription,处理背压和取消逻辑
    private static class StdinSubscription implements Subscription {
        private final Subscriber<? super String> subscriber;
        private final BufferedReader reader;
        private final AtomicBoolean isCancelled = new AtomicBoolean(false);
        private final AtomicLong requestedCount = new AtomicLong(0);

        public StdinSubscription(Subscriber<? super String> subscriber) {
            this.subscriber = subscriber;
            this.reader = new BufferedReader(new InputStreamReader(System.in));
        }

        @Override
        public void request(long n) {
            if (n <= 0) {
                subscriber.onError(new IllegalArgumentException("请求数量必须大于0"));
                return;
            }
            // 累加订阅者请求的数量
            requestedCount.addAndGet(n);
            // 唤醒读取线程发送数据
            Thread.yield();
        }

        @Override
        public void cancel() {
            if (isCancelled.compareAndSet(false, true)) {
                try {
                    reader.close();
                } catch (Exception ignored) {}
            }
        }

        public void startReading() {
            try {
                String inputLine;
                while ((inputLine = reader.readLine()) != null && !isCancelled.get()) {
                    // 等待订阅者的请求,实现背压
                    while (requestedCount.get() == 0 && !isCancelled.get()) {
                        Thread.yield();
                    }
                    if (isCancelled.get()) {
                        break;
                    }
                    subscriber.onNext(inputLine);
                    // 减少剩余请求数量
                    requestedCount.decrementAndGet();
                }
                if (!isCancelled.get()) {
                    subscriber.onComplete();
                }
            } catch (Exception e) {
                if (!isCancelled.get()) {
                    subscriber.onError(e);
                }
            } finally {
                try {
                    reader.close();
                } catch (Exception ignored) {}
            }
        }
    }

    public static void main(String[] args) {
        Publisher<String> stdinPublisher = new StdinPublisher();
        reactor.core.publisher.Flux.from(stdinPublisher)
            .subscribe(
                line -> System.out.println("收到输入: " + line),
                error -> System.err.println("读取错误: " + error.getMessage()),
                () -> System.out.println("输入流关闭")
            );

        // 保持主线程存活
        try {
            Thread.currentThread().join();
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
}

关键细节解释

  • 背压处理:用AtomicLong跟踪订阅者的请求数量,只有当订阅者有请求时才发送数据,避免数据积压压垮订阅者。
  • 订阅管理:Subscription的cancel()方法负责关闭输入流并标记订阅状态,确保资源能被及时释放。
  • 线程安全:用原子变量保证多线程环境下的状态安全,避免竞态条件。

额外注意事项

  • 不要在Schedulers.parallel()或Schedulers.single()这些非阻塞线程池里执行readLine(),这些线程是为CPU密集型任务设计的,阻塞会导致线程池耗尽。
  • 如果需要对输入内容进行处理,可以在Flux后面链式调用操作符,比如:
    createStdinFlux()
        .map(String::trim) // 去除首尾空格
        .filter(line -> !line.isEmpty()) // 过滤空行
        .subscribe(line -> System.out.println("处理后: " + line));
    

内容的提问来源于stack exchange,提问作者Bullet-tooth

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 07:25:50