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
相关产品推荐
相关产品推荐

