乐于分享
好东西不私藏

Flink源码学习系列课程 01:getExecutionEnvironment 环境自动发现机制

Flink源码学习系列课程 01:getExecutionEnvironment 环境自动发现机制

核心文件StreamExecutionEnvironment.javaUtils.javaDefaultExecutorServiceLoader.java

从哪里调用

这一课是 WordCount 程序的第一行代码——你在 IDE 里写下的起点。每次你写 StreamExecutionEnvironment.getExecutionEnvironment(),Flink 就开始了一轮"环境侦探"。

这一课解决什么问题

你在 WordCount 里写了:

final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

同样的代码,在 IDE 里直接 Run 和在 flink run 命令行提交,Flink 会走完全不同的执行路径——IDE 里是本地多线程模拟,集群上是真正的分布式执行。

我们要搞清楚:这一行代码如何做到自动识别运行环境? 你可能会想,是不是有某个 if/else 判断?其实比这精巧得多——它用了两级工厂 + ThreadLocal 的设计。

第一关:两级工厂

ThreadLocal 优先,静态变量兜底

打开 StreamExecutionEnvironment.java,跳到第 154-158 行,看这两个字段:

privatestatic StreamExecutionEnvironmentFactory contextEnvironmentFactory = null;

privatestaticfinal ThreadLocal<StreamExecutionEnvironmentFactory>
        threadLocalContextEnvironmentFactory = new ThreadLocal<>();

为什么有两个变量? 我们来推理一下场景:contextEnvironmentFactory 是静态变量,JVM 全局共享。如果多个线程同时调用 getExecutionEnvironment(),它们看到的是同一个值。但 flink run 的 CLI 在 main 线程注入工厂,而用户的 main() 方法可能在另一个线程被调用——如果用静态变量,跨线程的值传递依赖内存可见性的时序。

threadLocalContextEnvironmentFactory 是 ThreadLocal,每个线程独立。不同线程可以有不同的工厂,互不干扰。

所以两级的设计逻辑是:ThreadLocal 先看当前线程有没有工厂,没有再去看 JVM 全局的静态变量,再没有就说明没人注入——走本地模式。

关键细节①contextEnvironmentFactory 没有 volatile 修饰。你可能会担心多线程可见性问题。但 Flink 的设计保证了"写入一定在读取之前"——initializeContextEnvironment() 在 CLI 启动阶段(单线程)调用,之后才是用户代码的多线程执行。这种"单线程写、多线程读"的模式不需要 volatile。

第二关:核心入口

Java 8 Optional 的链式美学

跳到 StreamExecutionEnvironment.java 第 2194-2198 行:

publicstatic StreamExecutionEnvironment getExecutionEnvironment(Configuration configuration){
return Utils.resolveFactory(
                threadLocalContextEnvironmentFactory,  // ThreadLocal 优先
                contextEnvironmentFactory              // 静态变量兜底
            )
            .map(factory -> factory.createExecutionEnvironment(configuration))
            .orElseGet(() -> StreamExecutionEnvironment.createLocalEnvironment(configuration));
}

一眼看过去,就三行链式调用。我们来拆解:

  • Utils.resolveFactory(...) 返回 Optional<StreamExecutionEnvironmentFactory>
  • .map(...):如果 Optional 有值,调工厂的 createExecutionEnvironment(),得到具体环境
  • .orElseGet(...):如果 Optional 为空(两个工厂都是 null),直接创建本地环境

为什么这样设计? Flink 实现了"约定优于配置"——你不需要在代码里写"我要本地模式"还是"我要集群模式"。IDE 里跑默认就是本地;CLI 环境下工厂已经被注入,就走集群模式。同一份代码在不同环境下自动适配,这对开发者体验来说是巨大的提升。

实践环节

在你的 IDE 里,Ctrl/Cmd + 点击resolveFactory,跳进去看:

// Utils.java:331-337
publicstatic <T> Optional<T> resolveFactory(
        ThreadLocal<T> threadLocalFactory, @Nullable T staticFactory)
{
final T localFactory = threadLocalFactory.get();
final T factory = localFactory == null ? staticFactory : localFactory;
return Optional.ofNullable(factory);
}

这里有个让你会心一笑的设计——这个方法放在 flink-core 的 Utils 类,而不是 StreamExecutionEnvironment 内部。为什么?因为同一个模式在 Flink 的多处被复用(ClusterClientServiceLoader 的文件注释里直接写了 "almost identical",几乎是复制粘贴的)。好的工具方法要放在能被广泛引用的位置,这是代码复用的基本嗅觉。

第三关(关键!):工厂的注入与清理

ThreadLocal 内存泄漏的预防

翻到 StreamExecutionEnvironment.java 第 2368-2376 行:

protectedstaticvoidinitializeContextEnvironment(StreamExecutionEnvironmentFactory ctx){
    contextEnvironmentFactory = ctx;
    threadLocalContextEnvironmentFactory.set(ctx);
}

protectedstaticvoidresetContextEnvironment(){
    contextEnvironmentFactory = null;
    threadLocalContextEnvironmentFactory.remove();
}

