Flink SQL JAR 语句完全指南:ADD JAR / SHOW JARS / REMOVE JAR 的用法与原理
2026/9/24 13:38:10 网站建设 项目流程
  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】flink

项目地址:https://gitcode.com/gh_mirrors/fli/flink
点击查看免费下载

Flink SQL 的 JAR 语句(ADD JARSHOW JARSREMOVE JAR)用于在运行时向会话的 classpath 中动态添加、查看或移除用户 JAR(例如自定义 UDF、连接器等),让开发者无需重启集群即可扩展作业能力。本文以 Apache Flink 当前仓库的官方文档为主体,结合 AddJarOperation.java、ResourceManager.java 等源码实现,带你掌握 JAR 语句的完整语法、SQL CLI 实操示例、底层资源注册机制以及使用时的关键限制。

JAR 语句概览

JAR 语句用于在运行时将用户 JAR 添加到 classpath、从 classpath 移除 JAR,或查看当前 classpath 中已添加的 JAR 列表。

Flink SQL 目前支持以下三种 JAR 语句:

语句作用
ADD JAR将用户 JAR 添加到资源列表(即 classpath)
SHOW JARS展示所有通过ADD JAR添加的 JAR
REMOVE JAR移除通过ADD JAR添加的指定 JAR

在源码层面,这三种语句分别对应 flink-table-api-java 模块org.apache.flink.table.operations.command包下的三个 Operation 类:AddJarOperationShowJarsOperationRemoveJarOperation。它们均实现了Operation接口,其中AddJarOperation还实现了ExecutableOperation,说明它携带真正的执行逻辑,而ShowJarsOperation实现了ShowOperation,用于展示查询结果。

运行 JAR 语句

JAR 语句的典型运行环境是 SQL CLI(Flink SQL 客户端)。在 SQL CLI 中输入 JAR 语句并回车即可执行,每条语句执行后会返回执行结果提示。

以下示例展示了在 SQL CLI 中依次执行ADD JARSHOW JARSREMOVE JAR的完整过程:

Flink SQL> ADD JAR '/path/hello.jar'; [INFO] Execute statement succeeded. Flink SQL> ADD JAR 'hdfs:///udf/common-udf.jar'; [INFO] Execute statement succeeded. Flink SQL> SHOW JARS; +----------------------------+ | jars | +----------------------------+ | /path/hello.jar | | hdfs:///udf/common-udf.jar | +----------------------------+ Flink SQL> REMOVE JAR '/path/hello.jar'; [INFO] The specified jar is removed from session classloader.

从示例可以看到:

  • 添加本地路径 JAR 与远程文件系统 JAR 均返回Execute statement succeeded
  • SHOW JARS以表格形式输出所有已添加 JAR 的路径(列名为jars),顺序与添加顺序一致;
  • REMOVE JAR成功后提示该 JAR 已从session classloader(会话类加载器)中移除。

在 SQL CLI 的帮助命令中,这三种语句的定义也能看到其语义描述(见 CliStrings.java):

  • ADD JAR:将指定 JAR 文件添加到提交作业的 classloader,语法为ADD JAR '<path_to_filename>.jar'
  • REMOVE JAR:从提交作业的 classloader 中移除指定 JAR,语法为REMOVE JAR '<path_to_filename>.jar'
  • SHOW JARS:展示用户指定的 JAR 依赖列表,该列表受--jar--library启动选项以及ADD/REMOVE JAR命令的共同影响。

ADD JAR

ADD JAR '<path_to_filename>.jar'

ADD JAR将一个 JAR 文件添加到资源列表中。它支持添加位于本地或远程文件系统的 JAR(关于远程文件系统的支持范围,可参考 文件系统总览)。添加成功后,可以通过SHOW JARS语句查看已添加的 JAR。

语法要点

  • 路径参数必须用单引号包裹;
  • 路径后缀应为.jar
  • 支持file://hdfs://等 scheme 前缀,例如ADD JAR 'hdfs:///udf/common-udf.jar'
  • 远程路径的 JAR 会在注册时被下载到本地临时目录(详见下文源码解析)。

源码实现与执行链路

ADD JAR的执行由AddJarOperation.execute(Context ctx)驱动(见 AddJarOperation.java):

@Override public TableResultInternal execute(Context ctx) { ResourceUri resourceUri = new ResourceUri(ResourceType.JAR, getPath()); try { ctx.getResourceManager().registerJarResources(Collections.singletonList(resourceUri)); return TableResultImpl.TABLE_RESULT_OK; } catch (IOException e) { throw new TableException( String.format("Could not register the specified resource [%s].", resourceUri), e); } }

执行链路的关键步骤如下:

  1. 将 JAR 路径包装为ResourceUri(资源类型为ResourceType.JAR);
  2. 调用ResourceManager.registerJarResources()完成注册;
  3. 注册成功返回TABLE_RESULT_OK,失败则抛出TableException并携带资源信息。

在 ResourceManager.java 中,registerJarResources的实现进一步揭示了两点底层细节:

