Flink Plugins 插件机制详解:文件系统与 Metric Reporter 的隔离加载与实战部署
2026/9/20 20:24:53 网站建设 项目流程

Flink Plugins 插件机制详解:文件系统与 Metric Reporter 的隔离加载与实战部署

【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink

Apache Flink 从 1.9 版本引入 Plugins(插件)机制,通过受限的类加载器(ClassLoader)实现代码的严格隔离,让同一个 Flink 发行版可以安全地加载依赖冲突、版本不同的库(例如不同厂商的云文件系统客户端),而无需做类重定位(shading)或强行统一版本。本文以 Flink 官方文档《Plugins》为主体,结合本仓库中flink-core的插件实现源码,系统讲解插件机制的隔离原理、目录结构、plugins目录下的部署实操(文件系统与 Metric Reporter),以及底层插件发现与加载的完整链路,帮助你正确地在生产环境中启用 S3、OSS、Azure、GCS 等文件系统插件和各类指标上报器。

什么是 Plugins:为什么需要严格的类加载器隔离

插件机制的初衷非常明确:通过受限的类加载器实现代码的严格分离。一个插件无法访问其他插件中的类,也无法访问 Flink 中未被显式白名单(whitelist)放行的类。这种严格隔离带来一个关键收益:插件之间可以安全地持有同一库的冲突版本,既不需要在打包 fat jar 时重定位类(relocation),也不需要把依赖收敛到公共版本。

典型场景就是文件系统客户端:flink-s3-fs-hadoop(基于 Hadoop S3A)与flink-azure-fs-hadoop(基于 Hadoop ABFS/WASB)各自依赖不同版本的 Hadoop 生态库,二者如果同时放在lib目录下很容易产生类冲突,而插件机制让它们各自运行在自己的类加载器中,互不干扰。

目前 Flink 中**可插拔(pluggable)**的组件包括:

  • 文件系统(File Systems)
  • 指标上报器(Metric Reporters)

从仓库源码的设计注释看,未来 Connector、Format 甚至用户代码也都计划支持插件化。

隔离与插件目录结构

标准目录布局

插件位于 Flink 发行版的plugins根目录下,每个插件独占一个子目录,子目录内可以放一个或多个 JAR 文件。插件目录的名称是任意的(目录名会成为插件 ID)。标准发行版的布局如下:

flink-dist ├── conf ├── lib ... └── plugins ├── s3 │ ├── aws-credential-provider.jar │ └── flink-s3-fs-hadoop.jar └── azure └── flink-azure-fs-hadoop.jar

这个布局在发行版构建时已经预留好:仓库中 flink-dist/src/main/flink-bin/plugins/README.txt 给出了与上述一致的结构说明——"一个插件文件夹包含该插件所属的全部资源(JAR 文件),插件文件夹的名字即插件 ID"。

每个插件一个 ClassLoader

每个插件都通过自己的类加载器加载,与其他插件完全隔离。因此flink-s3-fs-hadoopflink-azure-fs-hadoop可以依赖不同的、互相冲突的库版本,且无需在打 fat jar 时重定位任何类。

这一隔离机制在源码中有清晰的实现:DefaultPluginManager内部维护了一张Map<String, PluginLoader>pluginId -> PluginLoader),每个插件 ID 对应一个独立的PluginLoaderPluginLoader通过PluginClassLoader(继承自URLClassLoader)加载该插件目录下的所有 JAR,见 flink-core/src/main/java/org/apache/flink/core/plugin/PluginLoader.java 中的createPluginClassLoadercreate方法。

白名单与 SPI 单例:为什么不会出现两份 FileSystem

插件可以访问 Flinklib/目录中被白名单放行的特定包。尤其重要的是:所有必要的服务提供者接口(SPI)都通过系统类加载器加载,从而保证任何时刻都不会存在两个版本的org.apache.flink.core.fs.FileSystem,即使有人在 fat jar 中意外打包了一份也不行。

这个"单例类"要求是严格必要的,因为 Flink 运行时需要一个确定的入口点进入插件。服务类的发现基于 JDK 标准的java.util.ServiceLoader,因此shading 打包时必须保留META-INF/services中的服务定义文件,否则插件在运行时无法被发现。

注意:当前 SPI 体系仍在完善中,从 Flink 核心泄露给插件的类比最终设计要多(源码中的原话是 "Currently, more Flink core classes are still accessible from plugins as we flesh out the SPI system")。编写插件实现时,应尽量只依赖@Public@PublicEvolving标注的接口。

