前置知识: Java、Java、Java、Java

Java 函数式编程

21 min中级

Lambda、Stream、函数式接口与函数式编程范式的系统性深度剖析

前置知识

学习目标

  • 掌握「0. 本节阅读指引(先读这一节)」的核心机制、典型用法与常见陷阱
  • 掌握「1. 历史动机与演化」的核心机制、典型用法与常见陷阱
  • 掌握「2. 形式化定义」的核心机制、典型用法与常见陷阱
  • 掌握「3. 理论推导:函数式编程的内部机制」的核心机制、典型用法与常见陷阱
  • 掌握「4. 代码示例:从入门到进阶的完整实战」的核心机制、典型用法与常见陷阱

0. 本节阅读指引(先读这一节)

本篇是「Java 函数式编程」进阶文档,含较多理论章节。

第一遍只读:4. 代码示例、附录 A(java.util.function 核心接口速查表)、附录 B(Stream 操作分类速查);先会写会用。

可跳过:1-3 节(历史、形式化定义、理论推导)中的公式;5-8 节(对比、陷阱、工程实践、案例研究)第二遍细读。

前置:027 Lambda 与函数式编程、025 Stream API。

Java 函数式编程深度指南

函数式编程(Functional Programming, FP)作为一种起源于 λ 演算的编程范式,自 LISP(1958)诞生以来深刻影响了计算机科学的发展。Java 在 2014 年发布的 Java 8 中正式引入 Lambda 表达式、Stream API 与 java.util.function 包,标志着这门以面向对象为核心的语言完成了”对象 + 函数”的双范式融合。本文将从 λ 演算的形式化基础出发,深入剖析 Java 函数式接口的字节码本质、Stream 的惰性求值机制、并行流的 Fork/Join 调度原理,并通过完整的工程案例展示函数式思维如何重构传统命令式代码。读者将不仅学会”如何使用 Lambda”,更能理解”为何 invokedynamic 是 Java 函数式实现的基石”,从而在架构设计与性能优化层面做出有依据的决策。


1. 历史动机与演化

1.1 函数式编程的数学渊源

函数式编程的根基可追溯至 Alonzo Church 于 1936 年提出的 λ 演算(Lambda Calculus)。λ 演算是一个形式系统,用于研究函数定义、函数应用和递归。其核心语法仅有三条规则:

表达式 ::= 变量 | λ.变量.表达式 | 表达式 表达式

例如,恒等函数写作 λx.x,应用于参数写作 (λx.x) y,通过 β-归约(β-reduction)得到 y。Church 证明 λ 演算与图灵机等价,奠定了可计算性理论的基础。

λ 演算的三个核心特性直接映射到现代函数式编程:

  1. 函数是一等公民(First-class citizen):函数可以作为参数传递、作为返回值、赋值给变量。
  2. 纯函数(Pure function):相同输入始终产生相同输出,无副作用。
  3. 高阶函数(Higher-order function):接受函数作为参数或返回函数的函数。

1958 年,John McCarthy 基于 λ 演算设计了 LISP,开启了函数式编程的工程化实践。此后 ML(1973)、Haskell(1990)、Erlang(1986)等语言进一步发展了类型系统、惰性求值、模式匹配等特性。

1.2 Java 引入函数式编程的动机

Java 8 之前,函数式编程在 Java 中只能通过匿名内部类”模拟”,但语法冗长、性能开销大。考虑以下排序示例:

// Java 7 匿名内部类写法
List<String> names = Arrays.asList("Alice", "Bob", "Charlie");
Collections.sort(names, new Comparator<String>() {
    @Override
    public int compare(String s1, String s2) {
        return Integer.compare(s1.length(), s2.length());
    }
});

// Java 8 Lambda 写法
names.sort((s1, s2) -> Integer.compare(s1.length(), s2.length()));

// Java 8 方法引用写法
names.sort(Comparator.comparingInt(String::length));

从 6 行代码压缩到 1 行,不仅是语法糖,更是抽象层次的提升。Java 设计者引入函数式编程的核心动机包括:

  1. 并行化需求:多核 CPU 普及后,命令式的可变状态循环难以并行化,而 Stream 的无状态操作天然适合并行。
  2. 代码表达力:声明式风格让代码更接近业务意图,减少样板代码。
  3. API 设计灵活性:函数式接口允许 API 接受行为参数,如 forEach、map、filter。
  4. 生态竞争:Scala、Groovy 等 JVM 语言已支持函数式特性,Java 需保持竞争力。

1.3 Java 函数式编程的版本演进

版本年份关键特性JEP
Java 82014Lambda、Stream、java.util.function、接口默认方法、OptionalJEP 126
Java 92017Stream.takeWhile/dropWhile/ofNullable、Optional.ifPresentOrElse/or/streamJEP 269
Java 102018var 局部变量类型推断(间接简化 Lambda)JEP 286
Java 122019Collectors.teeing 收集器组合JEP 334
Java 162021Stream.toList() 简化收集JEP 395
Java 172021 sealed 类与模式匹配(增强 Stream 的类型安全)JEP 409
Java 212023虚拟线程与 Stream 的协同、模式匹配 switchJEP 444

1.4 设计哲学:为何 Java 不采用”纯”函数式

Java 选择”对象 + 函数”的混合范式,而非 Haskell 式的纯函数式,原因如下:

  1. 向后兼容:Java 拥有海量遗留代码,激进范式转变会破坏生态。
  2. 可变状态的现实需求:I/O、UI、数据库操作本质上是带副作用的,纯函数式需通过 Monad 等抽象处理,对工程师要求高。
  3. 性能考量:纯函数式语言的不可变数据结构在密集计算场景下可能产生大量临时对象,Java 选择让开发者自行权衡。
  4. 渐进式采纳:开发者可在需要时使用函数式风格,无需全盘接受。

这种”务实主义”使 Java 函数式编程既保留了 OOP 的封装与多态,又获得了 FP 的表达力与并行能力。


2. 形式化定义

2.1 函数式接口的形式化定义

定义 3.1(函数式接口):一个接口 II 是函数式接口,当且仅当 II 恰好声明了一个抽象方法 mm。形式化地:

I is functional  ⟺  ∣{m∈I∣m is abstract}∣=1I \text{ is functional} \iff |\{ m \in I \mid m \text{ is abstract} \}| = 1

其中,来自 Object 类的公开方法(如 equals、toString)不计入抽象方法数。

注:@FunctionalInterface 注解是可选的,仅用于编译期检查。即使不标注,满足上述条件的接口仍可作为函数式接口使用(如 Comparator)。

2.2 Lambda 表达式的类型推断

Lambda 表达式的类型由目标类型(Target Type)推断。设目标类型为 T=IαT = I_\alpha(函数式接口),其函数描述符(Functional Descriptor)为 τ1→τ2\tau_1 \to \tau_2。Lambda 表达式 (x)→e(x) \to e 的类型检查规则为:

Γ,x:τ1⊢e:τ2T=Iα (functional)descriptor(Iα)=τ1→τ2Γ⊢(x)→e:T\frac{\Gamma, x : \tau_1 \vdash e : \tau_2 \quad T = I_\alpha \text{ (functional)} \quad \text{descriptor}(I_\alpha) = \tau_1 \to \tau_2}{\Gamma \vdash (x) \to e : T}

例如,Comparator<String> 的描述符为 (String, String) -> int,因此 (s1, s2) -> s1.length() - s2.length() 的类型可推断为 Comparator<String>。

2.3 Stream 流水线的代数结构

Stream 流水线可形式化为一个单子(Monad)。定义 Stream<T> 为类型构造子,其满足以下单子定律:

  1. 左单位律(Left identity):Stream.of(x).flatMap(f) ≡ f.apply(x)
  2. 右单位律(Right identity):s.flatMap(Stream::of) ≡ s
  3. 结合律(Associativity):s.flatMap(f).flatMap(g) ≡ s.flatMap(x -> f.apply(x).flatMap(g))

flatMap 即单子的 bind 操作(写作 >>=),of 即 return/pure。这保证了 Stream 流水线的组合正确性。

2.4 纯函数的数学性质

定义 3.2(纯函数):函数 f:A→Bf : A \to B 是纯函数,当且仅当:

  1. 确定性:∀x∈A,∀i,j,fi(x)=fj(x)\forall x \in A, \forall i, j, f_i(x) = f_j(x)(多次调用结果相同)
  2. 无副作用:ff 的执行不修改任何外部状态

纯函数满足引用透明性(Referential Transparency):表达式 ee 与其求值结果 vv 可互换而不影响程序语义。形式化地:

∀C[⋅],C[e]≡C[v]where e⇓v\forall C[\cdot], \quad C[e] \equiv C[v] \quad \text{where } e \Downarrow v

这是函数式编程可推理性、可测试性、可并行化的数学基础。

2.5 柯里化的形式化

定义 3.3(柯里化):将 nn 元函数 f:A1×A2×⋯×An→Bf : A_1 \times A_2 \times \cdots \times A_n \to B 转换为一元函数链的过程:

curry(f):A1→(A2→(⋯→(An→B)⋯ ))\text{curry}(f) : A_1 \to (A_2 \to (\cdots \to (A_n \to B)\cdots))

Java 中的柯里化通过返回 Function 链实现:

// 二元函数
BiFunction<Integer, Integer, Integer> add = (a, b) -> a + b;

