流组件非线程安全,多线程读写或异步传递易越界;需区分同步/异步上下文,避免共享状态污染;响应式流须明确线程调度与生命周期;集合边界须独立校验;异常须显式处理。
流组件本身不自动线程安全,哪怕只是个简单数据管道,只要多个线程同时读写、或在异步链中跨线程传递状态,边界没控住,就会出问题——不是“组件坏了”,而是访问逻辑越界了。
流操作必须区分同步与异步执行上下文
初学者常把 .map()、.filter() 当作纯函数调用,忽略其执行时机。在 Java Stream 中,串行流默认在当前线程执行;并行流(parallelStream())则交由 ForkJoinPool 调度,线程不可预测。若流中调用了含共享状态的操作(如静态计数器、ThreadLocal 缓存),就极易污染。
- 避免在流操作里直接修改外部变量(尤其是静态或单例字段)
- 若需累积状态,优先用
collect()配合线程安全容器(如ConcurrentHashMap或AtomicInteger) - 不要在
parallelStream()中调用非线程安全的工具类方法(如SimpleDateFormat)
订阅式流(如 RxJava、Reactor)必须明确生命周期与线程调度
像 Flux 或 Observable 这类响应式流,数据推送和订阅者处理可能不在同一线程。初学者常误以为 .subscribe() 写在哪,就在哪执行——其实默认使用 Schedulers.immediate() 或平台默认线程池,一旦涉及 I/O(如数据库、HTTP),必须显式指定线程模型。
- 阻塞操作(如 JDBC 查询)必须切到
boundedElastic()或自定义 IO 线程池,不能留在主线程或计算线程池 -
publishOn()控制下游执行线程,subscribeOn()控制上游源头线程,二者不可混淆 - 资源型流(如文件流、Socket 流)必须确保
onComplete或onError后释放,否则线程卡死或句柄泄漏
流中的集合边界必须独立校验,不能依赖外部保障
很多流组件底层封装了 List、Map 或数组,但“流安全”不等于“集合安全”。例如:从数据库查出空列表后调用 list.stream().findFirst().get(),会抛 NoSuchElementException;又或用索引流 IntStream.range(0, list.size()).mapToObj(list::get),若 list 在流执行中途被其他线程清空,就会触发 IndexOutOfBoundsException。
- 对可能为空的源,统一用
Optional包装再操作,如list.stream().findFirst().orElse(null) - 避免在流中直接通过下标访问集合,改用迭代语义(
forEach、reduce) - 若必须索引访问,先做防御性拷贝:
new ArrayList(list).stream()…
错误传播必须显式终止,不能靠流自动兜底
流组件通常不捕获运行时异常,一旦中间操作抛出 RuntimeException(如 NPE、ClassCastException),整个流会立即中断,且错误不会自动上报或记录。初学者常以为“流会吞掉异常”,结果线上静默失败。
- 关键业务流务必用
doOnError()记录日志,用onErrorResume()提供降级值 - 不要在
map中写 try-catch 包裹业务逻辑——这会让错误消失,破坏可观测性 - 对于 Reactor 的
Mono,onErrorMap()可将原始异常转为业务异常,便于统一拦截











