拓冰建站拓冰建站
首页 / 资讯中心 / 正文

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

Flink Plugins 插件机制详解文件系统与 Metric Reporter 的隔离加载与实战部署【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flinkApache 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-hadoop和flink-azure-fs-hadoop可以依赖不同的、互相冲突的库版本且无需在打 fat jar 时重定位任何类。这一隔离机制在源码中有清晰的实现DefaultPluginManager内部维护了一张MapString, PluginLoaderpluginId - PluginLoader每个插件 ID 对应一个独立的PluginLoaderPluginLoader通过PluginClassLoader继承自URLClassLoader加载该插件目录下的所有 JAR见 flink-core/src/main/java/org/apache/flink/core/plugin/PluginLoader.java 中的createPluginClassLoader与create方法。白名单与 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 S3flink-s3-fs-presto/flink-s3-fs-hadoop、阿里云 OSSflink-oss-fs-hadoop、Azure Blob Storageflink-azure-fs-hadoop、Google Cloud Storagegcs-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-presto和flink-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.jarMetric Reporter 插件Flink 提供的所有 Metric Reporter 都可以作为插件使用指标上报器的完整配置方式参见指标上报器文档。发行版构建时Flink 已经通过 flink-dist/src/main/assemblies/plugins.xml 把开箱即用的指标插件打包进plugins/目录包括插件目录对应 JAR源码模块plugins/metrics-jmx/flink-metrics-jmxflink-metrics-jmxplugins/metrics-graphite/flink-metrics-graphiteflink-metrics-graphiteplugins/metrics-influx/flink-metrics-influxdbflink-metrics-influxdbplugins/metrics-prometheus/flink-metrics-prometheusflink-metrics-prometheusplugins/metrics-statsd/flink-metrics-statsdflink-metrics-statsdplugins/metrics-datadog/flink-metrics-datadogflink-metrics-datadogplugins/metrics-slf4j/flink-metrics-slf4jflink-metrics-slf4j此外plugins.xml还展示了插件机制的一个扩展用例——GPU 外部资源插件plugins/external-resource-gpu/其中除了 JAR 还打包了gpu-discovery-common.sh与nvidia-gpu-discovery.sh两个发现脚本见 flink-external-resources/flink-external-resource-gpu说明插件目录中并非只能放 JAR也可以放置插件运行所需的资源文件。源码级剖析插件发现与加载的完整链路1. 插件目录如何确定环境变量与默认值插件根目录的解析逻辑在 PluginConfig.java优先读取环境变量FLINK_PLUGINS_DIR未设置时回退到默认值plugins相对于 Flink 发行版根目录。对应的常量定义见 ConfigConstants.javaENV_FLINK_PLUGINS_DIR/DEFAULT_FLINK_PLUGINS_DIRS。如果该目录不存在PluginConfig只记录一条警告日志并返回空插件系统以无插件状态继续启动。2. 目录扫描DirectoryBasedPluginFinderPluginUtils.createPluginManagerFromRootFolder(configuration)是插件系统的入口见 PluginUtils.java当插件目录存在时用DirectoryBasedPluginFinder扫描目录。扫描规则见 DirectoryBasedPluginFinder.java非常直接只把顶层子目录识别为一个插件子目录名成为插件 ID子目录中的 JAR 文件匹配glob:**.jar按 URL 字典序排序后组成该插件的资源 URL 列表如果某个插件子目录中一个 JAR 都没有会抛出IOException提示补齐 JAR 或删除目录——所以不要留下空的插件目录。每个子目录最终被封装为一个PluginDescriptor(pluginId, urls, excludePatterns)。3. 类加载与 SPI 发现PluginManager 与 PluginLoaderDefaultPluginManager见 DefaultPluginManager.java维护插件 ID 到PluginLoader的映射同一个插件 ID 的PluginLoader只会创建一次并复用用ReentrantLock保证线程安全load(ClassP service)对每个已知插件调用PluginLoader.load(service)把各插件返回的迭代器拼接成总迭代器。PluginLoader见 PluginLoader.java则把PluginClassLoader与ServiceLoader组合起来在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 等上报器直接按指标上报器文档配置启用即可。理解隔离原理每个插件一个PluginClassLoaderSPI 接口走系统类加载器保证单例服务实现通过ServiceLoader按META-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),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门