// 柯里化后
Function<Integer, Function<Integer, Integer>> curriedAdd = a -> b -> a + b;

// 部分应用
Function<Integer, Integer> add5 = curriedAdd.apply(5);
System.out.println(add5.apply(3)); // 8

柯里化使部分应用(Partial Application)成为可能,是函数组合的核心技术。


3. 理论推导:函数式编程的内部机制

3.1 invokedynamic 与 Lambda 的字节码实现

Java 8 的 Lambda 表达式并不编译为匿名内部类,而是通过 invokedynamic(JSR 292)实现。这一设计决策由 Brian Goetz 主导,核心动机是:

  1. 避免运行时类生成开销:匿名内部类会在编译期生成 .class 文件,增加类加载开销。
  2. 保持实现灵活性:将 Lambda 的实际实现策略延迟到运行时,由 LambdaMetafactory 决定。
  3. 性能优化空间:未来可采用方法句柄(MethodHandle)的内联优化。

考虑以下 Lambda:

Comparator<String> cmp = (s1, s2) -> Integer.compare(s1.length(), s2.length());

编译后的字节码大致为:

INVOKEDYNAMIC compare()Ljava/util/Comparator; [
  // bootstrap method
  LambdaMetafactory.metafactory(...),
  // method handle to synthetic lambda method
  lambda$0(Ljava/lang/String;Ljava/lang/String;)I
]

其中 lambda$0 是编译器生成的合成方法(Synthetic Method),包含 Lambda 体的实际逻辑。LambdaMetafactory.metafactory 在首次调用时生成一个实现 Comparator 的动态类,将 lambda$0 包装为具体实现。

首次调用 vs 后续调用:

  • 首次调用:metafactory 解析调用点(CallSite),生成 LambdaForm,绑定方法句柄。
  • 后续调用:直接调用已绑定的方法句柄,性能接近直接方法调用。

通过 -Djdk.internal.lambda.dumpProxyClasses=/tmp/lambda 可导出生成的动态类,验证其结构。

3.2 Stream 流水线的惰性求值机制

Stream 的核心设计是惰性求值(Lazy Evaluation):中间操作(Intermediate Operation)不会立即执行,只有在终端操作(Terminal Operation)触发时才回溯执行整个流水线。

惰性求值的实现原理:

  1. 每个中间操作返回一个新的 Stream 实例(如 ReferencePipeline),保存上游 Stage 引用与操作逻辑。
  2. 终端操作触发 evaluate 方法,从 Sink 链的尾端向头端反向回溯。
  3. 每个元素从源头流向 Sink 链,依次被 filter、map 等处理。

考虑以下代码:

List<String> result = Stream.of("a", "bb", "ccc", "dddd")
    .filter(s -> s.length() > 1)
    .map(String::toUpperCase)
    .collect(Collectors.toList());

执行时 Sink 链的结构:

Source -> filterSink -> mapSink -> collectSink
            ↑           ↑           ↑
        接收元素       接收过滤后的   收集到结果
        判断长度>1     转大写

惰性求值的关键优势:

  1. 短路优化:findFirst、anyMatch 等操作可提前终止。
  2. 融合优化:多个操作可融合为单次遍历,避免中间集合。
  3. 无限流支持:Stream.iterate、Stream.generate 可表示无限序列。

短路示例:

// 只处理到第一个匹配元素即停止
Optional<String> first = Stream.of("a", "bb", "ccc", "dddd")
    .peek(s -> System.out.println("处理: " + s))  // 调试用
    .filter(s -> s.length() > 2)
    .findFirst();

// 输出:
// 处理: a
// 处理: bb
// 处理: ccc
// 结果: Optional[CCC]

可以看到 dddd 未被处理,体现了短路优化。

3.3 并行流的 Fork/Join 调度

parallelStream 与 stream().parallel() 基于 ForkJoinPool.commonPool() 实现并行。其工作流程:

  1. 分割(Fork):将源数据递归分割为子任务,直到达到阈值(默认每个任务约 1024 个元素)。
  2. 执行(Compute):每个子任务在工作线程中独立执行流水线。
  3. 合并(Join):将子任务结果合并为最终结果。

Spliterator 的分割语义:

// ArrayList 的 Spliterator 支持 ORDERED | SIZED | SUBSIZED
// 可精确分割
Spliterator<String> spliterator = list.spliterator();
Spliterator<String> prefix = spliterator.trySplit();  // 分割前半部分

并行流的适用条件:

  1. 数据源支持高效分割(ArrayList 优于 LinkedList)。
  2. 操作无状态(不依赖外部可变变量)。
  3. 操作无副作用(不修改共享状态)。
  4. 数据量足够大(通常 > 10000 元素)。
  5. 操作本身计算量足够大(简单 map 难以抵消并行开销)。

3.4 函数式接口的字节码与运行时

@FunctionalInterface 注解在运行时通过 @Retention(RetentionPolicy.RUNTIME) 保留,但 JVM 不依赖它判断函数式接口。实际判断逻辑在 LambdaMetafactory 中:

// 简化的 LambdaMetafactory.metafactory 逻辑
public static CallSite metafactory(...) {
    // 1. 校验目标接口是否为函数式接口
    if (!isFunctionalInterface(targetType)) {
        throw new LambdaConversionException("Not a functional interface");
    }
    // 2. 提取函数描述符
    MethodType descriptor = targetType.getFunctionalDescriptor();
    // 3. 校验方法句柄类型与描述符兼容
    if (!isAdaptCompatible(actualMethodType, descriptor)) {
        throw new LambdaConversionException("Type mismatch");
    }
    // 4. 生成动态代理类
    Class<?> implClass = generateLambdaClass(...);
    // 5. 绑定调用点
    return new ConstantCallSite(MethodHandles.constant(implClass, implClass.newInstance()));
}

这一机制使得 Lambda 的实现策略可在未来演进(如 GraalVM 的部分求值优化),而无需修改字节码。

3.5 方法引用的四种形式

方法引用(Method Reference)是 Lambda 的语法糖,编译器将其转换为方法句柄。四种形式:

形式语法等价 Lambda示例
静态方法引用ClassName::staticMethodx -> ClassName.staticMethod(x)Integer::parseInt
特定对象的实例方法引用instance::methodx -> instance.method(x)System.out::println
任意对象的实例方法引用ClassName::method(obj, x) -> obj.method(x)String::length
构造方法引用ClassName::newx -> new ClassName(x)ArrayList::new

任意对象实例方法引用的语义:第一个参数成为接收者(receiver),其余参数作为方法参数。例如 String::concat 等价于 (s1, s2) -> s1.concat(s2)。

3.6 Collectors 的归约代数

Collector 接口定义了五个组件:

public interface Collector<T, A, R> {
    Supplier<A> supplier();           // 创建累加器初始值
    BiConsumer<A, T> accumulator();   // 累加元素到累加器
    BinaryOperator<A> combiner();     // 合并两个累加器(并行用)
    Function<A, R> finisher();        // 最终转换
    Set<Characteristics> characteristics();  // 特征标识
}

其归约过程形式化为:

collect(s,c)=finisherc(reduce(s,supplierc,accumulatorc,combinerc))\text{collect}(s, c) = \text{finisher}_c\left(\text{reduce}\left(s, \text{supplier}_c, \text{accumulator}_c, \text{combiner}_c\right)\right)

其中 reduce 在并行场景下分治执行:

  1. 将流分为 nn 个子流。
  2. 每个子流独立执行 supplier → accumulator。
  3. 通过 combiner 两两合并子结果。
  4. finisher 转换为最终结果。

Characteristics.IDENTITY_FINISH 表示 finisher 是恒等函数,可省略。UNORDERED 表示结果与顺序无关,允许并行合并无序化。CONCURRENT 表示累加器线程安全,可并行累加同一容器(罕见)。


4. 代码示例:从入门到进阶的完整实战

4.1 入门:Lambda 基础语法

package com.example.fp.basics;

import java.util.*;
import java.util.function.*;

public class LambdaBasics {

    public static void main(String[] args) {
        // 无参数 Lambda(Runnable 的描述符为 () -> void)
        Runnable r = () -> System.out.println("Hello, Lambda");
        r.run();

        // 单参数 Lambda(Consumer<String> 的描述符为 (String) -> void)
        Consumer<String> print = s -> System.out.println(s);
        print.accept("Hello, Consumer");

        // 类型推断省略
        Consumer<String> printInferred = s -> System.out.println(s);

        // 多参数 Lambda(BiFunction 的描述符为 (Integer, Integer) -> Integer)
        BiFunction<Integer, Integer, Integer> add = (a, b) -> a + b;
        System.out.println(add.apply(3, 5));  // 8

        // 代码块 Lambda(需要 return 语句)
        Comparator<String> cmp = (s1, s2) -> {
            int diff = s1.length() - s2.length();
            return diff != 0 ? diff : s1.compareTo(s2);
        };

        // 方法引用:四种形式
        Consumer<String> println = System.out::println;  // 特定对象实例方法引用
        Function<String, Integer> len = String::length;  // 任意对象实例方法引用
        Function<String, Integer> parse = Integer::parseInt;  // 静态方法引用
        Supplier<ArrayList<String>> factory = ArrayList::new;  // 构造方法引用

        // 变量捕获(effectively final 才能捕获)
        String prefix = "User: ";
        Consumer<String> greet = name -> System.out.println(prefix + name);
        greet.accept("Alice");
        // prefix = "Admin: ";  // 编译错误:被捕获变量必须 effectively final
    }
}

