- 大数据
- 流处理
- 批处理
- 数据工程
【免费下载链接】flink
Flink SQL 的 JAR 语句(ADD JAR、SHOW JARS、REMOVE 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 类:AddJarOperation、ShowJarsOperation与RemoveJarOperation。它们均实现了Operation接口,其中AddJarOperation还实现了ExecutableOperation,说明它携带真正的执行逻辑,而ShowJarsOperation实现了ShowOperation,用于展示查询结果。
运行 JAR 语句
JAR 语句的典型运行环境是 SQL CLI(Flink SQL 客户端)。在 SQL CLI 中输入 JAR 语句并回车即可执行,每条语句执行后会返回执行结果提示。
以下示例展示了在 SQL CLI 中依次执行ADD JAR、SHOW JARS、REMOVE 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); } }执行链路的关键步骤如下:
- 将 JAR 路径包装为
ResourceUri(资源类型为ResourceType.JAR); - 调用
ResourceManager.registerJarResources()完成注册; - 注册成功返回
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 JAR后SHOW 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 JAR、SHOW JARS、REMOVE JAR构成了 Flink SQL 在运行时动态管理用户依赖的完整闭环:
ADD JAR负责将本地或远程 JAR 校验、暂存并注册到会话 classpath,底层由ResourceManager与MutableURLClassLoader支撑;SHOW JARS实时展示当前会话已注册的 JAR 列表;REMOVE JAR从会话 classloader 中移除指定 JAR,但当前仅 SQL CLI 支持。
掌握这三个语句,你可以在不重启集群的前提下灵活扩展 SQL 作业的 UDF、连接器与格式库,同时注意规避 Hive 集成与 SQL Gateway 场景下的已知限制。
- 大数据
- 流处理
- 批处理
- 数据工程
【免费下载链接】flink
相关推荐
Flink SQL JAR 语句实战指南:ADD JAR / SHOW JARS / REMOVE JAR 用法与底层原理
Flink SQL JAR 语句实战指南:ADD JAR / SHOW JARS / REMOVE JAR 用法与底层原理 JAR 语句是 Flink SQL
大数据流处理批处理数据工程Flink Hive 方言 ADD 语句完全指南:用 ADD JAR 向类路径动态加载用户 JAR
Flink Hive 方言 ADD 语句完全指南:用 ADD JAR 向类路径动态加载用户 JAR Hive 方言(Hive Dialect)是 Flink T
大数据流处理批处理数据工程Apache Spark SQL 资源管理语句全解析:ADD/LIST FILE、JAR 与 ARCHIVE 的语法、参数与底层实现
Apache Spark SQL 资源管理语句全解析:ADD/LIST FILE、JAR 与 ARCHIVE 的语法、参数与底层实现 导读 在 Apache S
大数据数据分析批处理流处理机器学习图计算
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考