Administrator
发布于 2019-10-12 / 3312 阅读
45

Stream 并行流慎用:一次 parallelStream 导致的数据错乱

对账少了 3.7 万元:一行 parallelStream 惹的祸

10 月 8 号,财务来找我,说 9 月 30 号的日报表金额对不上,系统算出来 148.2 万,实际银行流水 151.9 万,少了 3.7 万。而且奇怪的是,同一个任务重跑一遍,出来的数字还不一样:第二次是 149.6 万。

同样的输入,两次跑出不同结果,那基本就是并发问题。我去翻那段对账代码,一眼就看到了:

public DailyReport buildReport(LocalDate date) {
    List<Settlement> settlements = settlementMapper.listByDate(date);   // 约 4.2 万条

    // 为了"快一点",前任同事改成了并行流
    List<ReportItem> items = new ArrayList<>();          // 普通 ArrayList
    settlements.parallelStream().forEach(s -> {
        if (s.getStatus() == SettlementStatus.SUCCESS) {
            items.add(convert(s));                        // 多线程并发 add!
        }
    });

    return new DailyReport(items);
}

ArrayList 不是线程安全的,多线程并发 add 会丢数据。这段代码跑了三个月,平时丢几条看不出来,月底数据量大才暴露。

复现一下

我写了个最小复现,跑 10 次:

public class ParallelStreamBug {
    public static void main(String[] args) {
        for (int round = 0; round < 10; round++) {
            List<Integer> source = IntStream.range(0, 100000).boxed().collect(Collectors.toList());

            List<Integer> result = new ArrayList<>();
            source.parallelStream().forEach(result::add);

            System.out.println("expected 100000, actual " + result.size()
                    + ", containsNull=" + result.contains(null));
        }
    }
}

输出:

expected 100000, actual 87241, containsNull=true
expected 100000, actual 91033, containsNull=true
expected 100000, actual 78620, containsNull=false
expected 100000, actual 93417, containsNull=false
expected 100000, actual 88155, containsNull=true
...

每次都少,而且少的数量不固定,有时还会混入 null

为什么会有 null?看 ArrayList.add 的实现:

public boolean add(E e) {
    ensureCapacityInternal(size + 1);   // 检查容量,不够就扩容
    elementData[size++] = e;            // 先赋值,再 size++
    return true;
}

elementData[size++] = e 这行不是原子的,它分三步:读 size、写 elementData[size]、size 加一。两个线程同时读到 size=5,都往下标 5 写,一个覆盖另一个;然后 size 变成 7,但下标 6 从来没被写过,就是 null。

扩容时的竞争更糟,grow()Arrays.copyOf 期间另一个线程的写入会直接丢掉。数据量大的时候丢得更多,跟我们的现象一致(月底数据量大,丢得多)。

parallelStream 用的是哪个线程池

这是我之前没搞清楚的点。看 AbstractPipeline 的源码:

final <R> R evaluate(TerminalOp<E_OUT, R> terminalOp) {
    // ...
    if (isParallel()) {
        return terminalOp.evaluateParallel(this, sourceSpliterator());
    }
}

最终会走到 ForkJoinTask,提交给 ForkJoinPool.commonPool()。这是个JVM 全局共享的静态池

// ForkJoinPool 里的静态字段
private static final ForkJoinPool common;

// 默认并行度
static final int COMMON_PARALLELISM = Runtime.getRuntime().availableProcessors() - 1;

重点是这两条:

  • 并行度默认 = CPU 核数 - 1。我们生产机器是 4 核,所以并行度只有 3。
  • 整个 JVM 里所有 parallelStream 共用这一个池。Tomcat 的请求线程、定时任务、批处理,谁用并行流都挤在这 3 个线程上。

所以有个很坑的现象:你的并行流可能跑得比串行还慢,因为公共池被别的任务占满了。我就遇到过——对账任务是凌晨 2 点跑的,同一时间还有三个定时任务也在用并行流,四个任务抢 3 个线程,互相拖累。

想改并行度可以用 JVM 参数:

-Djava.util.concurrent.ForkJoinPool.common.parallelism=8

但这是全局的,会影响所有并行流。JDK 没有提供"给单个流指定线程池"的 API(至少 JDK 8 没有)。真要隔离,只能绕开 parallelStream,手动用自定义线程池 + CompletableFuture,或者用这个 hack:

// 在自定义的 ForkJoinPool 里提交,让并行流在这个池里跑
ForkJoinPool customPool = new ForkJoinPool(8);
customPool.submit(() ->
    list.parallelStream().forEach(...)
).get();

这个技巧能work,因为 ForkJoinTask.fork() 会判断当前线程是不是 ForkJoinWorkerThread,是就提交到它所属的池。官方文档没承诺过这个行为,属于实现细节,我只在离线任务里用过,线上主流程不敢用。