4.2 进阶:自定义函数式接口与组合

package com.example.fp.advanced;

import java.util.function.*;

/**
 * 自定义函数式接口,支持链式组合
 */
@FunctionalInterface
public interface Transformer<T, R> {
    R transform(T input);

    /**
     * 与后续 Function 组合,等价于 Function.andThen
     */
    default <V> Transformer<T, V> andThen(Function<R, V> after) {
        return input -> after.apply(this.transform(input));
    }

    /**
     * 与前置 Function 组合,等价于 Function.compose
     */
    default <V> Transformer<V, R> compose(Function<V, T> before) {
        return input -> this.transform(before.apply(input));
    }
}

class TransformerDemo {
    public static void main(String[] args) {
        Transformer<String, Integer> strLen = String::length;
        Transformer<String, String> upper = String::toUpperCase;

        // 组合:先转大写,再求长度,再格式化
        Transformer<String, String> pipeline = upper
            .andThen(len -> "长度=" + len)
            .andThen(s -> "[" + s + "]");

        System.out.println(pipeline.transform("hello"));  // [长度=5]
    }
}

4.3 实战:Stream 数据处理管道

package com.example.fp.stream;

import java.util.*;
import java.util.stream.*;
import java.math.BigDecimal;
import java.time.LocalDateTime;

public class StreamPipeline {

    record Order(Long id, String customer, BigDecimal amount,
                 String status, LocalDateTime createdAt) {}

    public static void main(String[] args) {
        List<Order> orders = List.of(
            new Order(1L, "Alice", new BigDecimal("100.50"), "COMPLETED",
                LocalDateTime.of(2024, 1, 15, 10, 0)),
            new Order(2L, "Bob", new BigDecimal("250.00"), "COMPLETED",
                LocalDateTime.of(2024, 1, 20, 14, 30)),
            new Order(3L, "Alice", new BigDecimal("75.25"), "PENDING",
                LocalDateTime.of(2024, 2, 1, 9, 15)),
            new Order(4L, "Charlie", new BigDecimal("500.00"), "COMPLETED",
                LocalDateTime.of(2024, 2, 5, 16, 45)),
            new Order(5L, "Bob", new BigDecimal("180.75"), "CANCELLED",
                LocalDateTime.of(2024, 2, 10, 11, 20))
        );

        // 1. 基本过滤与映射
        List<String> completedCustomers = orders.stream()
            .filter(o -> "COMPLETED".equals(o.status()))
            .map(Order::customer)
            .distinct()
            .sorted()
            .collect(Collectors.toList());
        System.out.println("已完成订单客户: " + completedCustomers);

        // 2. 分组统计
        Map<String, Long> countByCustomer = orders.stream()
            .collect(Collectors.groupingBy(Order::customer, Collectors.counting()));
        System.out.println("客户订单数: " + countByCustomer);

        // 3. 分组求和
        Map<String, BigDecimal> sumByCustomer = orders.stream()
            .filter(o -> "COMPLETED".equals(o.status()))
            .collect(Collectors.groupingBy(
                Order::customer,
                Collectors.reducing(
                    BigDecimal.ZERO,
                    Order::amount,
                    BigDecimal::add
                )
            ));
        System.out.println("客户已完成总额: " + sumByCustomer);

        // 4. 分组后取最大值
        Map<String, Optional<Order>> maxByCustomer = orders.stream()
            .collect(Collectors.groupingBy(
                Order::customer,
                Collectors.maxBy(Comparator.comparing(Order::amount))
            ));
        maxByCustomer.forEach((k, v) ->
            System.out.println(k + " 最大订单: " + v.map(Order::amount).orElse(null)));

        // 5. 多级分组:客户 + 状态
        Map<String, Map<String, List<Order>>> multiGrouped = orders.stream()
            .collect(Collectors.groupingBy(
                Order::customer,
                Collectors.groupingBy(Order::status)
            ));
        System.out.println("多级分组: " + multiGrouped);

        // 6. 分区
        Map<Boolean, List<Order>> partitioned = orders.stream()
            .collect(Collectors.partitioningBy(o ->
                o.amount().compareTo(new BigDecimal("200")) > 0));
        System.out.println("大额订单数: " + partitioned.get(true).size());

        // 7. 统计摘要
        DoubleSummaryStatistics stats = orders.stream()
            .filter(o -> "COMPLETED".equals(o.status()))
            .mapToDouble(o -> o.amount().doubleValue())
            .summaryStatistics();
        System.out.printf("已完成订单统计: count=%d, sum=%.2f, avg=%.2f, max=%.2f%n",
            stats.getCount(), stats.getSum(), stats.getAverage(), stats.getMax());

        // 8. 字符串拼接
        String customerList = orders.stream()
            .map(Order::customer)
            .distinct()
            .collect(Collectors.joining(", ", "[", "]"));
        System.out.println("客户列表: " + customerList);

        // 9. 不可变收集
        List<String> immutable = orders.stream()
            .map(Order::customer)
            .distinct()
            .collect(Collectors.toUnmodifiableList());

        // 10. Java 16+ 简化收集
        List<Order> completedOrders = orders.stream()
            .filter(o -> "COMPLETED".equals(o.status()))
            .toList();  // 返回不可变 List
    }
}

4.4 实战:自定义 Collector

package com.example.fp.collector;

import java.util.*;
import java.util.function.*;
import java.util.stream.Collector;
import java.util.stream.Collectors;

/**
 * 自定义收集器:将流元素收集为以指定分隔符连接、带前后缀的字符串
 */
public class JoiningCollector
        implements Collector<String, StringJoiner, String> {

    private final String delimiter;
    private final String prefix;
    private final String suffix;

    public JoiningCollector(String delimiter, String prefix, String suffix) {
        this.delimiter = delimiter;
        this.prefix = prefix;
        this.suffix = suffix;
    }

    @Override
    public Supplier<StringJoiner> supplier() {
        return () -> new StringJoiner(delimiter, prefix, suffix);
    }

    @Override
    public BiConsumer<StringJoiner, String> accumulator() {
        return StringJoiner::add;
    }

    @Override
    public BinaryOperator<StringJoiner> combiner() {
        return StringJoiner::merge;
    }

    @Override
    public Function<StringJoiner, String> finisher() {
        return StringJoiner::toString;
    }

    @Override
    public Set<Characteristics> characteristics() {
        // 不标记 IDENTITY_FINISH,因为 StringJoiner -> String 是有损转换
        return Set.of(Characteristics.CONCURRENT);
    }

    public static void main(String[] args) {
        String result = List.of("Alice", "Bob", "Charlie").stream()
            .collect(new JoiningCollector(", ", "{", "}"));
        System.out.println(result);  // {Alice, Bob, Charlie}
    }
}

/**
 * 自定义收集器:分位数计算
 */
class PercentileCollector implements Collector<Double, List<Double>, Map<String, Double>> {

    @Override
    public Supplier<List<Double>> supplier() {
        return ArrayList::new;
    }

    @Override
    public BiConsumer<List<Double>, Double> accumulator() {
        return List::add;
    }

    @Override
    public BinaryOperator<List<Double>> combiner() {
        return (l1, l2) -> {
            l1.addAll(l2);
            return l1;
        };
    }

    @Override
    public Function<List<Double>, Map<String, Double>> finisher() {
        return list -> {
            list.sort(Double::compare);
            int n = list.size();
            Map<String, Double> result = new LinkedHashMap<>();
            result.put("p50", percentile(list, 0.5));
            result.put("p90", percentile(list, 0.9));
            result.put("p99", percentile(list, 0.99));
            return result;
        };
    }

    private double percentile(List<Double> sorted, double p) {
        int index = (int) Math.ceil(p * sorted.size()) - 1;
        return sorted.get(Math.max(0, index));
    }

    @Override
    public Set<Characteristics> characteristics() {
        return Set.of();  // 有 finisher,不标记 IDENTITY_FINISH
    }

    public static void main(String[] args) {
        List<Double> data = java.util.stream.DoubleStream
            .generate(() -> Math.random() * 100)
            .limit(1000)
            .boxed()
            .collect(Collectors.toList());

        Map<String, Double> percentiles = data.stream()
            .collect(new PercentileCollector());
        System.out.println("分位数: " + percentiles);
    }
}

4.5 实战:并行流与自定义 ForkJoinPool

package com.example.fp.parallel;

import java.util.*;
import java.util.concurrent.*;
import java.util.stream.*;

public class ParallelStreamDemo {