谁在调用它们? 在你的 IDE 里 Ctrl/Cmd + 点击initializeContextEnvironment,查看调用方——在 flink-clients 模块的 CliFrontend 中:

// 简化示意
publicstaticvoidmain(String[] args){
// 1. 解析参数,确定目标集群
// 2. 创建对应的 StreamExecutionEnvironmentFactory
    StreamExecutionEnvironmentFactory factory = ...;
// 3. 注入
    StreamExecutionEnvironment.initializeContextEnvironment(factory);
try {
// 4. 执行用户的 main() 方法
        userMain(args);
    } finally {
// 5. 清理——不论成功失败,一定会执行!
        StreamExecutionEnvironment.resetContextEnvironment();
    }
}

关键细节②resetContextEnvironment() 为什么放在 finally 块里?不是因为代码要"优雅",而是如果不清理,ThreadLocal 会一直持有工厂的引用,导致这个线程无法被 GC——在应用服务器这种长生命周期的场景下就是内存泄漏。用 finally 保证即使用户代码抛异常,环境也会被清理。

实践环节

在 resetContextEnvironment() 处打个断点,然后跑一次 flink run WordCount.jar,你会发现它确实在 finally 中被调用——两件事都做了:contextEnvironmentFactory = null 和 threadLocalContextEnvironmentFactory.remove()。前者释放静态引用,后者清除当前线程的 ThreadLocal 变量。

第四关:本地模式的并行度

你的 CPU 核数决定一切

StreamExecutionEnvironment.java 第 160-161 行:

privatestaticint defaultLocalParallelism = Runtime.getRuntime().availableProcessors();

这意味着什么? 如果你用的是一台 8 核机器,env.fromSource(...).flatMap(...).sum(...).print() 会创建 8 个并行实例,每个实例在独立线程中执行(MiniCluster 用多线程模拟分布式)。

你在 IDE 里跑 WordCount,开 8 个线程处理一个小文件——本质上就是本地多线程模拟分布式。如果你想调试时只看一个线程的逻辑,就 env.setParallelism(1)

调用链追踪

用户代码:
  StreamExecutionEnvironment.getExecutionEnvironment()
    │
    ▼
  getExecutionEnvironment(new Configuration())               // :2178
    │ 传入空的 Configuration(Map<String,String> = {})
    │
    ▼
  getExecutionEnvironment(configuration)                     // :2194
    │
    ├─→ Utils.resolveFactory(threadLocalFactory, staticFactory)  // :331
    │     │
    │     ├─→ threadLocalFactory.get()                       // ThreadLocal 取值
    │     │     返回: StreamExecutionEnvironmentFactory 或 null
    │     │
    │     ├─→ localFactory == null ? staticFactory : localFactory
    │     │     如果 ThreadLocal 有值 → 用 ThreadLocal 的
    │     │     如果 ThreadLocal 是 null → 用静态变量(也可能是 null)
    │     │
    │     └─→ return Optional.ofNullable(factory)
    │           如果 factory 非 null → Optional.of(factory)
    │           如果 factory 是 null → Optional.empty()
    │
    ├─→ .map(factory → factory.createExecutionEnvironment(configuration))
    │     如果 Optional 有值 → 调用工厂方法 → 返回 Remote/Local/其他环境
    │
    └─→ .orElseGet(() → createLocalEnvironment(configuration))
          如果 Optional 为空(两个工厂都是 null)
          创建 LocalStreamEnvironment,并行度 = CPU 核数

边界条件速查

场景
行为
IDE 中直接跑 main()
两个工厂都是 null → createLocalEnvironment(),并行度 = CPU 核数
flink run
 提交
CLI 在 main() 前注入工厂 → 工厂非 null → 创建 RemoteStreamEnvironment
单元测试中跑
取决于测试框架是否注入了工厂,默认走 Local
并发调用 getExecutionEnvironment()
ThreadLocal 保证线程安全,每个线程各自判断
main 中调两次
两个独立对象,各有一套 transformations List

总结

┌──────────────────────────────────────────────────────────────────┐
│  getExecutionEnvironment() 核心要点:                             │
│                                                                   │
│  1. 两级工厂查找:ThreadLocal 优先,静态变量兜底,都为空则 Local    │
│  2. 工厂在 CLI 启动时通过 initializeContextEnvironment() 注入      │
│  3. ThreadLocal 设计提供了线程级别的环境隔离                      │
│  4. resetContextEnvironment() 在 finally 中清理,防止内存泄漏      │
│  5. Java 8 Optional 链式调用:resolve → map → orElseGet            │
│  6. 本地模式默认并行度 = CPU 核数(Runtime.availableProcessors()) │
│  7. 设计模式:静态工厂方法 + 策略模式(环境策略由工厂决定)         │
└──────────────────────────────────────────────────────────────────┘

下一课预告

拿到了 StreamExecutionEnvironment 对象,它的构造函数初始化了哪些关键成员?特别是 transformations List——整个 DAG 的骨架——它是如何为惰性求值做准备的?