public void registerJarResources(List<ResourceUri> resourceUris) throws IOException { registerResources( prepareStagingResources( resourceUris, ResourceType.JAR, true, url -> { try { JarUtils.checkJarFile(url); } catch (IOException e) { throw new ValidationException( String.format("Failed to register jar resource [%s]", url), e); } }, false), true); }
  • 事务式注册:注册过程采用"先暂存(staging)再真正注册"的两阶段方式,一旦列表中有任意资源注册失败,整个注册过程会回滚,避免产生半注册状态;
  • 合法性校验:注册前会对 JAR 执行JarUtils.checkJarFile(url)校验,确保其是合法的 JAR 文件,非法文件会抛出ValidationException

此外,ResourceManager的构造逻辑(见 ResourceManager.java)说明了远程 JAR 的处理方式:本地下载目录由配置项table.resources.download-dir决定,默认值为System.getProperty("java.io.tmpdir")(即系统临时目录,定义见 TableConfigOptions.java)。从源码结构可以推断,远程 JAR 会被下载到该目录下以flink-table-<UUID>命名的子目录中,并交由MutableURLClassLoader动态加载到 classpath——这就是"无需重启即可加载新 JAR"的实现基础。

使用限制

请勿使用ADD JAR语句来加载 Hive 的 source/sink/function/catalog。这是 Hive 连接器的一个已知限制(known limitation),将在未来版本中修复。目前,建议按照 Hive 连接器依赖配置指南 中的说明来搭建 Hive 集成。

SHOW JARS

SHOW JARS

展示所有通过ADD JAR语句添加的 JAR。该语句的底层实现是ShowJarsOperation(见 ShowJarsOperation.java):

@Override public TableResultInternal execute(Context ctx) { String[] jars = ctx.getResourceManager().getResources().keySet().stream() .map(ResourceUri::getUri) .toArray(String[]::new); return buildStringArrayResult("jars", jars); }

其内部逻辑为:从ResourceManager.getResources()(返回不可变视图,见 ResourceManager.java)中取出所有已注册资源的 URI,并以列名为jars的字符串数组结果返回。因此SHOW JARS输出中的jars列名,直接来源于源码中的buildStringArrayResult("jars", jars)调用。

REMOVE JAR

REMOVE JAR '<path_to_filename>.jar'

移除通过ADD JAR语句添加的指定 JAR。路径参数同样需要与添加时完全一致的字符串(包含 scheme 前缀),例如添加时写hdfs:///udf/common-udf.jar,移除时也必须写hdfs:///udf/common-udf.jar

在源码层面,RemoveJarOperation(见 RemoveJarOperation.java)仅实现了Operation接口,其asSummaryString()返回REMOVE JAR '<path>'形式的语句摘要,而实际的移除动作由 SQL CLI 的会话类加载器完成,这也与 REMOVE 成功后返回的提示语"The specified jar is removed from session classloader"相互印证。

使用限制

注意REMOVE JAR语句仅在 SQL CLI 中可用

从源码实现看,这一限制在 SQL Gateway 侧有明确印证:OperationExecutor.callRemoveJar方法直接抛出UnsupportedOperationException("SQL Gateway doesn't support REMOVE JAR syntax now.")(见 OperationExecutor.java),即 SQL Gateway 当前不支持REMOVE JAR语法。

常见问题与最佳实践

1. 为什么ADD JARSHOW JARS看不到?

请检查:

  • 是否使用了不同的路径写法(如file:///path/a.jar/path/a.jar会被视为不同资源);
  • 是否在添加时使用了相对路径,导致解析后的绝对路径不一致;
  • 是否在 SQL Gateway 会话中操作(不同会话的 classloader 相互隔离)。

2. 远程 JAR 的加载时机

从 ResourceManager.java 的实现看,ADD JAR执行时即会校验并注册资源,注册成功的 JAR 会立即进入当前会话的 classloader 可见范围。对于远程文件系统(如 HDFS、S3、OSS),请确保目标文件系统已按 文件系统插件指南 正确配置。

3. 需要动态加载 Hive 组件怎么办?

如文档所述,ADD JAR目前不能用于加载 Hive source/sink/function/catalog。正确的做法是按照 Hive 连接器依赖配置指南 中关于依赖(dependencies)的说明,在启动时通过--jar--library等选项或在lib/目录中放置依赖来搭建 Hive 集成。

4. 确认当前会话已加载的 JAR

可以直接在 SQL CLI 中执行SHOW JARS查看当前会话的 JAR 列表。注意该列表不仅包含ADD JAR添加的资源,还包含通过--jar--library启动选项传入的依赖(见 CliStrings.java 对SHOW JARS的描述)。

总结

ADD JARSHOW JARSREMOVE JAR构成了 Flink SQL 在运行时动态管理用户依赖的完整闭环:

  • ADD JAR负责将本地或远程 JAR 校验、暂存并注册到会话 classpath,底层由ResourceManagerMutableURLClassLoader支撑;
  • SHOW JARS实时展示当前会话已注册的 JAR 列表;
  • REMOVE JAR从会话 classloader 中移除指定 JAR,但当前仅 SQL CLI 支持。

掌握这三个语句,你可以在不重启集群的前提下灵活扩展 SQL 作业的 UDF、连接器与格式库,同时注意规避 Hive 集成与 SQL Gateway 场景下的已知限制。

  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】flink

项目地址:https://gitcode.com/gh_mirrors/fli/flink
点击查看免费下载

相关推荐

上一篇:Tartube开发者指南:深入理解Python GUI应用架构
下一篇:9cc项目架构与代码组织:如何构建易于理解的编译器源代码结构

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

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

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

立即咨询