    public static void main(String[] args) throws Exception {
        // 大数据集
        List<Integer> bigData = IntStream.rangeClosed(1, 10_000_000)
            .boxed()
            .collect(Collectors.toList());

        // 1. 串行流
        long start = System.currentTimeMillis();
        long sumSequential = bigData.stream()
            .mapToLong(x -> (long) x * x)
            .sum();
        System.out.println("串行耗时: " + (System.currentTimeMillis() - start) + "ms");

        // 2. 并行流(使用默认 commonPool)
        start = System.currentTimeMillis();
        long sumParallel = bigData.parallelStream()
            .mapToLong(x -> (long) x * x)
            .sum();
        System.out.println("并行耗时: " + (System.currentTimeMillis() - start) + "ms");

        // 3. 自定义并行度的 ForkJoinPool
        ForkJoinPool customPool = new ForkJoinPool(4);
        try {
            start = System.currentTimeMillis();
            long sumCustom = customPool.submit(() ->
                bigData.parallelStream()
                    .mapToLong(x -> (long) x * x)
                    .sum()
            ).get();
            System.out.println("自定义池并行耗时: " + (System.currentTimeMillis() - start) + "ms");
        } finally {
            customPool.shutdown();
        }

        // 4. 并行流的危险:共享可变状态
        List<Integer> unsafeList = new ArrayList<>();
        IntStream.rangeClosed(1, 1000).parallel()
            .forEach(unsafeList::add);  // 危险!ArrayList 非线程安全
        System.out.println("不安全列表大小(可能丢失): " + unsafeList.size());

        // 5. 安全的并行收集
        List<Integer> safeList = IntStream.rangeClosed(1, 1000).parallel()
            .boxed()
            .collect(Collectors.toList());
        System.out.println("安全列表大小: " + safeList.size());

        // 6. 并行排序(注意顺序)
        List<Integer> sorted = bigData.parallelStream()
            .sorted(Comparator.reverseOrder())
            .collect(Collectors.toList());
        System.out.println("降序前10: " + sorted.subList(0, 10));
    }
}

4.6 实战:函数组合与柯里化

package com.example.fp.composition;

import java.util.*;
import java.util.function.*;

public class FunctionComposition {

    public static void main(String[] args) {
        // 基本函数组合
        Function<Integer, Integer> doubleIt = x -> x * 2;
        Function<Integer, Integer> addOne = x -> x + 1;
        Function<Integer, Integer> square = x -> x * x;

        // compose: g.compose(f).apply(x) = g(f(x))
        Function<Integer, Integer> doubleThenAdd = addOne.compose(doubleIt);
        System.out.println(doubleThenAdd.apply(5));  // (5*2)+1 = 11

        // andThen: f.andThen(g).apply(x) = g(f(x))
        Function<Integer, Integer> addThenDouble = doubleIt.andThen(addOne);
        System.out.println(addThenDouble.apply(5));  // (5+1)*2 = 12

        // 多函数组合
        Function<Integer, Integer> pipeline = doubleIt
            .andThen(addOne)
            .andThen(square);
        System.out.println(pipeline.apply(5));  // ((5*2)+1)^2 = 121

        // 柯里化:将 BiFunction 转为 Function 链
        BiFunction<Integer, Integer, Integer> multiply = (a, b) -> a * b;
        Function<Integer, Function<Integer, Integer>> curriedMultiply =
            a -> b -> a * b;

        // 部分应用
        Function<Integer, Integer> triple = curriedMultiply.apply(3);
        Function<Integer, Integer> quadruple = curriedMultiply.apply(4);
        System.out.println(triple.apply(5));     // 15
        System.out.println(quadruple.apply(5));  // 20

        // 高阶函数:函数作为参数
        Function<Integer, Integer> composed = composeAll(
            Arrays.asList(doubleIt, addOne, square)
        );
        System.out.println(composed.apply(5));  // ((5*2)+1)^2 = 121

        // Predicate 组合
        Predicate<Integer> isPositive = x -> x > 0;
        Predicate<Integer> isEven = x -> x % 2 == 0;
        Predicate<Integer> isPositiveEven = isPositive.and(isEven);
        Predicate<Integer> isPositiveOrEven = isPositive.or(isEven);
        Predicate<Integer> isNegative = isPositive.negate();

        System.out.println(isPositiveEven.test(4));   // true
        System.out.println(isPositiveEven.test(3));   // false
        System.out.println(isNegative.test(-5));      // true

        // Consumer 组合
        Consumer<String> printUpper = s -> System.out.println(s.toUpperCase());
        Consumer<String> printLength = s -> System.out.println(s.length());
        Consumer<String> combined = printUpper.andThen(printLength);
        combined.accept("Hello");  // HELLO \n 5
    }

    /**
     * 将多个 Function 组合为一个
     */
    @SafeVarargs
    public static <T> Function<T, T> composeAll(Function<T, T>... functions) {
        return Arrays.stream(functions)
            .reduce(Function.identity(), Function::andThen);
    }

    /**
     * 重载:从 List 组合
     */
    public static <T> Function<T, T> composeAll(List<Function<T, T>> functions) {
        return functions.stream()
            .reduce(Function.identity(), Function::andThen);
    }
}

4.7 实战:Optional 的函数式用法

package com.example.fp.optional;

import java.util.*;
import java.util.function.*;
import java.util.stream.*;

public class OptionalDemo {

    public static void main(String[] args) {
        // 创建 Optional
        Optional<String> present = Optional.of("Hello");
        Optional<String> empty = Optional.empty();
        Optional<String> nullable = Optional.ofNullable(null);

        // 函数式消费
        present.map(String::toUpperCase)
            .ifPresent(System.out::println);  // HELLO

        // 链式操作
        String result = present
            .filter(s -> s.length() > 3)
            .map(String::toUpperCase)
            .orElse("DEFAULT");
        System.out.println(result);  // HELLO

        // flatMap 处理嵌套 Optional
        Optional<Optional<String>> nested = present.map(s -> Optional.of(s + "!"));
        Optional<String> flattened = present.flatMap(s -> Optional.of(s + "!"));
        System.out.println(flattened.orElse(""));  // Hello!

        // Java 9+ ifPresentOrElse
        present.ifPresentOrElse(
            s -> System.out.println("存在: " + s),
            () -> System.out.println("不存在")
        );

        // Java 9+ or
        Optional<String> fallback = empty.or(() -> Optional.of("默认值"));
        System.out.println(fallback.get());  // 默认值

        // Java 9+ stream
        long count = Stream.of(
                Optional.of("a"),
                Optional.empty(),
                Optional.of("b")
            )
            .flatMap(Optional::stream)
            .count();
        System.out.println("非空数量: " + count);  // 2

        // 实战:避免 null 检查地狱
        User user = findUser(1L).orElseThrow(() -> new RuntimeException("用户不存在"));
        System.out.println(user.name());
    }

    record User(Long id, String name) {}

    static Optional<User> findUser(Long id) {
        if (id == 1L) {
            return Optional.of(new User(1L, "Alice"));
        }
        return Optional.empty();
    }
}

4.8 实战:原始类型流

package com.example.fp.primitive;

import java.util.*;
import java.util.stream.*;

public class PrimitiveStreamDemo {

    public static void main(String[] args) {
        // IntStream 避免装箱开销
        int[] numbers = IntStream.rangeClosed(1, 100).toArray();

        // 统计摘要
        IntSummaryStatistics stats = IntStream.of(numbers)
            .filter(n -> n % 2 == 0)
            .summaryStatistics();
        System.out.printf("偶数统计: count=%d, sum=%d, avg=%.2f, min=%d, max=%d%n",
            stats.getCount(), stats.getSum(), stats.getAverage(),
            stats.getMin(), stats.getMax());

        // 数值计算
        long factorial = LongStream.rangeClosed(1, 20)
            .reduce(1L, (a, b) -> a * b);
        System.out.println("20! = " + factorial);

        // 斐波那契数列(iterate)
        List<Long> fibonacci = Stream.iterate(
                new long[]{0, 1},
                pair -> new long[]{pair[1], pair[0] + pair[1]}
            )
            .limit(20)
            .map(pair -> pair[0])
            .collect(Collectors.toList());
        System.out.println("斐波那契前20项: " + fibonacci);

        // 勾股数
        record PythagoreanTriple(int a, int b, int c) {}
        List<PythagoreanTriple> triples = IntStream.rangeClosed(1, 100)
            .boxed()
            .flatMap(a -> IntStream.rangeClosed(a, 100)
                .filter(b -> Math.sqrt(a * a + b * b) % 1 == 0)
                .mapToObj(b -> new PythagoreanTriple(a, b,
                    (int) Math.sqrt(a * a + b * b))))
            .collect(Collectors.toList());
        System.out.println("勾股数(前5组): " + triples.subList(0, 5));

        // DoubleStream 浮点计算
        double pi = DoubleStream.iterate(1.0, n -> n + 2)
            .limit(1000)
            .map(n -> 4 * Math.pow(-1, (n - 1) / 2) / n)
            .sum();
        System.out.printf("π ≈ %.6f%n", pi);
    }
}

4.9 实战:Stream 高级操作

package com.example.fp.advanced;

import java.util.*;
import java.util.function.*;
import java.util.stream.*;

public class AdvancedStreamOps {