日志框架白名单

此外,最常见的日志框架(如 SLF4J、Log4j 等)也被白名单放行,因此 Flink 核心、插件和用户代码可以统一使用同一套日志配置输出日志。这一点也体现在CoreOptions的类加载模式配置中(详见下文"Parent-first 类加载配置")。

文件系统插件:启用与部署

所有文件系统都是可插拔的,这意味着它们应该以插件方式使用。完整的文件系统列表与用法参见文件系统概览:本地文件系统默认可用;Amazon S3(flink-s3-fs-presto/flink-s3-fs-hadoop)、阿里云 OSS(flink-oss-fs-hadoop)、Azure Blob Storage(flink-azure-fs-hadoop)、Google Cloud Storage(gcs-connector)等外部文件系统均以插件形式提供。

以 S3 为例,启用插件的操作步骤是:在启动 Flink 之前,将opt目录下对应的 JAR 文件复制到发行版plugins目录下的某个子目录中:

mkdir ./plugins/s3-fs-hadoop cp ./opt/flink-s3-fs-hadoop-<version>.jar ./plugins/s3-fs-hadoop/

其中<version>替换为当前 Flink 发行版的实际版本号(即flink-s3-fs-hadoop-{{< version >}}.jar形式的文件名)。命令执行完后,S3 路径s3://<your-bucket>/<endpoint>即可用于读取、写入以及作为 checkpoint 存储,详见 Amazon S3 文档。

仓库中对应的插件源码模块包括 flink-s3-fs-hadoop、flink-s3-fs-presto、flink-azure-fs-hadoop、flink-oss-fs-hadoop、flink-gs-fs-hadoop 等,均位于 flink-filesystems 聚合模块下。

重要警告一:S3 插件只能以插件方式使用

flink-s3-fs-prestoflink-s3-fs-hadoop只能作为插件使用,因为仓库中已经移除了这些插件的类重定位(relocation)。把它们放到lib目录会导致系统启动失败。这一约束在 文件系统概览 中也有说明:Flink 1.9 引入插件机制,1.10 起 S3 插件不再隐藏/重定位类,旧机制(放入lib目录)不可再用;官方建议未来版本的 Flink 将不再支持通过lib目录加载文件系统组件。

重要警告二:凭证提供者需要放入插件目录

由于严格的类隔离,文件系统插件不再能访问lib目录中的凭证提供者(credential provider)。如果使用 S3、OSS 等对象存储需要额外的凭证提供者 JAR(例如自定义的aws-credential-provider.jar),请把它们一并放入对应的插件目录,例如:

plugins/s3-fs-hadoop/ ├── aws-credential-provider.jar └── flink-s3-fs-hadoop.jar

Metric Reporter 插件

Flink 提供的所有 Metric Reporter 都可以作为插件使用,指标上报器的完整配置方式参见指标上报器文档。

发行版构建时,Flink 已经通过 flink-dist/src/main/assemblies/plugins.xml 把开箱即用的指标插件打包进plugins/目录,包括:

插件目录对应 JAR(源码模块)
plugins/metrics-jmx/flink-metrics-jmx(flink-metrics-jmx)
plugins/metrics-graphite/flink-metrics-graphite(flink-metrics-graphite)
plugins/metrics-influx/flink-metrics-influxdb(flink-metrics-influxdb)
plugins/metrics-prometheus/flink-metrics-prometheus(flink-metrics-prometheus)
plugins/metrics-statsd/flink-metrics-statsd(flink-metrics-statsd)
plugins/metrics-datadog/flink-metrics-datadog(flink-metrics-datadog)
plugins/metrics-slf4j/flink-metrics-slf4j(flink-metrics-slf4j)

此外,plugins.xml还展示了插件机制的一个扩展用例——GPU 外部资源插件plugins/external-resource-gpu/,其中除了 JAR 还打包了gpu-discovery-common.shnvidia-gpu-discovery.sh两个发现脚本(见 flink-external-resources/flink-external-resource-gpu),说明插件目录中并非只能放 JAR,也可以放置插件运行所需的资源文件。

源码级剖析:插件发现与加载的完整链路

1. 插件目录如何确定:环境变量与默认值

