프로시저
프로시저 (Procedures)
Flink Table API & SQL은 사용자가 프로시저(procedure)로 데이터 조작과 관리 작업을 수행할 수 있게 해줘요. 프로시저는 제공된 StreamExecutionEnvironment로 FLINK 작업을 실행할 수 있어 더 강력하고 유연해요.
출처: 문서
본문
구현 가이드 (Implementation Guide)
프로시저를 호출하려면 카탈로그(catalog)에 프로시저가 존재해야 해요. 카탈로그에 프로시저를 제공하려면 프로시저를 구현한 다음 Catalog.getProcedure(ObjectPath procedurePath) 메서드로 반환해야 해요.
다음 단계들은 카탈로그에서 프로시저를 구현하고 제공하는 방법을 안내해요.
프로시저 클래스 (Procedure Class)
구현 클래스는 인터페이스 org.apache.flink.table.procedures.Procedure를 구현해야 해요.
클래스는 public으로 선언되어야 하고, abstract가 아니어야 하며, 전역적으로 접근 가능해야 해요. 따라서 비정적 내부 클래스(non-static inner)나 익명 클래스는 허용되지 않아요.
호출 메서드 (Call Methods)
인터페이스는 어떤 메서드도 제공하지 않으므로, 프로시저의 로직을 구현할 call이라는 이름의 메서드를 정의해야 해요.
메서드는 public으로 선언되어야 하며 잘 정의된 인자 집합을 가져야 해요.
다음 사항에 주의하세요:
call메서드의 첫 번째 파라미터는 항상ProcedureContext여야 해요. 이것은 Flink Job을 실행하기 위한StreamExecutionEnvironment를 얻는 메서드getExecutionEnvironment를 제공해요.- 반환 타입은 항상 배열이어야 해요. 예:
int[],String[]등.
더 자세한 내용은 클래스 org.apache.flink.table.procedures.Procedure의 Java doc에서 찾을 수 있어요.
일반적인 JVM 메서드 호출 의미론이 적용돼요. 따라서 다음이 가능해요:
call(ProcedureContext, Integer)과call(ProcedureContext, LocalDateTime)같은 오버로드된 메서드 구현call(ProcedureContext, Integer...)같은 var-args 사용LocalDateTime과Integer를 모두 받는call(ProcedureContext, Object)같은 객체 상속 사용- 모든 종류의 인자를 받는
call(ProcedureContext, Object...)같은 위의 조합
Scala로 프로시저를 구현하려면 가변 인자(varargs)의 경우 scala.annotation.varargs 어노테이션을 추가하세요. 또한 NULL을 지원하려면 박스형(박스형) 원시 타입(예: Int 대신 java.lang.Integer)을 사용하는 것이 권장돼요.
다음 스니펫은 오버로드된 프로시저의 예시를 보여줘요:
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 (CloseableIterator<Long> result = env.fromSequence(0, n - 1).executeAndCollect()) {
while (result.hasNext()) {
sequenceN[i++] = result.next();
}
}
return sequenceN;
}
}
import org.apache.flink.table.procedure.ProcedureContext
import org.apache.flink.table.procedures.Procedure
import scala.annotation.varargs
// procedures with overloaded call methods
class GenerateSequenceProcedure extends Procedure {
def call(context: ProcedureContext, a: Integer, b: Integer): Array[Integer] = {
Array(a + b)
}
def call(context: ProcedureContext, a: String, b: String): Array[Integer] = {
Array(Integer.valueOf(a) + Integer.valueOf(b))
}
@varargs // generate var-args like Java
def call(context: ProcedureContext, d: Double*): Array[Integer] = {
Array(d.sum.toInt)
}
}
타입 추론 (Type Inference)
테이블 생태계(SQL 표준과 유사)는 강한 타입의 API예요. 따라서 프로시저 파라미터와 반환 타입 모두 데이터 타입에 매핑되어야 해요.
논리적 관점에서 플래너는 기대되는 타입, 정밀도(precision), 스케일(scale)에 대한 정보가 필요해요. JVM 관점에서 플래너는 프로시저를 호출할 때 내부 데이터 구조가 JVM 객체로 어떻게 표현되는지에 대한 정보가 필요해요.
입력 인자를 검증하고 프로시저의 파라미터와 결과에 대한 데이터 타입을 파생하는 로직은 type inference라는 용어로 요약돼요.
Flink의 프로시저는 리플렉션(reflection)을 통해 프로시저의 클래스와 call 메서드에서 데이터 타입을 파생하는 자동 타입 추론 추출을 구현해요. 이 암시적 리플렉션 추출 접근이 성공하지 못하면, @DataTypeHint와 @ProcedureHint로 영향받는 파라미터, 클래스 또는 메서드에 주석을 달아 추출 과정을 지원할 수 있어요. 프로시저에 주석을 다는 더 많은 예시는 아래에 나와 있어요.
참고: call 메서드에서 반환 타입은 배열 타입 T[]여야 하지만, @DataTypeHint로 반환 타입에 주석을 달면 실제로는 배열 타입의 컴포넌트 타입인 T에 주석을 다는 것이 기대돼요.
자동 타입 추론 (Automatic Type Inference)
자동 타입 추론은 프로시저의 클래스와 call 메서드를 검사해 프로시저의 인자와 결과에 대한 데이터 타입을 파생해요. @DataTypeHint와 @ProcedureHint 주석은 자동 추출을 지원해요.
데이터 타입에 암시적으로 매핑될 수 있는 클래스의 전체 목록은 데이터 타입 추출 섹션을 참조하세요.
@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;
// procedure with overloaded call methods
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("ROW")
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)};
}
}
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
// procedure with overloaded call methods
class OverloadedProcedure extends Procedure {
// no hint required
def call(context: ProcedureContext, a: Long, b: Long): Array[Long] = {
Array(a + b)
}
// define the precision and scale of a decimal
@DataTypeHint("DECIMAL(12, 3)")
def call(context: ProcedureContext, a: Double, b: Double): Array[BigDecimal] = {
Array(BigDecimal.valueOf(a + b))
}
// define a nested data type
@DataTypeHint("ROW")
def call(context: ProcedureContext, i: Integer): Array[Row] = {
Row.of(java.lang.String.valueOf(i), java.time.Instant.ofEpochSecond(i))
}
// allow wildcard input and custom serialized output
@DataTypeHint(value = "RAW", bridgedTo = classOf[java.nio.ByteBuffer])
def call(context: ProcedureContext, @DataTypeHint(inputGroup = InputGroup.ANY) o: Object): Array[java.nio.ByteBuffer] = {
Array[MyUtils.serializeToByteBuffer(o)]
}
}
@ProcedureHint
어떤 시나리오에서는 하나의 call 메서드가 동시에 여러 서로 다른 데이터 타입을 처리하는 것이 바람직해요. 또한 오버로드된 call 메서드가 한 번만 선언되어야 하는 공통 결과 타입을 갖는 시나리오도 있어요.
@ProcedureHint 주석은 인자 데이터 타입에서 결과 데이터 타입으로의 매핑을 제공할 수 있어요. 이것은 입력과 결과 데이터 타입에 대해 전체 프로시저 클래스나 call 메서드에 주석을 다는 것을 가능하게 해요. 클래스 상단이나 프로시저 시그니처를 오버로딩하기 위해 각 call 메서드에 하나 이상의 주석을 개별적으로 선언할 수 있어요. 모든 힌트 파라미터는 선택적이에요. 파라미터가 정의되지 않으면 기본 리플렉션 기반 추출이 사용돼요. 프로시저 클래스 상단에 정의된 힌트 파라미터는 모든 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("ROW"))
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]};
}
}
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
import scala.annotation.varargs
// procedure with overloaded call methods
// but globally defined output type
@ProcedureHint(output = new DataTypeHint("ROW"))
class OverloadedFunction extends Procedure {
def call(context: ProcedureContext, a: Int, b: Int): Array[Row] = {
Array(Row.of("Sum", Int.box(a + b)))
}
// overloading of arguments is still possible
def call(context: ProcedureContext): Array[Row] = {
Array(Row.of("Empty args", Int.box(-1)))
}
}
// decouples the type inference from call methods,
// the type inference is entirely determined by the function hints
@ProcedureHint(
input = Array(new DataTypeHint("INT"), new DataTypeHint("INT")),
output = new DataTypeHint("INT")
)
@ProcedureHint(
input = Array(new DataTypeHint("BIGINT"), new DataTypeHint("BIGINT")),
output = new DataTypeHint("BIGINT")
)
@ProcedureHint(
input = Array(),
output = new DataTypeHint("BOOLEAN")
)
class OverloadedProcedure extends Procedure {
// an implementer just needs to make sure that a method exists
// that can be called by the JVM
@varargs
def call(context: ProcedureContext, o: AnyRef*): Array[AnyRef]= {
if (o.length == 0) {
Array(Boolean.box(false))
}
Array(o(0))
}
}
명명된 파라미터 (Named Parameters)
프로시저를 호출할 때 파라미터 이름을 사용해 파라미터 값을 지정할 수 있어요. 명명된 파라미터(named parameters)를 사용하면 파라미터 이름과 값을 동시에 프로시저에 전달할 수 있어, 잘못된 파라미터 순서로 인한 혼란을 피하고 코드의 가독성과 유지 관리성을 향상시켜요. 또한 명명된 파라미터는 필수가 아닌 파라미터를 생략할 수 있으며, 생략된 파라미터는 기본적으로 null로 채워져요. @ArgumentHint 주석을 사용해 파라미터의 이름, 타입, 필수 여부를 지정할 수 있어요.
다음 세 가지 예시는 서로 다른 스코프에서 @ArgumentHint를 사용하는 방법을 보여줘요. 더 많은 정보는 주석 클래스의 문서를 참조하세요.
- 프로시저의
call메서드 파라미터에@ArgumentHint주석 사용.
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;
public static class NamedParameterProcedure extends 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)};
}
}
class NamedParameterProcedure extends Procedure {
def call(context: ProcedureContext, @ArgumentHint(name = "param1", isOptional = true) a: Integer, @ArgumentHint(name = "param2") b: Integer): Array[Integer] = {
Array(a + (if (b == null) 0 else b))
}
}
- 프로시저의
call메서드에@ArgumentHint주석 사용.
public static class NamedParameterProcedure extends Procedure {
@ProcedureHint(
arguments = {
@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)};
}
}
- 프로시저의 클래스에
@ArgumentHint주석 사용.
@ProcedureHint(
arguments = {
@ArgumentHint(name = "param1", type = @DataTypeHint("INTEGER"), isOptional = false),
@ArgumentHint(name = "param2", type = @DataTypeHint("INTEGER"), isOptional = true)
}
)
public static class NamedParameterProcedure extends Procedure {
public @DataTypeHint("INT") Integer[] call(ProcedureContext context, Integer a, Integer b) {
return new Integer[] {a + (b == null ? 0 : b)};
}
}
@ArgumentHint주석은 이미@DataTypeHint주석을 포함하므로,@ProcedureHint안에서@DataTypeHint와 함께 사용할 수 없어요. 함수 파라미터에 적용할 때,@ArgumentHint는@DataTypeHint와 동시에 사용할 수 없으며,@ArgumentHint를 사용하는 것이 권장돼요.- 명명된 파라미터는 해당 프로시저 클래스에 오버로드 함수와 가변 파라미터 함수가 포함되지 않을 때만 적용되며, 그렇지 않으면 명명된 파라미터를 사용하면 오류가 발생해요.
카탈로그에서 프로시저 반환 (Return Procedure in Catalog)
프로시저를 구현한 후 카탈로그는 메서드 Catalog.getProcedure(ObjectPath procedurePath)에서 프로시저를 반환할 수 있어요. 다음 예시는 카탈로그에서 프로시저를 반환하는 방법을 보여줘요.
또한 메서드 Catalog.listProcedures(String dbName)에서 모든 프로시저를 나열하는 것이 기대돼요.
import 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 List<String> 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);
}
}
}
예시 (Examples)
다음 예시는 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 (CloseableIterator<Long> 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)");
// first implement a procedure
class GenerateSequenceProcedure extends Procedure {
def call(context: ProcedureContext, n: Integer): Array[Long] = {
val env = context.getExecutionEnvironment
val sequenceN = Array[Long]
var i = 0;
env.fromSequence(0, n - 1).executeAndCollect()
.forEachRemaining(r => {
sequenceN(i) = r
i = i + 1
})
sequenceN;
}
}
// then provide the procedure in a custom catalog
class CatalogWithBuiltInProcedure(name: String) extends GenericInMemoryCatalog(name) {
val PROCEDURE_MAP = collection.immutable.HashMap[ObjectPath, Procedure](ObjectPath.fromString("system.generate_n"),
new GenerateSequenceProcedure());
@throws(classOf[ProcedureNotExistException])
override def getProcedure(procedurePath: ObjectPath): Procedure = {
if (PROCEDURE_MAP.contains(procedurePath)) {
PROCEDURE_MAP(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)")