    public static void main(String[] args) {
        // 1. flatMap 展平嵌套结构
        List<List<Integer>> nested = List.of(
            List.of(1, 2, 3),
            List.of(4, 5),
            List.of(6, 7, 8, 9)
        );
        List<Integer> flattened = nested.stream()
            .flatMap(List::stream)
            .collect(Collectors.toList());
        System.out.println("展平: " + flattened);

        // 2. flatMap 处理字符串
        List<String> words = List.of("Hello World", "Java Stream", "Flat Map");
        List<String> chars = words.stream()
            .flatMap(s -> Arrays.stream(s.split(" ")))
            .collect(Collectors.toList());
        System.out.println("分词: " + chars);

        // 3. reduce 归约
        Optional<Integer> sum = Stream.of(1, 2, 3, 4, 5)
            .reduce(Integer::sum);
        System.out.println("求和: " + sum.orElse(0));

        // 4. 带初始值的 reduce
        int product = Stream.of(1, 2, 3, 4, 5)
            .reduce(1, (a, b) -> a * b);
        System.out.println("阶乘: " + product);

        // 5. 复杂对象归约
        record Acc(int sum, int count) {}
        Acc result = Stream.of(1, 2, 3, 4, 5)
            .reduce(
                new Acc(0, 0),
                (acc, n) -> new Acc(acc.sum() + n, acc.count() + 1),
                (a1, a2) -> new Acc(a1.sum() + a2.sum(), a1.count() + a2.count())
            );
        System.out.println("平均: " + (double) result.sum() / result.count());

        // 6. takeWhile / dropWhile (Java 9+)
        List<Integer> taken = Stream.of(1, 2, 3, 4, 5, 2, 1)
            .takeWhile(n -> n < 4)
            .collect(Collectors.toList());
        System.out.println("takeWhile: " + taken);  // [1, 2, 3]

        List<Integer> dropped = Stream.of(1, 2, 3, 4, 5, 2, 1)
            .dropWhile(n -> n < 4)
            .collect(Collectors.toList());
        System.out.println("dropWhile: " + dropped);  // [4, 5, 2, 1]

        // 7. ofNullable (Java 9+)
        long count = Stream.of(1, 2, null, 3, null, 4)
            .flatMap(Stream::ofNullable)
            .count();
        System.out.println("非空数量: " + count);  // 4

        // 8. Collectors.teeing (Java 12+)
        record Stats(double sum, int count) {}
        Stats stats = Stream.of(1.0, 2.0, 3.0, 4.0, 5.0)
            .collect(Collectors.teeing(
                Collectors.summingDouble(Double::doubleValue),
                Collectors.counting(),
                (s, c) -> new Stats(s, c.intValue())
            ));
        System.out.println("统计: " + stats);

        // 9. groupingBy 下游收集器组合
        record Sale(String category, double amount) {}
        Map<String, Double> avgByCategory = Stream.of(
            new Sale("Books", 50),
            new Sale("Books", 30),
            new Sale("Electronics", 200),
            new Sale("Electronics", 300)
        ).collect(Collectors.groupingBy(
            Sale::category,
            Collectors.averagingDouble(Sale::amount)
        ));
        System.out.println("分类平均: " + avgByCategory);

        // 10. collectingAndThen 后处理
        Map<String, List<Sale>> unmodifiable = Stream.of(
            new Sale("Books", 50),
            new Sale("Electronics", 200)
        ).collect(Collectors.collectingAndThen(
            Collectors.groupingBy(Sale::category),
            Collections::unmodifiableMap
        ));
    }
}

4.10 完整案例:基于函数式的领域建模

package com.example.fp.domain;

import java.util.*;
import java.util.function.*;
import java.util.stream.*;
import java.math.BigDecimal;
import java.time.LocalDateTime;

/**
 * 基于函数式思维的电商领域模型
 */
public class FunctionalECommerce {

    // 不可变值对象
    record Product(Long id, String name, BigDecimal price, String category) {}
    record Customer(Long id, String name, String tier) {}
    record OrderItem(Product product, int quantity) {
        BigDecimal subtotal() {
            return product.price().multiply(BigDecimal.valueOf(quantity));
        }
    }
    record Order(Long id, Customer customer, List<OrderItem> items,
                 LocalDateTime createdAt, String status) {
        BigDecimal total() {
            return items.stream()
                .map(OrderItem::subtotal)
                .reduce(BigDecimal.ZERO, BigDecimal::add);
        }
    }

    // 函数式服务
    static class OrderService {
        private final List<Order> orders;

        OrderService(List<Order> orders) {
            this.orders = orders;
        }

        // 高阶查询函数
        public List<Order> queryOrders(Predicate<Order> predicate,
                                       Comparator<Order> sorter,
                                       int limit) {
            return orders.stream()
                .filter(predicate)
                .sorted(sorter)
                .limit(limit)
                .collect(Collectors.toList());
        }

        // 函数组合的查询构建器
        public static class OrderQueryBuilder {
            private Predicate<Order> predicate = o -> true;
            private Comparator<Order> sorter = Comparator.comparing(Order::createdAt);
            private int limit = Integer.MAX_VALUE;

            public OrderQueryBuilder where(Predicate<Order> p) {
                this.predicate = this.predicate.and(p);
                return this;
            }

            public OrderQueryBuilder orderBy(Comparator<Order> s) {
                this.sorter = s;
                return this;
            }

            public OrderQueryBuilder limit(int n) {
                this.limit = n;
                return this;
            }

            public Predicate<Order> buildPredicate() {
                return predicate;
            }

            public Comparator<Order> buildSorter() {
                return sorter;
            }

            public int buildLimit() {
                return limit;
            }
        }

        // 统计报告
        public Map<String, BigDecimal> revenueByCategory() {
            return orders.stream()
                .filter(o -> "COMPLETED".equals(o.status()))
                .flatMap(o -> o.items().stream())
                .collect(Collectors.groupingBy(
                    item -> item.product().category(),
                    Collectors.reducing(
                        BigDecimal.ZERO,
                        OrderItem::subtotal,
                        BigDecimal::add
                    )
                ));
        }

        // 客户消费排名
        public List<Map.Entry<String, BigDecimal>> topCustomers(int n) {
            return orders.stream()
                .filter(o -> "COMPLETED".equals(o.status()))
                .collect(Collectors.groupingBy(
                    o -> o.customer().name(),
                    Collectors.reducing(
                        BigDecimal.ZERO,
                        Order::total,
                        BigDecimal::add
                    )
                ))
                .entrySet().stream()
                .sorted(Map.Entry.<String, BigDecimal>comparingByValue().reversed())
                .limit(n)
                .collect(Collectors.toList());
        }
    }

    public static void main(String[] args) {
        Product book = new Product(1L, "Java编程思想", new BigDecimal("99.00"), "Books");
        Product laptop = new Product(2L, "MacBook Pro", new BigDecimal("12999.00"), "Electronics");
        Product phone = new Product(3L, "iPhone", new BigDecimal("6999.00"), "Electronics");

        Customer alice = new Customer(1L, "Alice", "GOLD");
        Customer bob = new Customer(2L, "Bob", "SILVER");

        List<Order> orders = List.of(
            new Order(1L, alice, List.of(
                new OrderItem(book, 2),
                new OrderItem(laptop, 1)
            ), LocalDateTime.now().minusDays(5), "COMPLETED"),
            new Order(2L, bob, List.of(
                new OrderItem(phone, 1)
            ), LocalDateTime.now().minusDays(3), "COMPLETED"),
            new Order(3L, alice, List.of(
                new OrderItem(book, 1)
            ), LocalDateTime.now().minusDays(1), "PENDING")
        );

        OrderService service = new OrderService(orders);

        // 使用查询构建器
        OrderService.OrderQueryBuilder builder = new OrderService.OrderQueryBuilder()
            .where(o -> "COMPLETED".equals(o.status()))
            .where(o -> o.customer().tier().equals("GOLD"))
            .orderBy(Comparator.comparing(Order::total).reversed())
            .limit(10);

        List<Order> goldCompleted = service.queryOrders(
            builder.buildPredicate(),
            builder.buildSorter(),
            builder.buildLimit()
        );
        System.out.println("Gold 客户已完成订单: " + goldCompleted.size());

        // 分类收入
        System.out.println("分类收入: " + service.revenueByCategory());

        // Top 客户
        System.out.println("Top 客户: " + service.topCustomers(5));
    }
}

5. 对比分析:函数式 vs 命令式

5.1 代码风格对比

场景:统计每个部门的平均薪资,并按薪资降序排列。

// 命令式风格
Map<String, List<Employee>> byDept = new HashMap<>();
for (Employee e : employees) {
    byDept.computeIfAbsent(e.getDepartment(), k -> new ArrayList<>()).add(e);
}
Map<String, Double> avgSalary = new HashMap<>();
for (Map.Entry<String, List<Employee>> entry : byDept.entrySet()) {
    double sum = 0;
    for (Employee e : entry.getValue()) {
        sum += e.getSalary();
    }
    avgSalary.put(entry.getKey(), sum / entry.getValue().size());
}
List<Map.Entry<String, Double>> sorted = new ArrayList<>(avgSalary.entrySet());
sorted.sort((e1, e2) -> Double.compare(e2.getValue(), e1.getValue()));

// 函数式风格
List<Map.Entry<String, Double>> result = employees.stream()
    .collect(Collectors.groupingBy(
        Employee::getDepartment,
        Collectors.averagingDouble(Employee::getSalary)
    ))
    .entrySet().stream()
    .sorted(Map.Entry.<String, Double>comparingByValue().reversed())
    .collect(Collectors.toList());

对比维度:

维度命令式函数式
代码行数~15 行~6 行
抽象层级怎么做(How)做什么(What)
可变状态多处修改 byDept、avgSalary无可变状态
并行化需手工同步仅改 .parallelStream()
可读性直白但冗长简洁但需熟悉 API
性能单次遍历效率高有中间对象开销

5.2 性能对比:Stream vs for 循环

// 测试数据:1,000,000 个整数
List<Integer> data = IntStream.rangeClosed(1, 1_000_000).boxed().collect(Collectors.toList());

// 1. 传统 for 循环
long sumFor = 0;
for (int n : data) {
    if (n % 2 == 0) {
        sumFor += n * n;
    }
}

