核心文件:StreamExecutionEnvironment.java, Utils.java, DefaultExecutorServiceLoader.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 核数
边界条件速查
createLocalEnvironment(),并行度 = CPU 核数 | |
flink run | |
getExecutionEnvironment() | |
总结
┌──────────────────────────────────────────────────────────────────┐
│ 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 的骨架——它是如何为惰性求值做准备的?
夜雨聆风