正确的写法

回到对账那个方法。改成不依赖共享可变状态,用 collect

public DailyReport buildReport(LocalDate date) {
    List<Settlement> settlements = settlementMapper.listByDate(date);

    List<ReportItem> items = settlements.parallelStream()
            .filter(s -> s.getStatus() == SettlementStatus.SUCCESS)
            .map(this::convert)
            .collect(Collectors.toList());      // collect 内部会做分片合并,线程安全

    return new DailyReport(items);
}

collect 之所以安全,是因为 Stream 框架对每个分片创建独立的容器(supplier()),最后用 combiner() 合并。每个线程操作自己的容器,没有共享。这是并行流的正确用法:无状态、无副作用

聚合计算也一样,别用外部变量累加:

// 错误:共享可变状态
BigDecimal[] total = {BigDecimal.ZERO};
list.parallelStream().forEach(s -> total[0] = total[0].add(s.getAmount()));

// 正确:用 reduce
BigDecimal total = list.parallelStream()
        .map(Settlement::getAmount)
        .reduce(BigDecimal.ZERO, BigDecimal::add);

// 或者用 collectingAndThen / summarizingXXX
LongSummaryStatistics stat = list.parallelStream()
        .collect(Collectors.summarizingLong(Settlement::getCount));

还有 forEachforEachOrdered 的区别。forEach 在并行流里不保证顺序forEachOrdered 保证按原始顺序处理,但会牺牲并行带来的大部分收益:

// 不保证顺序
list.parallelStream().forEach(System.out::println);

// 保证顺序,但基本退化成串行
list.parallelStream().forEachOrdered(System.out::println);

我们的对账报表要求按时间排序输出,所以最后在外面加了一次 .sorted(),而不是用 forEachOrdered

什么时候该用并行流

我拿真实数据测了一下。4.2 万条结算记录,每条做一次金额计算和对象转换(纯 CPU,无 IO),机器 4 核:

数据量串行 stream并行 parallelStream加速比
1,000 条3 ms8 ms0.38x(更慢)
10,000 条21 ms14 ms1.5x
100,000 条187 ms72 ms2.6x
1,000,000 条1,742 ms596 ms2.9x

数据量小于 1 万条时,并行流是负优化——拆分成 ForkJoinTask、提交到线程池、合并结果,这套开销远大于计算本身。而且加速比最高也就 2.9 倍,不是理想中的 4 倍(4 核),因为合并阶段是串行的,还有公共池的调度开销。

再测一下"单个元素处理耗时"的影响。同样是 10 万条,但每条的处理时间不同:

单条处理耗时串行并行加速比
0.001 ms(简单取值)187 ms72 ms2.6x
0.1 ms(复杂计算)10,420 ms2,731 ms3.8x
1 ms(含一次 DB 查询)103,000 ms101,000 ms1.02x

第三行很说明问题:处理里有 IO(DB 查询)时,并行基本没收益。因为瓶颈在数据库连接池(我们配的 20 个),不在 CPU。而且并行流会让 3 个线程同时抢连接池,反而增加了锁竞争。

所以我的判断标准:

适合用并行流:

  • 数据量大于 1 万条
  • 单个元素的处理是纯 CPU 计算,耗时在微秒级以上
  • 处理逻辑无状态、无副作用,不修改共享变量
  • 不需要保证处理顺序(或者最后统一排序)
  • 运行在离线任务里,不和其他并行流抢公共池

不要用:

  • 处理里有 IO(查库、调接口、读写文件)
  • 数据量小于 1 万
  • 需要往外部容器 add、需要累加到外部变量
  • 在 Web 请求的主流程里(会占用公共池,影响其他并行任务)
  • 代码会被不熟并发的人维护(可读性差,容易改坏)

我们最后的对账任务改成了:用 collect 收集结果,数据量 4.2 万条走并行,耗时从串行的 91 毫秒降到 38 毫秒。其实这点收益意义不大,真正耗时的是前面那次数据库查询(800 毫秒)。所以严格说,这个并行流是完全没必要的优化——它带来的风险(丢数据)远超收益(53 毫秒)。

怎么避免以后再犯

我在项目里加了两条:

  1. 代码规范里写明:parallelStream 里禁止出现对外部变量的写操作。Code Review 时重点看这一点。
  2. 用 SonarQube 的规则。它有一条 squid:S3959,专门检测"在并行流里使用有状态的 lambda",能扫出大部分这类写法。

另外我养成了一个习惯:写并行流之前先问一句"串行跑一次要多久"。如果答案是几十毫秒,就没必要并行。我们项目里 90% 的 parallelStream 都属于这种没必要的优化。

先到这

《Stream 并行流慎用:一次 parallelStream 导致的数据错乱》这块我前前后后踩了不止一次。今天先写这些,后面想到新的再补。

参考