// 2. 串行 Stream
long sumStream = data.stream()
    .filter(n -> n % 2 == 0)
    .mapToLong(n -> (long) n * n)
    .sum();

// 3. 并行 Stream
long sumParallel = data.parallelStream()
    .filter(n -> n % 2 == 0)
    .mapToLong(n -> (long) n * n)
    .sum();

性能数据(JMH 基准测试,吞吐量 ops/ms):

数据规模for 循环串行 Stream并行 Stream
1,0005,2001,800450
10,00052018045
100,00052185.5
1,000,0005.21.80.9

结论:

  1. 小数据量:for 循环最快,Stream 有装箱与流水线开销。
  2. 大数据量:并行 Stream 接近 for 循环,但需考虑并行开销。
  3. 原始类型流(IntStream)能显著降低装箱开销,性能接近 for 循环。

5.3 Lambda vs 匿名内部类

// 匿名内部类
Comparator<String> cmp1 = new Comparator<String>() {
    @Override
    public int compare(String s1, String s2) {
        return Integer.compare(s1.length(), s2.length());
    }
};

// Lambda
Comparator<String> cmp2 = (s1, s2) -> Integer.compare(s1.length(), s2.length());

底层差异:

维度匿名内部类Lambda
编译产物生成独立 .class 文件仅生成合成方法
类加载每次使用需加载类首次调用动态生成
this 引用指向匿名类实例指向外围类实例
内存占用每个 new 创建新实例单例(无捕获)或多例(有捕获)
序列化可序列化默认不可(避免泄露)

5.4 函数式接口 vs 自定义接口

// 方案 A:使用标准函数式接口
Function<User, String> getName = User::getName;

// 方案 B:自定义接口
@FunctionalInterface
interface UserToString extends Function<User, String> {
    // 可添加默认方法
}
UserToString getName2 = User::getName;

何时使用自定义接口:

  1. 需要附加语义(如 UserToString 表达”将用户转为字符串”的领域语义)。
  2. 需要扩展默认方法(如组合、校验)。
  3. 需要在 API 文档中明确类型约束。
  4. 自定义接口的缺点:与标准库互操作需适配。

5.5 Stream vs Reactive Streams

维度StreamReactor Flux/Mono
数据模型同步、有限异步、可无限
背压不支持支持
错误传播抛异常onError 信号
时间维度即时时间感知(interval、delay)
适用场景集合处理异步 I/O、事件流

6. 陷阱与反模式

6.1 反模式:Stream 重复消费

// 反模式:Stream 只能消费一次
Stream<String> stream = Stream.of("a", "b", "c");
stream.forEach(System.out::println);
stream.map(String::toUpperCase).forEach(System.out::println);  // IllegalStateException!

// 正确做法:使用 Supplier 延迟创建
Supplier<Stream<String>> streamSupplier = () -> Stream.of("a", "b", "c");
streamSupplier.get().forEach(System.out::println);
streamSupplier.get().map(String::toUpperCase).forEach(System.out::println);

6.2 反模式:并行流共享可变状态

// 反模式:并行流修改共享集合
List<Integer> results = new ArrayList<>();
IntStream.rangeClosed(1, 1000).parallel()
    .forEach(results::add);  // ArrayList 非线程安全,可能丢失数据

// 正确做法 1:使用线程安全集合
List<Integer> safeResults = Collections.synchronizedList(new ArrayList<>());
IntStream.rangeClosed(1, 1000).parallel()
    .forEach(safeResults::add);

// 正确做法 2:使用 collect(推荐)
List<Integer> collected = IntStream.rangeClosed(1, 1000).parallel()
    .boxed()
    .collect(Collectors.toList());

// 正确做法 3:使用 reduce
List<Integer> reduced = IntStream.rangeClosed(1, 1000).parallel()
    .boxed()
    .reduce(
        new ArrayList<>(),
        (list, n) -> { list.add(n); return list; },
        (l1, l2) -> { l1.addAll(l2); return l1; }
    );

6.3 反模式:Lambda 中的副作用

// 反模式:Lambda 修改外部可变状态
List<String> names = new ArrayList<>();
users.stream()
    .filter(u -> u.getAge() > 18)
    .forEach(u -> names.add(u.getName()));  // 副作用

// 正确做法:使用 collect
List<String> adultNames = users.stream()
    .filter(u -> u.getAge() > 18)
    .map(User::getName)
    .collect(Collectors.toList());

6.4 反模式:过度嵌套的 Stream

// 反模式:过度嵌套,可读性差
List<String> result = data.stream()
    .flatMap(a -> a.getSubItems().stream()
        .flatMap(b -> b.getTags().stream()
            .filter(tag -> tag.length() > 3)
            .map(tag -> tag.toUpperCase())))
    .distinct()
    .collect(Collectors.toList());

// 正确做法:抽取中间方法
List<String> result = data.stream()
    .map(this::extractTags)  // 抽取为方法
    .flatMap(List::stream)
    .filter(tag -> tag.length() > 3)
    .map(String::toUpperCase)
    .distinct()
    .collect(Collectors.toList());

6.5 反模式:在 forEach 中执行 I/O

// 反模式:forEach 中执行数据库操作
users.stream().forEach(u -> {
    saveToDatabase(u);  // 串行 I/O,性能差
});

// 正确做法 1:批量处理
List<User> batch = users.stream().collect(Collectors.toList());
saveAllToDatabase(batch);

// 正确做法 2:并行流(如果 I/O 可并行)
users.parallelStream().forEach(this::saveToDatabase);
// 或使用 CompletableFuture
List<CompletableFuture<Void>> futures = users.stream()
    .map(u -> CompletableFuture.runAsync(() -> saveToDatabase(u)))
    .collect(Collectors.toList());
futures.forEach(CompletableFuture::join);

6.6 反模式:使用 peek 修改状态

// 反模式:peek 用于副作用(语义不明确)
List<User> users = rawUsers.stream()
    .peek(u -> u.setProcessed(true))  // 修改对象状态
    .collect(Collectors.toList());

// 正确做法:使用 map 进行转换
List<User> users = rawUsers.stream()
    .map(u -> u.withProcessed(true))  // 返回新对象
    .collect(Collectors.toList());

6.7 反模式:忽视受检异常

// 反模式:Lambda 中抛出受检异常
files.stream()
    .map(Files::readString)  // 抛出 IOException,编译错误!

// 正案 1:包装为非受检异常
files.stream()
    .map(file -> {
        try {
            return Files.readString(file);
        } catch (IOException e) {
            throw new UncheckedIOException(e);
        }
    })
    .collect(Collectors.toList());

// 方案 2:抽取工具方法
files.stream()
    .map(this::readStringSafely)
    .collect(Collectors.toList());

// 方案 3:使用 Either Monad(需要 Vavr 等库)

6.8 反模式:Collectors.toMap 键冲突

// 反模式:toMap 遇到重复键抛异常
Map<String, Integer> nameToAge = users.stream()
    .collect(Collectors.toMap(User::getName, User::getAge));  // DuplicateKeyException

// 正确做法 1:提供合并函数
Map<String, Integer> nameToAge = users.stream()
    .collect(Collectors.toMap(
        User::getName,
        User::getAge,
        (a, b) -> a  // 保留第一个
    ));

// 正确做法 2:分组为 List
Map<String, List<Integer>> nameToAges = users.stream()
    .collect(Collectors.groupingBy(
        User::getName,
        Collectors.mapping(User::getAge, Collectors.toList())
    ));

6.9 反模式:Optional 滥用

// 反模式 1:Optional 作为字段
class User {
    private Optional<String> nickname;  // 序列化问题,内存浪费
}

// 反模式 2:Optional 作为参数
void setName(Optional<String> name) {  // 增加调用复杂度
    this.name = name.orElse("匿名");
}

// 正确做法:使用方法重载或 nullable
void setName(String name) {
    this.name = name != null ? name : "匿名";
}
void setName() {
    setName("匿名");
}

6.10 反模式:方法引用过度使用

// 反模式:方法引用降低可读性
list.stream()
    .map(Object::toString)
    .map(String::trim)
    .map(String::toLowerCase)
    .forEach(System.out::println);

// 正确做法:复杂场景使用 Lambda 提升可读性
list.stream()
    .map(obj -> obj.toString().trim().toLowerCase())
    .forEach(System.out::println);

7. 工程实践:函数式编程的项目落地

7.1 项目结构建议

flowchart TD
    T0["src/main/java/com/example/"]
    T1["domain/                  # 领域模型(不可变值对象)"]
    T2["User.java            # record User(...)"]
    T3["Order.java"]
    T4["service/                 # 函数式服务"]
    T5["UserService.java"]
    T6["OrderService.java"]
    T7["functional/              # 自定义函数式接口"]
    T8["Transformer.java"]
    T9["ThrowingFunction.java  # 处理受检异常"]
    T10["collector/               # 自定义收集器"]
    T11["PercentileCollector.java"]
    T12["util/                    # 函数式工具"]
    T13["StreamUtils.java"]
    T0 --> T1
    T3 --> T4
    T6 --> T7
    T9 --> T10
    T11 --> T12
    T12 --> T13

7.2 处理受检异常的工具

@FunctionalInterface
public interface ThrowingFunction<T, R, E extends Exception> {
    R apply(T t) throws E;