插件根目录的解析逻辑在 PluginConfig.java:优先读取环境变量FLINK_PLUGINS_DIR,未设置时回退到默认值plugins(相对于 Flink 发行版根目录)。对应的常量定义见 ConfigConstants.java(ENV_FLINK_PLUGINS_DIR/DEFAULT_FLINK_PLUGINS_DIRS)。如果该目录不存在,PluginConfig只记录一条警告日志并返回空,插件系统以"无插件"状态继续启动。

2. 目录扫描:DirectoryBasedPluginFinder

PluginUtils.createPluginManagerFromRootFolder(configuration)是插件系统的入口,见 PluginUtils.java:当插件目录存在时,用DirectoryBasedPluginFinder扫描目录。

扫描规则(见 DirectoryBasedPluginFinder.java)非常直接:

  • 只把顶层子目录识别为一个插件,子目录名成为插件 ID;
  • 子目录中的 JAR 文件(匹配glob:**.jar)按 URL 字典序排序后,组成该插件的资源 URL 列表;
  • 如果某个插件子目录中一个 JAR 都没有,会抛出IOException提示补齐 JAR 或删除目录——所以不要留下空的插件目录。

每个子目录最终被封装为一个PluginDescriptor(pluginId, urls, excludePatterns)

3. 类加载与 SPI 发现:PluginManager 与 PluginLoader

DefaultPluginManager(见 DefaultPluginManager.java)维护插件 ID 到PluginLoader的映射:同一个插件 ID 的PluginLoader只会创建一次并复用(用ReentrantLock保证线程安全);load(Class<P> service)对每个已知插件调用PluginLoader.load(service),把各插件返回的迭代器拼接成总迭代器。

PluginLoader(见 PluginLoader.java)则把PluginClassLoaderServiceLoader组合起来:在TemporaryClassLoaderContext中临时把线程上下文类加载器切换为插件类加载器,然后通过ServiceLoader.load(service, pluginClassLoader)META-INF/services/<SPI>中声明的实现类名实例化插件服务。文件系统的 SPI 入口是org.apache.flink.core.fs.FileSystemFactory,新文件系统的接入方式(继承FileSystem/ 实现FileSystemFactory/ 添加 service entry)详见文件系统概览中的"添加新的外部文件系统实现"一节。

由于 SPI 接口本身由系统类加载器加载,而实现类由插件类加载器加载,因此保证了FileSystem在全进程中的单例性,也解释了为什么插件实现中不应使用Thread.currentThread().getContextClassLoader()——运行期间上下文类加载器会被插件机制动态切换。

4. Parent-first 类加载配置:白名单的可配置实现

文档中提到的"白名单"在实践中对应一个可配置项。CoreOptions(见 CoreOptions.java)提供了两个配置:

  • plugin.classloader.parent-first-patterns:始终从插件父类加载器解析的类名前缀列表(默认值即日志框架等白名单模式),官方建议一般不要修改;
  • plugin.classloader.parent-first-patterns.additional:追加额外的 parent-first 模式,用于缓解插件机制的意外副作用(该配置项被标注为"实现细节",仅在需要时使用)。

两者通过getPluginParentFirstLoaderPatterns(Configuration)合并后传入PluginLoader,最终与插件自身的 exclude 模式拼接,构造PluginClassLoader(见 PluginLoader.java)。这也是日志框架能在核心、插件、用户代码间统一配置的底层实现。

小结与最佳实践

  • 部署文件系统插件:启动前把opt/下对应 JAR 复制到plugins/<任意目录名>/;S3 的两个插件严禁放入lib;凭证提供者 JAR 必须与插件放在同一目录。
  • 部署指标插件:发行版默认已在plugins/下打包好 JMX、Prometheus、Graphite、InfluxDB、StatsD、Datadog、SLF4J 等上报器,直接按指标上报器文档配置启用即可。
  • 理解隔离原理:每个插件一个PluginClassLoader,SPI 接口走系统类加载器保证单例,服务实现通过ServiceLoaderMETA-INF/services发现,parent-first 模式可通过plugin.classloader.parent-first-patterns.additional扩展。
  • 编写自定义插件:只使用@Public/@PublicEvolving的类与接口,保留META-INF/services服务定义,避免在实现中依赖线程上下文类加载器;插件目录不能为空。

如需进一步了解插件所承载的各类文件系统能力(S3 凭证配置、OSS endpoint、Azure 凭据、GCS 集成等),可继续阅读文件系统概览及其下的 S3、OSS、Azure 等专题文档。

【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询