的语法、执行方式与原理)
大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载CALL语句是 Flink Table API SQL 中用于调用存储过程Procedure的专用 SQL 语句通常被用来执行数据操纵或管理类任务例如清理/重写数据文件、回滚快照等。本文以 docs/content/docs/dev/table/sql/call.md 为核心脉络结合 Flink 源码完整讲解CALL语句的语法、在 Java/Scala/Python/SQL CLI 四种环境下的执行方式、底层调用链路与结果处理机制以及存储过程的实现与注册方式。读完本文你将能够在自己的 Flink 项目中编写、注册并调用存储过程并理解其类型推导与参数绑定的内部原理。注意CALL语句要求被调用的过程procedure必须已经存在于对应的 catalog 中。如果过程不存在执行时会抛出异常。关于某个 catalog 提供哪些过程需要参考该 catalog 的文档关于如何实现一个过程请参考 Procedures 实现指南。一、CALL 语句概述CALLCall Statement语句用于调用一个存储过程该过程通常由第三方提供用于执行数据操纵data manipulation或管理administrative任务。它与常规的SELECT查询不同SELECT描述一个数据查询/变换逻辑惰性执行CALL直接调用过程立即执行其中的逻辑并返回一个TableResult关联该过程产生的结果集。从源码结构看CALL语句在 planner 层会被解析为一个CallProcedureOperationplanner 侧的包装实现是 PlannerCallProcedureOperation其execute(Context ctx)方法完成了参数转换 → 反射调用call方法 → 结果转 TableResult的完整链路。过程Procedure是什么一个存储过程在 Flink 中表现为实现了接口org.apache.flink.table.procedures.Procedure的类。该接口本身不声明任何方法见 Procedure.java你需要自行定义一个名为call的公共方法来实现过程逻辑。Procedure接口的 Javadoc 给出了两个典型场景的伪代码示例IcebergRewriteDataFilesProcedure接受STRING表名返回ROWrewritten_data_files_count STRING, added_data_files_count STRING数组用于重写数据文件RollbackToSnapShotProcedure接受STRING表名和LONG快照 ID返回String[]用于回滚快照。这些正是CALL语句面向的数据操纵与管理任务的典型形态。底层机制的完整说明见 Procedures 实现指南。二、CALL 语句语法CALL语句的完整语法如下CALL [catalog_name.][database_name.]procedure_name ([ expression [, expression]* ] )语法要点catalog_name可选过程所属 catalog省略时使用当前 catalogdatabase_name可选过程所属 database省略时使用当前 databaseprocedure_name过程名称expression一个或多个逗号分隔过程参数可以是字面量、表达式如1 2、带cast的类型转换表达式如cast(1 as bigint)、时间戳/区间表达式等整个调用需要加一对括号参数可以为空即CALL proc()。与函数/表函数名称一样当 catalog 名、database 名或过程名包含特殊字符或需要精确指定时可以使用反引号包裹例如CALL system.generate_n(4)表示调用当前 catalog下systemdatabase 中的generate_n过程CALL my_catalog.system.generate_n(5)则表示显式指定 catalog 为my_catalog。在 SqlNodeToCallOperationTest 中可以看到各种语法的解析验证callsystem.primitive_arg(1, 2)→ 解析为CALL PROCEDURE: (procedureIdentifier: [p1.system.primitive_arg], inputTypes: [INT NOT NULL, BIGINT NOT NULL], outputTypes: [INT NOT NULL], arguments: [1, 2])表达式参数callsystem.var_arg(1, 2, 1 2)→ 参数值被规约为3参数必须是可规约为字面量的表达式否则抛异常例如cast((1.2 2.4) as decimal)这种不可规约表达式会报 cant be converted to literal命名参数写法call system.named_args(d1,c2)其参数顺序按名称绑定为arguments: [2, 1]签名不匹配时如callsystem.primitive_arg(1)抛出 No match found for function signature primitive_arg( )。三、如何执行一条 CALL 语句CALL语句可通过四种方式执行Java、Scala、Python 的TableEnvironmentAPI以及 SQL CLI。无论哪种方式executeSql/execute_sql都会立即调用过程并返回一个TableResult实例。3.1 Java使用TableEnvironment.executeSql()方法TableEnvironment tEnv TableEnvironment.create(...); // assuming the procedure generate_n has existed in system database of the current catalog tEnv.executeSql(CALL system.generate_n(4)).print();输出效果与 SQL CLI 类似打印出 0~3 四行结果。3.2 Scala使用TableEnvironment.executeSql()方法val tEnv TableEnvironment.create(...) // assuming the procedure generate_n has existed in system database of the current catalog tEnv.executeSql(CALL system.generate_n(4)).print()3.3 PythonPython 端使用execute_sql()方法TableEnvironment的 Python 版本table_env TableEnvironment.create(...) # assuming the procedure generate_n has existed in system database of the current catalog table_env.execute_sql(CALL system.generate_n(4)).print()3.4 SQL CLI在 SQL CLI 中直接输入CALL语句// assuming the procedure generate_n has existed in system database of the current catalog Flink SQL CALL system.generate_n(4); -------- | result | -------- | 0 | | 1 | | 2 | | 3 | -------- 4 rows in set !ok可以看到当过程的返回类型是原子类型非复合类型时结果会被包装为单列单行的表列名统一为result。四、存储过程的实现与注册CALL 的前置条件CALL语句能成功执行的前提是被调用的过程已由某个 catalog 提供。本节按 Procedures 实现指南 的步骤从实现到注册给出可复制的完整流程。4.1 实现 Procedure 类实现类必须实现接口org.apache.flink.table.procedures.Procedure类必须是public、非abstract、全局可访问不允许非静态内部类或匿名类自行定义一个名为call的public方法实现过程逻辑。call方法有两个硬性约定第一个参数必须是ProcedureContext它通过getExecutionEnvironment()提供StreamExecutionEnvironment用于在过程内部运行一个 Flink Job见 ProcedureContext.java。Flink 在每次过程调用时都会基于当前配置新建一个StreamExecutionEnvironment传入过程内部可以安全地修改它而不会泄漏到外部返回类型必须是数组如int[]、String[]、Row[]等。常规 JVM 方法调用语义同样适用因此call方法支持方法重载如call(ProcedureContext, Integer)与call(ProcedureContext, LocalDateTime)可变参数如call(ProcedureContext, Integer...)对象继承如call(ProcedureContext, Object)可同时接受LocalDateTime和Integer以上组合如call(ProcedureContext, Object...)接受任意类型的参数。若用 Scala 实现过程并打算使用可变参数需要添加scala.annotation.varargs注解此外推荐使用装箱类型如java.lang.Integer而非Int以支持NULL。一个带重载call方法的示例import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.table.procedure.ProcedureContext; import org.apache.flink.table.procedures.Procedure; // procedure with overloaded call methods public class GenerateSequenceProcedure implements Procedure { public long[] call(ProcedureContext context, int n) { return generate(context.getExecutionEnvironment(), n); } public long[] call(ProcedureContext context, String n) { return generate(context.getExecutionEnvironment(), Integer.parseInt(n)); } private long[] generate(StreamExecutionEnvironment env, int n) throws Exception { long[] sequenceN new long[n]; int i 0; try (CloseableIteratorLong result env.fromSequence(0, n - 1).executeAndCollect()) { while (result.hasNext()) { sequenceN[i] result.next(); } } return sequenceN; } }4.2 类型推导Type InferenceFlink 的表生态是强类型系统与 SQL 标准一致过程的参数与返回值都必须映射为 数据类型。planner 需要逻辑层面的类型、精度、标度信息也需要 JVM 层面内部数据结构如何表示为 JVM 对象的信息。验证输入参数、为参数与结果推导数据类型的过程统称为类型推导type inference。默认情况下Flink 通过反射自动从过程类和call方法推导数据类型。如果反射提取不成功可以用DataTypeHint与ProcedureHint注解辅助。注意虽然call方法返回类型必须是数组T[]但使用DataTypeHint注解返回类型时注解的其实是数组的元素类型T。DataTypeHint行内类型提示用于在参数或返回类型上直接声明数据类型import org.apache.flink.table.annotation.DataTypeHint; import org.apache.flink.table.annotation.InputGroup; import org.apache.flink.table.procedure.ProcedureContext; import org.apache.flink.table.procedures.Procedure; import org.apache.flink.types.Row; public static class OverloadedProcedure implements Procedure { // no hint required public Long[] call(ProcedureContext context, long a, long b) { return new Long[] {a b}; } // define the precision and scale of a decimal public DataTypeHint(DECIMAL(12, 3)) BigDecimal[] call(ProcedureContext context, double a, double b) { return new BigDecimal[] {BigDecimal.valueOf(a b)}; } // define a nested data type DataTypeHint(ROWs STRING, t TIMESTAMP_LTZ(3)) public Row[] call(ProcedureContext context, int i) { return new Row[] {Row.of(String.valueOf(i), Instant.ofEpochSecond(i))}; } // allow wildcard input and custom serialized output DataTypeHint(value RAW, bridgedTo ByteBuffer.class) public ByteBuffer[] call(ProcedureContext context, DataTypeHint(inputGroup InputGroup.ANY) Object o) { return new ByteBuffer[] {MyUtils.serializeToByteBuffer(o)}; } }ProcedureHint过程级/方法级类型提示当希望一个call方法同时处理多种数据类型或多个重载call方法共享同一个结果类型时使用ProcedureHint。它可以在类上或每个call方法上声明一个或多个注解为输入与结果数据类型建立映射。所有 hint 参数都是可选的未定义的部分使用默认的反射提取定义在类上的 hint 会被所有call方法继承。import org.apache.flink.table.annotation.DataTypeHint; import org.apache.flink.table.annotation.ProcedureHint; import org.apache.flink.table.procedure.ProcedureContext; import org.apache.flink.table.procedures.Procedure; import org.apache.flink.types.Row; // procedure with overloaded call methods // but globally defined output type ProcedureHint(output DataTypeHint(ROWs STRING, i INT)) public static class OverloadedProcedure implements Procedure { public Row[] call(ProcedureContext context, int a, int b) { return new Row[] {Row.of(Sum, a b)}; } // overloading of arguments is still possible public Row[] call(ProcedureContext context) { return new Row[] {Row.of(Empty args, -1)}; } } // decouples the type inference from call methods, // the type inference is entirely determined by the procedure hints ProcedureHint( input {DataTypeHint(INT), DataTypeHint(INT)}, output DataTypeHint(INT) ) ProcedureHint( input {DataTypeHint(BIGINT), DataTypeHint(BIGINT)}, output DataTypeHint(BIGINT) ) ProcedureHint( input {}, output DataTypeHint(BOOLEAN) ) public static class OverloadedProcedure implements Procedure { // an implementer just needs to make sure that a method exists // that can be called by the JVM public Object[] call(ProcedureContext context, Object... o) { if (o.length 0) { return new Object[] {false}; } return new Object[] {o[0]}; } }上述第一个示例中两个重载call方法共享了类级声明的输出类型ROWs STRING, i INT第二个示例则完全用ProcedureHint声明了三组输入输出映射call方法本身只需能被 JVM 调用即可类型推导完全由 hint 决定。4.3 命名参数Named Parameters调用过程时可以使用参数名来指定参数值这样既能避免参数顺序写错导致的混乱也能省略非必填参数省略项默认填充null提升代码可读性与可维护性。使用ArgumentHint注解可以指定参数的名字、类型以及是否必填。ArgumentHint可以在三个作用域上使用① 在call方法的参数上使用public static class NamedParameterProcedure implements Procedure { public DataTypeHint(INT) Integer[] call(ProcedureContext context, ArgumentHint(name a, isOption true) Integer a, ArgumentHint(name b) Integer b) { return new Integer[] {a (b null ? 0 : b)}; } }② 在call方法上使用public static class NamedParameterProcedure extends Procedure { ProcedureHint( argument {ArgumentHint(name param1, type DataTypeHint(INTEGER), isOptional false), ArgumentHint(name param2, type DataTypeHint(INTEGER), isOptional true)} ) public DataTypeHint(INT) Integer[] call(ProcedureContext context, Integer a, Integer b) { return new Integer[] {a (b null ? 0 : b)}; } }③ 在过程类上使用ProcedureHint( argument {ArgumentHint(name param1, type DataTypeHint(INTEGER), isOptional false), ArgumentHint(name param2, type DataTypeHint(INTEGER), isOptional true)} ) public static class NamedParameterProcedure implements Procedure { public DataTypeHint(INT) Integer[] call(ProcedureContext context, Integer a, Integer b) { return new Integer[] {a (b null ? 0 : b)}; } }使用命名参数时需注意ArgumentHint注解本身已包含DataTypeHint语义因此不能与DataTypeHint在ProcedureHint中同时使用应用到函数参数上时也不能与DataTypeHint并用推荐统一使用ArgumentHint命名参数仅当过程类不包含重载函数和可变参数函数时生效否则使用命名参数会报错。对应 SQL 侧命名参数的调用形如call system.named_args(d1,c2)在 SqlNodeToCallOperationTest 中有完整测试用例其过程用ProcedureHint(argumentNames {c, d})声明参数名。4.4 在 Catalog 中返回过程实现过程后catalog 需要重写Catalog.getProcedure(ObjectPath procedurePath)方法返回过程同时建议在Catalog.listProcedures(String dbName)中列出所有过程。下面的示例基于GenericInMemoryCatalog提供一个内置了system.generate_n过程的 catalogimport org.apache.flink.table.catalog.Catalog; import org.apache.flink.table.catalog.GenericInMemoryCatalog; import org.apache.flink.table.catalog.ObjectPath; import org.apache.flink.table.catalog.exceptions.CatalogException; import org.apache.flink.table.catalog.exceptions.DatabaseNotExistException; import org.apache.flink.table.catalog.exceptions.ProcedureNotExistException; import org.apache.flink.table.procedure.ProcedureContext; import org.apache.flink.table.procedures.Procedure; import java.util.HashMap; import java.util.Map; // catalog with built-in procedures public static class CatalogWithBuiltInProcedure extends GenericInMemoryCatalog { static { PROCEDURE_MAP.put(ObjectPath.fromString(system.generate_n), new GenerateSequenceProcedure()); } public CatalogWithBuiltInProcedure(String name) { super(name); } Override public ListString listProcedures(String dbName) throws DatabaseNotExistException, CatalogException { if (!databaseExists(dbName)) { throw new DatabaseNotExistException(getName(), dbName); } return PROCEDURE_MAP.keySet().stream().filter(procedurePath - procedurePath.getDatabaseName().equals(dbName)) .map(ObjectPath::getObjectName).collect(Collectors.toList()); } Override public Procedure getProcedure(ObjectPath procedurePath) throws ProcedureNotExistException, CatalogException { if (PROCEDURE_MAP.containsKey(procedurePath)) { return PROCEDURE_MAP.get(procedurePath); } else { throw new ProcedureNotExistException(getName(), procedurePath); } } }示例依赖前文 4.1 中实现的GenerateSequenceProcedurePROCEDURE_MAP为静态的MapObjectPath, Procedure。五、完整示例实现 → 注册 → CALL 调用将上面的步骤串起来在 Java 中完成实现过程 → 放入自定义 catalog → 注册 catalog → 用 CALL 语句调用的端到端流程import org.apache.flink.table.catalog.GenericInMemoryCatalog; import org.apache.flink.table.catalog.ObjectPath; import org.apache.flink.table.catalog.exceptions.CatalogException; import org.apache.flink.table.catalog.exceptions.ProcedureNotExistException; import org.apache.flink.table.procedure.ProcedureContext; import org.apache.flink.table.procedures.Procedure; // first implement a procedure public class GenerateSequenceProcedure implements Procedure { public long[] call(ProcedureContext context, int n) { long[] sequenceN new long[n]; int i 0; try (CloseableIteratorLong result env.fromSequence(0, n - 1).executeAndCollect()) { while (result.hasNext()) { sequenceN[i] result.next(); } } return sequenceN; } } // then provide the procedure in a custom catalog public static class CatalogWithBuiltInProcedure extends GenericInMemoryCatalog { static { PROCEDURE_MAP.put(ObjectPath.fromString(system.generate_n), new GenerateSequenceProcedure()); } // emit some methods // ... Override public Procedure getProcedure(ObjectPath procedurePath) throws ProcedureNotExistException, CatalogException { if (PROCEDURE_MAP.containsKey(procedurePath)) { return PROCEDURE_MAP.get(procedurePath); } else { throw new ProcedureNotExistException(getName(), procedurePath); } } } TableEnvironment tEnv TableEnvironment.create(...); // register the catalog tEnv.registerCatalog(my_catalog, new CatalogWithBuiltInProcedure()); // call the procedure with CALL statement tEnv.executeSql(call my_catalog.system.generate_n(5));注意GenerateSequenceProcedure.call中通过context.getExecutionEnvironment()拿到StreamExecutionEnvironment用env.fromSequence(...).executeAndCollect()生成并收集序列——这正是ProcedureContext的核心用途让过程可以在内部运行一个真实的 Flink Job。六、底层原理CALL 语句的执行链路从源码看一条CALL语句的执行分为解析 → 构建 Operation → 执行三个阶段核心实现在 PlannerCallProcedureOperation.java① 参数组装getConvertedArgumentValues把 planner 解析得到的内部参数internalInputArguments与类型inputTypes组装成调用参数数组数组首位固定放入ProcedureContext实例。该实例由getProcedureContext基于当前TableConfig的根配置新建Configuration再通过StreamExecutionEnvironment.getExecutionEnvironment(configuration)创建包装为DefaultProcedureContext传入。若参数类型不是内部类型还会用DataStructureConverters.getConverter(inputType)把 Flink 内部值转换为外部值toExternal。② 反射调用call方法callProcedure/invokeCallMethod通过ExtractionUtils.collectMethods收集过程类中所有名为call常量ProcedureDefinition.PROCEDURE_CALL的方法再过滤出可被当前签名调用isInvokable、返回类型是数组、且数组元素类型与输出类型转换类兼容的方法。若找不到匹配方法抛出ValidationException提示 Could not find an implementation method call ... matches the following signature若匹配到多个仅调用第一个并打印 WARN 日志。对于可变参数方法invokeCallMethod会把参数按 varargs 索引重组为数组再调用见源码中的 varargs 调整逻辑。③ 结果转换procedureResultToTableResult若输出类型不是复合类型先用DataTypes.ROW(DataTypes.FIELD(result, ...))包装成单列result表这就是 3.4 节 CLI 输出列名为result的原因若输出是结构化类型STRUCTURED_TYPE则用RowRowConverter完成 RowData 与 Row 的互转最后构建一个CallProcedureResultProvider内部将结果数组逐元素转换为 RowData/Row并支持把基本类型数组int[]包装为Integer[]配合ResultKind.SUCCESS_WITH_CONTENT组装成TableResultImpl返回。asSummaryString方法输出形如CALL PROCEDURE: (procedureIdentifier: [...], inputTypes: [...], outputTypes: [...], arguments: [...])的摘要与测试中的断言一致。与之配套SqlNodeToCallOperationTest 用PlannerCallProcedureOperation实例断言了各种 CALL 场景基本类型参数、不同签名映射、可变参数、Row 结果、POJO 结果、时间戳参数、命名参数、异常场景是理解 CALL 语句行为最直接的测试样本。七、常见错误与排查过程不存在CALL要求过程已存在于对应 catalog否则抛出ProcedureNotExistException或类似的找不到过程的异常。请检查 catalog 的getProcedure(ObjectPath)是否已注册该过程以及 SQL 中 catalog/database 路径是否正确如CALL my_catalog.system.generate_n(5)。签名不匹配SQL 参数的类型与任意call方法都不匹配时抛出ValidationExceptionNo match found for function signature ...。可尝试用cast显式转换参数类型例如cast(1 as bigint)。参数不可规约为字面量CALL 的参数必须是可规约到字面量的表达式否则报 The argument at position N ... cant be converted to literal。命名参数与重载/可变参数冲突过程类存在重载方法或可变参数方法时使用命名参数会报错。call方法不满足约定第一个参数不是ProcedureContext、方法不是public、返回类型不是数组都会导致调用阶段找不到匹配方法。参考文档CALL 语句原始文档Procedures 实现指南SQL CLI数据类型核心实现PlannerCallProcedureOperation.java核心接口Procedure.java、ProcedureContext.java测试用例SqlNodeToCallOperationTest.java赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Flink SQL CALL 语句完全指南存储过程调用、语法与实现原理Flink SQL CALL 语句完全指南存储过程调用、语法与实现原理 CALL 语句是 Flink Table API 与 SQL 中用于调用存储过程St大数据流处理批处理数据工程Presto CALL 语句完全指南语法、命名参数与连接器存储过程调用实战Presto CALL 语句完全指南语法、命名参数与连接器存储过程调用实战 导读 CALL 是 Presto分布式 SQL 查询引擎中用于调用 连接器C大数据数据库后端Flink 存储过程Procedure开发实战实现、类型推导与 CALL 调用全指南Flink 存储过程Procedure开发实战实现、类型推导与 CALL 调用全指南 Flink 在 Table API 与 SQL 中引入了存储过程S大数据流处理批处理数据工程上一篇OpenAI Plugins机器学习指南如何轻松集成AI模型到你的工作流下一篇解决魔兽争霸III在Windows 11上的兼容性问题从根源优化的完整方案创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考