    static <T, R, E extends Exception> Function<T, R> unchecked(
            ThrowingFunction<T, R, E> function) {
        return t -> {
            try {
                return function.apply(t);
            } catch (Exception e) {
                throw new RuntimeException(e);
            }
        };
    }
}

// 使用
List<String> contents = files.stream()
    .map(ThrowingFunction.unchecked(Files::readString))
    .collect(Collectors.toList());

7.3 函数式缓存

public class MemoizedFunction<T, R> implements Function<T, R> {
    private final Function<T, R> delegate;
    private final Map<T, R> cache = new ConcurrentHashMap<>();

    public MemoizedFunction(Function<T, R> delegate) {
        this.delegate = delegate;
    }

    @Override
    public R apply(T t) {
        return cache.computeIfAbsent(t, delegate);
    }
}

// 使用
Function<Integer, Integer> slowSquare = x -> {
    try { Thread.sleep(100); } catch (InterruptedException e) {}
    return x * x;
};
Function<Integer, Integer> memoized = new MemoizedFunction<>(slowSquare);

long start = System.currentTimeMillis();
memoized.apply(5);  // 100ms
memoized.apply(5);  // <1ms(缓存命中)
memoized.apply(6);  // 100ms

7.4 函数式重试

public class Retry {
    public static <T> Supplier<T> withRetry(Supplier<T> action, int maxRetries) {
        return () -> {
            RuntimeException lastException = null;
            for (int i = 0; i <= maxRetries; i++) {
                try {
                    return action.get();
                } catch (RuntimeException e) {
                    lastException = e;
                    if (i < maxRetries) {
                        try {
                            Thread.sleep((long) Math.pow(2, i) * 100);
                        } catch (InterruptedException ie) {
                            Thread.currentThread().interrupt();
                        }
                    }
                }
            }
            throw lastException;
        };
    }
}

// 使用
String result = Retry.withRetry(() -> callRemoteService(), 3).get();

7.5 函数式配置

public class FunctionalRouter {
    private final Map<String, Function<Request, Response>> routes = new HashMap<>();

    public FunctionalRouter route(String path, Function<Request, Response> handler) {
        routes.put(path, handler);
        return this;
    }

    public Response handle(Request request) {
        return routes.getOrDefault(request.path(), req ->
            new Response(404, "Not Found")).apply(request);
    }
}

// 使用
FunctionalRouter router = new FunctionalRouter()
    .route("/users", this::listUsers)
    .route("/orders", this::listOrders)
    .route("/health", req -> new Response(200, "OK"));

7.6 测试函数式代码

class FunctionalTest {
    @Test
    void testPipeline() {
        // 给定
        List<User> users = List.of(
            new User("Alice", 25),
            new User("Bob", 17),
            new User("Charlie", 30)
        );

        // 当
        List<String> result = users.stream()
            .filter(u -> u.age() > 18)
            .sorted(Comparator.comparing(User::name))
            .map(User::name)
            .collect(Collectors.toList());

        // 则
        assertEquals(List.of("Alice", "Charlie"), result);
    }

    @Test
    void testPureFunction() {
        Function<Integer, Integer> square = x -> x * x;

        // 引用透明性:多次调用结果相同
        assertEquals(25, square.apply(5));
        assertEquals(25, square.apply(5));
        assertEquals(25, square.apply(5));
    }

    @Test
    void testParallelStream() {
        List<Integer> data = IntStream.rangeClosed(1, 10000).boxed().collect(Collectors.toList());

        long sum = data.parallelStream()
            .mapToLong(Integer::longValue)
            .sum();

        assertEquals(50005000L, sum);
    }
}

8. 案例研究:主流框架的函数式实践

8.1 Spring Framework 的函数式风格

Spring 5+ 大量采用函数式风格:

// 函数式路由(Spring WebFlux)
@Configuration
public class FunctionalRoutes {
    @Bean
    public RouterFunction<ServerResponse> userRoutes(UserHandler handler) {
        return RouterFunctions.route()
            .GET("/users/{id}", handler::getUser)
            .GET("/users", handler::listUsers)
            .POST("/users", handler::createUser)
            .build();
    }
}

// 函数式 Bean 注册(Spring 5+)
GenericApplicationContext context = new GenericApplicationContext();
context.registerBean(UserService.class);
context.refresh();

// Spring Security 函数式 DSL
@Configuration
@EnableWebSecurity
public class SecurityConfig {
    @Bean
    public SecurityFilterChain filterChain(HttpSecurity http) throws Exception {
        return http
            .authorizeHttpRequests(auth -> auth
                .requestMatchers("/public/**").permitAll()
                .anyRequest().authenticated()
            )
            .formLogin(form -> form.loginPage("/login"))
            .build();
    }
}

8.2 Reactor 的函数式响应式

// Reactor 的 Mono/Flux 完全基于函数式组合
Mono<User> getUser = userRepository.findById(id)
    .map(user -> user.withLastSeen(Instant.now()))
    .flatMap(this::enrichUserData)
    .onErrorResume(e -> Mono.just(User.DEFAULT))
    .switchIfEmpty(Mono.error(new UserNotFoundException(id)));

8.3 Vavr 的函数式增强

Vavr 库为 Java 提供了更完整的函数式工具:

// Vavr 的 Try 处理异常
Try<Integer> result = Try.of(() -> 1 / 0)
    .recover(ArithmeticException.class, e -> 0)
    .map(x -> x + 1);
System.out.println(result.get());  // 1

// Vavr 的 Either
Either<String, Integer> either = compute()
    .map(value -> value * 2)
    .filter(value -> value > 0, () -> "非正数");

// Vavr 的 Pattern Matching
String description = Match(user).of(
    Case($(u -> u.age() < 18), "未成年"),
    Case($(u -> u.age() >= 65), "老年"),
    Case($(), "成年")
);

8.4 JUnit 5 的函数式断言

// JUnit 5 的函数式断言
assertAll(
    () -> assertEquals(4, list.size()),
    () -> assertTrue(list.contains("Alice")),
    () -> assertNotNull(list.get(0))
);

// 动态测试
@TestFactory
Stream<DynamicTest> dynamicTests() {
    return Stream.of(1, 2, 3, 4, 5)
        .map(n -> DynamicTest.dynamicTest(
            "测试 " + n,
            () -> assertTrue(n > 0)
        ));
}

8.5 Collectors 在 Apache Commons 的应用

// Apache Commons 的 MultiValuedMap 与 Collectors
ListMultiValuedMap<String, Integer> multiMap = users.stream()
    .collect(Collectors.toMap(
        User::getDepartment,
        User::getSalary,
        (a, b) -> a,  // 不合并
        ArrayList::new  // 工厂方法
    ));

// Guava 的 ImmutableList 收集
ImmutableList<User> immutable = users.stream()
    .collect(ImmutableList.toImmutableList());

9.1 基础题(记忆与理解)

  1. 列举 java.util.function 包中至少 5 个核心函数式接口及其函数描述符。
  2. 解释 @FunctionalInterface 注解的作用,并说明是否必须标注。
  3. 描述 Lambda 表达式与匿名内部类在字节码层面的主要差异。
  4. 说明 Stream 中间操作与终端操作的区别,各举 3 个例子。
  5. 解释 flatMap 与 map 的语义差异,并各给出一个适用场景。

应用题知识点讲解

  1. 使用 Stream API 实现以下功能:给定一个字符串列表,统计每个单词的出现频率,并按频率降序输出前 10 个。
  2. 实现一个自定义 Collector,将订单流按月份分组,并计算每月的总金额、平均金额、最大金额。
  3. 使用 parallelStream 并行处理 100 万条日志记录,过滤出错误级别日志并写入文件。注意处理线程安全。
  4. 设计一个基于函数组合的验证器,支持对用户对象进行多重校验(用户名非空、密码长度、邮箱格式)。
  5. 使用 IntStream 计算圆周率 π 的近似值(莱布尼茨级数),并对比不同项数的精度。

9.3 分析题

  1. 以下代码有什么问题?请指出并修复:

    List<Integer> list = new ArrayList<>();
    IntStream.range(0, 100).parallel().forEach(list::add);
  2. 解释为何以下代码的输出可能少于 5:

    long count = Stream.of(1, 2, 3, 4, 5)
        .peek(n -> { if (n == 3) throw new RuntimeException(); })
        .count();
  3. 分析以下两段代码的性能差异,并说明何时选择哪种:

    // 代码 A
    int sum = 0;
    for (int n : data) sum += n * n;
    
    // 代码 B
    int sum = data.stream().mapToInt(n -> n * n).sum();

9.4 设计题

  1. 设计一个函数式的任务调度器,支持任务依赖、超时控制、错误重试。要求所有操作通过函数组合完成。
  2. 设计一个基于 Collector 的报表生成器,支持多维度聚合(按时间、地区、产品分类),并能导出为 CSV/JSON 格式。
  3. 重构以下命令式代码为函数式风格,保证功能等价:
    Map<String, Integer> result = new HashMap<>();
    for (User u : users) {
        if (u.getAge() > 18 && "CN".equals(u.getCountry())) {
            String key = u.getCity();
            result.put(key, result.getOrDefault(key, 0) + 1);
        }
    }

