Hive 함수
Hive 함수 (Hive Functions)
Flink에서 Hive 내장 함수와 사용자 정의 함수(UDF)를 사용하는 방법을 다룹니다.
출처: 문서
본문
HiveModule로 Hive 내장 함수 사용 (Use Hive Built-in Functions via HiveModule)
HiveModule은 Hive 내장 함수를 Flink 시스템(내장) 함수로 Flink SQL과 Table API 사용자에게 제공합니다. 자세한 내용은 HiveModule을 참고하세요.
Java
String name = "myhive";
String version = "2.3.4";
tableEnv.loadModue(name, new HiveModule(version));
Scala
val name = "myhive"
val version = "2.3.4"
tableEnv.loadModue(name, new HiveModule(version));
Python
from pyflink.table.module import HiveModule
name = "myhive"
version = "2.3.4"
t_env.load_module(name, HiveModule(version))
SQL Client
LOAD MODULE hive WITH ('hive-version' = '2.3.4');
일부 이전 버전의 Hive 내장 함수는 스레드 안전성 문제가 있습니다. 사용자 자신의 Hive를 패치하여 수정할 것을 권장합니다.
네이티브 Hive 집계 함수 사용 (Use Native Hive Aggregate Functions)
HiveModule이 CoreModule보다 높은 우선순위로 로드되면 Flink는 먼저 Hive 내장 함수를 사용하려고 합니다. 그리고 Hive 내장 집계 함수의 경우 Flink는 현재 정렬(sort) 기반 집계 연산자만 사용할 수 있습니다. Flink 1.17부터 해시(hash) 기반 집계 연산자로 실행할 수 있는 네이티브 Hive 집계 함수 중 일부를 도입했습니다. 현재는 sum/count/avg/min/max 5개의 함수만 지원하며, 향후 더 많은 집계 함수를 지원할 예정입니다. 사용자는 table.exec.hive.native-agg-function.enabled 옵션을 켜서 네이티브 집계 함수를 사용할 수 있으며, 이는 작업 성능을 크게 향상시킵니다.
| 키 | 기본값 | 타입 | 설명 |
|---|---|---|---|
| table.exec.hive.native-agg-function.enabled | false | Boolean | 네이티브 집계 함수 사용을 활성화합니다. 집계 성능을 향상시킬 수 있는 해시 기반 집계 전략을 사용할 수 있습니다. 작업(job) 레벨 옵션입니다. |
주의: 네이티브 집계 함수의 능력은 현재 Hive 내장 집계 함수와 완전히 일치하지 않습니다. 예를 들어 일부 데이터 타입은 지원되지 않습니다. 성능이 병목이 아니라면 이 옵션을 켤 필요는 없습니다. 또한
table.exec.hive.native-agg-function.enabled옵션은 SqlClient를 통해 사용할 때 작업별로 켤 수 없으며, 현재 모듈 레벨만 지원됩니다. 사용자는 이 옵션을 먼저 켠 후 HiveModule을 로드해야 합니다. 이 문제는 향후 수정될 예정입니다.
Hive 사용자 정의 함수 (Hive User Defined Functions)
사용자는 기존의 Hive 사용자 정의 함수를 Flink에서 사용할 수 있습니다. 지원되는 UDF 타입은 다음과 같습니다:
- UDF
- GenericUDF
- GenericUDTF
- UDAF
- GenericUDAFResolver2
쿼리 계획 및 실행 시 Hive의 UDF와 GenericUDF는 자동으로 Flink의 ScalarFunction으로, Hive의 GenericUDTF는 Flink의 TableFunction으로, Hive의 UDAF와 GenericUDAFResolver2는 Flink의 AggregateFunction으로 변환됩니다.
Hive 사용자 정의 함수를 사용하려면 사용자는:
- 해당 함수를 포함하는 Hive Metastore 기반
HiveCatalog를 세션의 현재 카탈로그로 설정해야 합니다. - 해당 함수를 포함하는 jar를 Flink의 classpath에 포함해야 합니다.
Hive 사용자 정의 함수 사용 (Using Hive User Defined Functions)
Hive Metastore에 다음과 같은 Hive 함수가 등록되어 있다고 가정합니다:
/**
* Test simple udf. Registered under name 'myudf'
*/
public class TestHiveSimpleUDF extends UDF {
public IntWritable evaluate(IntWritable i) {
return new IntWritable(i.get());
}
public Text evaluate(Text text) {
return new Text(text.toString());
}
}
/**
* Test generic udf. Registered under name 'mygenericudf'
*/
public class TestHiveGenericUDF extends GenericUDF {
@Override
public ObjectInspector initialize(ObjectInspector[] arguments) throws UDFArgumentException {
checkArgument(arguments.length == 2);
checkArgument(arguments[1] instanceof ConstantObjectInspector);
Object constant = ((ConstantObjectInspector) arguments[1]).getWritableConstantValue();
checkArgument(constant instanceof IntWritable);
checkArgument(((IntWritable) constant).get() == 1);
if (arguments[0] instanceof IntObjectInspector ||
arguments[0] instanceof StringObjectInspector) {
return arguments[0];
} else {
throw new RuntimeException("Not support argument: " + arguments[0]);
}
}
@Override
public Object evaluate(DeferredObject[] arguments) throws HiveException {
return arguments[0].get();
}
@Override
public String getDisplayString(String[] children) {
return "TestHiveGenericUDF";
}
}
Hive CLI에서 등록된 것을 확인할 수 있습니다:
hive> show functions;
OK
......
mygenericudf
myudf
myudtf
그런 다음 SQL에서 다음과 같이 사용할 수 있습니다:
Flink SQL> select mygenericudf(myudf(name), 1) as a, mygenericudf(myudf(age), 1) as b, s from mysourcetable, lateral table(myudtf(name, 1)) as T(s);