10.1 学术论文

  1. Church, A. (1936). An unsolvable problem of elementary number theory. American Journal of Mathematics.
  2. Backus, J. (1978). Can Programming Be Liberated from the von Neumann Style? Communications of the ACM.
  3. Hughes, J. (1989). Why Functional Programming Matters. Computer Journal.
  4. Wadler, P. (1990). Theorems for Free! FPCA.
  5. Goetz, B. (2013). Translation of Lambda Expressions. Java Specification Request 335.

10.2 规范与标准

  1. JSR 335: Lambda Expressions for the Java Programming Language.
  2. JEP 126: Lambda Expressions and Virtual Extension Methods.
  3. JEP 269: Convenience Factory Methods for Collections.
  4. JEP 395: Records (JDK 16).
  5. The Java Language Specification (JLS), Chapter 15.27: Lambda Expressions.

10.3 书籍

  1. Urma, R. G., Fusco, M., & Mycroft, A. (2018). Modern Java in Action. Manning.
  2. Goetz, B., et al. (2006). Java Concurrency in Practice. Addison-Wesley.
  3. Okasaki, C. (1998). Purely Functional Data Structures. Cambridge University Press.
  4. Lipovača, M. (2011). Learn You a Haskell for Great Good! No Starch Press.
  5. Abelson, H., & Sussman, G. J. (1996). Structure and Interpretation of Computer Programs (2nd ed.). MIT Press.

11.1 函数式编程理论

  • λ 演算入门:Hankin, C. (2004). An Introduction to Lambda Calculi for Computer Scientists.
  • 范畴论与函数式编程:Bartosz Milewski 的 Category Theory for Programmers 系列。
  • 类型系统:Pierce, B. C. (2002). Types and Programming Languages. MIT Press.

11.2 Java 函数式进阶

  • Java 函数式编程深度:Subramaniam, V. (2019). Functional Programming in Java (2nd ed.). Pragmatic Bookshelf.
  • Stream 内部机制:OpenJDK java.util.stream 包源码与文档。
  • invokedynamic 深入:Rose, J. (2009). Bytecodes meet Combinators: invokedynamic and the JVM.

11.3 相关技术

  • 响应式编程:Reactive Manifesto 与 Project Reactor 文档。
  • 模式匹配:JEP 441 (Pattern Matching for switch) 与 Scala 模式匹配对比。
  • 不可变数据结构:Clojure 的持久化数据结构与 Vavr 的实现。

11.4 实战项目

  • 函数式 Web 框架:Spring WebFlux、Vert.x、Javalin 的对比研究。
  • 函数式数据处理:Apache Spark 的 RDD API 与 Java Stream 的对比。
  • 函数式 DSL 设计:Gradle Kotlin DSL、Spock Framework 的设计思路。

附录 A:java.util.function 核心接口速查表

接口函数描述符主要用途常用方法
Function<T, R>T -> R转换apply, compose, andThen
Consumer<T>T -> void消费accept, andThen
Supplier<T>() -> T提供get
Predicate<T>T -> boolean判断test, and, or, negate
BiFunction<T, U, R>(T, U) -> R双参转换apply, andThen
BiConsumer<T, U>(T, U) -> void双参消费accept, andThen
BiPredicate<T, U>(T, U) -> boolean双参判断test, and, or, negate
UnaryOperator<T>T -> T一元操作apply (继承 Function)
BinaryOperator<T>(T, T) -> T二元操作apply, minBy, maxBy
IntFunction<R>int -> Rint 转换apply
IntPredicateint -> booleanint 判断test
IntConsumerint -> voidint 消费accept
IntSupplier() -> intint 提供getAsInt
ToIntFunction<T>T -> int转 intapplyAsInt

附录 B:Stream 操作分类速查

中间操作(Intermediate Operations)

操作类型描述有状态
filter无状态过滤元素否
map无状态转换元素否
mapToInt/Long/Double无状态转换为原始类型否
flatMap无状态展平嵌套流否
flatMapToInt/Long/Double无状态展平为原始流否
peek无状态调试否
distinct有状态去重是
sorted有状态排序是
limit有状态短路截取前 N 个是
skip有状态跳过前 N 个是
takeWhile (Java 9+)有状态短路取满足条件的前缀是
dropWhile (Java 9+)有状态丢弃满足条件的前缀是

终端操作(Terminal Operations)

操作类型描述
forEach消费遍历元素
forEachOrdered消费按顺序遍历
toArray转换转为数组
reduce归约归约为单值
collect归约收集为集合
toList (Java 16+)归约简化收集
min归约最小值
max归约最大值
count归约计数
anyMatch短路任一匹配
allMatch短路全部匹配
noneMatch短路无匹配
findFirst短路第一个元素
findAny短路任意元素
iterator转换转为迭代器

附录 C:函数式编程术语对照表

英文术语中文译名简要说明
Pure Function纯函数无副作用、引用透明
Higher-Order Function高阶函数接受或返回函数的函数
First-class Function一等公民函数可赋值、传参、返回
Lambda ExpressionLambda 表达式匿名函数字面量
Closure闭包捕获外部变量的 Lambda
Currying柯里化多参函数转为一元函数链
Partial Application部分应用固定部分参数
Monad单子提供 flatMap 的类型构造子
Functor函子提供 map 的类型构造子
Referential Transparency引用透明性表达式可替换为其值
Side Effect副作用修改外部状态
Immutability不可变性创建后不可修改
Lazy Evaluation惰性求值按需计算
Tail Recursion尾递归递归调用为最后操作
Pattern Matching模式匹配类型驱动的分支
Algebraic Data Type代数数据类型通过组合构造的类型

结语

函数式编程不仅是语法糖,更是一种思维范式。从 λ 演算的数学根基,到 invokedynamic 的字节码实现,再到 Stream 的惰性求值与并行调度,Java 函数式编程展现了”工程务实”与”理论严谨”的平衡。

本节以 12 个章节、10 个完整代码示例、10 个反模式剖析、5 个框架案例研究,系统性地呈现了 Java 函数式编程的全貌。读者通过本节的学习,应能:

  1. 理解原理:从字节码层面理解 Lambda 与 Stream 的实现机制。
  2. 掌握工具:熟练使用 java.util.function、Stream、Collectors、Optional。
  3. 规避陷阱:识别并行流的共享状态问题、Stream 重复消费、副作用滥用等反模式。
  4. 工程落地:将函数式思维应用于领域建模、错误处理、测试设计。

在虚拟线程与模式匹配普及的时代,函数式编程将继续作为 Java 生态的核心范式。掌握函数式思维,不仅是写出更简洁的代码,更是构建可推理、可测试、可并行系统的思维基石。希望本节能为读者的函数式之旅提供一个坚实的起点。

Function 函数

基本写法:定义 Function Function<<入参类型>, <返回类型>> <变量> = <Lambda>;

// 接收一个参数返回一个结果
Function<String, Integer> len = s -> s.length();

基本写法:应用函数 <function>.apply(<参数>);

// 执行函数并返回结果
int n = len.apply("hello");

基本写法:复合函数 <f1>.andThen(<f2>);

// 先执行 f1 再执行 f2
Function<String, String> upper = s -> s.toUpperCase();
Function<String, String> bang = upper.andThen(s -> s + "!");

基本写法:反向复合 <f2>.compose(<f1>);

// 先执行 f1 再执行 f2
Function<Integer, Integer> add1 = x -> x + 1;
Function<Integer, Integer> mul2 = x -> x * 2;
Function<Integer, Integer> h = mul2.compose(add1);

Predicate 断言

基本写法:定义 Predicate Predicate<<类型>> <变量> = <Lambda>;

// 返回 boolean 的判断函数
Predicate<String> nonEmpty = s -> s != null && !s.isEmpty();

基本写法:测试断言 <predicate>.test(<参数>);

// 对输入进行判断
boolean ok = nonEmpty.test("abc");

基本写法:与运算 <p1>.and(<p2>);

// 两个断言都为真
Predicate<Integer> positive = x -> x > 0;
Predicate<Integer> even = x -> x % 2 == 0;
Predicate<Integer> posEven = positive.and(even);

基本写法:取反 <predicate>.negate();

// 取反断言
Predicate<String> isEmpty = nonEmpty.negate();

Consumer 消费者

基本写法:定义 Consumer Consumer<<类型>> <变量> = <Lambda>;

// 接收一个参数无返回值
Consumer<String> printer = s -> System.out.println(s);

基本写法:链式消费 <c1>.andThen(<c2>);

// 先执行 c1 再执行 c2
Consumer<String> c1 = s -> System.out.print("[" + s);
Consumer<String> c2 = s -> System.out.println("]");
Consumer<String> chained = c1.andThen(c2);

Supplier 供应者

基本写法:定义 Supplier Supplier<<类型>> <变量> = <Lambda>;

// 无参数返回一个结果
Supplier<String> now = () -> java.time.Instant.now().toString();

基本写法:获取值 <supplier>.get();

// 执行并返回结果
String v = now.get();

BiFunction 双参函数

基本写法:定义 BiFunction BiFunction<<T>, <U>, <R>> <变量> = <Lambda>;

// 接收两个参数返回一个结果
BiFunction<Integer, Integer, Integer> add = (a, b) -> a + b;

原始类型特化

基本写法:IntFunction IntFunction<<R>> <变量> = <Lambda>;

// 接收 int 返回 R
IntFunction<String> idx = i -> "#" + i;

基本写法:IntPredicate IntPredicate <变量> = <Lambda>;

// 接收 int 返回 boolean
IntPredicate positive = i -> i > 0;