Topik ini menjelaskan cara mengembangkan dan menggunakan user-defined functions (UDF) berbasis Java.
Informasi latar belakang
Mulai versi v2.2.0, StarRocks mendukung user-defined functions (UDF) yang ditulis dalam Java.
Mulai versi v3.0, StarRocks mendukung UDF global. Anda cukup menambahkan kata kunci GLOBAL pada pernyataan SQL terkait (CREATE, SHOW, dan DROP) agar berlaku secara global, sehingga tidak perlu mengeksekusi pernyataan tersebut di setiap database secara terpisah. Dengan demikian, Anda dapat mengembangkan fungsi kustom sesuai skenario bisnis untuk memperluas kemampuan StarRocks.
StarRocks saat ini mendukung jenis UDF berikut:
-
scalar user-defined function (scalar UDF)
-
user-defined aggregate function (UDAF)
-
user-defined window function (UDWF)
-
user-defined table-valued function (UDTF)
Prasyarat
Sebelum menggunakan Java UDF di StarRocks, pastikan Anda memenuhi persyaratan berikut:
-
Anda harus menginstal Apache Maven untuk membuat dan mengembangkan proyek Java.
-
Anda harus menginstal JDK 1.8 di server Anda.
-
Aktifkan fitur UDF dengan mengatur item konfigurasi FE
enable_udfmenjadiTRUEpada halaman Instance Configuration. Kemudian, restart instans untuk menerapkan perubahan tersebut.
Pemetaan tipe data
|
Tipe SQL |
Tipe Java |
|
BOOLEAN |
java.lang.Boolean |
|
TINYINT |
java.lang.Byte |
|
SMALLINT |
java.lang.Short |
|
INT |
java.lang.Integer |
|
BIGINT |
java.lang.Long |
|
FLOAT |
java.lang.Float |
|
DOUBLE |
java.lang.Double |
|
STRING/VARCHAR |
java.lang.String |
Kembangkan dan gunakan UDF
Anda akan membuat proyek Maven dan menulis fungsi dalam Java.
Langkah 1: Buat proyek Maven
Buat proyek Maven dengan struktur direktori dasar berikut.
project
|--pom.xml
|--src
| |--main
| | |--java
| | |--resources
| |--test
|--target
Langkah 2: Tambahkan dependensi
Tambahkan dependensi berikut ke file pom.xml.
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>org.example</groupId>
<artifactId>udf</artifactId>
<version>1.0-SNAPSHOT</version>
<properties>
<maven.compiler.source>8</maven.compiler.source>
<maven.compiler.target>8</maven.compiler.target>
</properties>
<dependencies>
<dependency>
<groupId>com.alibaba</groupId>
<artifactId>fastjson</artifactId>
<version>1.2.76</version>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-dependency-plugin</artifactId>
<version>2.10</version>
<executions>
<execution>
<id>copy-dependencies</id>
<phase>package</phase>
<goals>
<goal>copy-dependencies</goal>
</goals>
<configuration>
<outputDirectory>${project.build.directory}/lib</outputDirectory>
</configuration>
</execution>
</executions>
</plugin>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-assembly-plugin</artifactId>
<version>3.3.0</version>
<executions>
<execution>
<id>make-assembly</id>
<phase>package</phase>
<goals>
<goal>single</goal>
</goals>
</execution>
</executions>
<configuration>
<descriptorRefs>
<descriptorRef>jar-with-dependencies</descriptorRef>
</descriptorRefs>
</configuration>
</plugin>
</plugins>
</build>
</project>
Langkah 3: Kembangkan UDF
Kembangkan UDF yang diperlukan dalam Java.
Kembangkan scalar UDF
Scalar UDF memproses satu baris dalam satu waktu dan mengembalikan satu nilai untuk setiap baris. Saat digunakan dalam kueri, setiap baris muncul dalam set hasil. Contoh fungsi scalar umum meliputi UPPER, LOWER, ROUND, dan ABS.
Contoh berikut menunjukkan cara mengekstraksi data dari objek JSON. Dalam beberapa skenario bisnis, nilai suatu bidang dalam data JSON mungkin berupa string JSON, bukan objek JSON. Untuk mengekstraksi string JSON bersarang tersebut, Anda harus menggunakan pemanggilan bersarang ke fungsi GET_JSON_STRING, seperti GET_JSON_STRING(GET_JSON_STRING('{"key":"{\\"k0\\":\\"v0\\"}"}', "$.key"), "$.k0").
Untuk menyederhanakan pernyataan SQL, Anda dapat mengembangkan UDF untuk mengekstraksi string JSON secara langsung, misalnya MY_UDF_JSON_GET('{"key":"{\\"k0\\":\\"v0\\"}"}', "$.key.k0").
package com.starrocks.udf.sample;
import com.alibaba.fastjson.JSONPath;
public class UDFJsonGet {
public final String evaluate(String jsonObj, String key) {
if (obj == null || key == null) return null;
try {
// Pustaka JSONPath dapat sepenuhnya memperluas path, bahkan jika nilai bidang adalah string berformat JSON.
return JSONPath.read(jsonObj, key).toString();
} catch (Exception e) {
return null;
}
}
}
Kelas kustom Anda harus mengimplementasikan metode berikut.
Tipe data parameter permintaan dan nilai kembali dalam metode tersebut harus sesuai dengan yang dideklarasikan dalam pernyataan CREATE FUNCTION pada Langkah 6. Pemetaan tipe data harus mengikuti aturan dalam Pemetaan tipe data.
|
Metode |
Deskripsi |
|
TYPE1 evaluate(TYPE2, ...) |
Metode |
Kembangkan UDAF
UDAF beroperasi pada sekelompok baris dan mengembalikan satu nilai. Contoh fungsi agregat umum meliputi SUM, COUNT, MAX, dan MIN. Fungsi-fungsi ini mengagregasi data dari beberapa baris dalam setiap kelompok GROUP BY dan menghasilkan satu hasil untuk setiap kelompok.
Contoh berikut menggunakan fungsi MY_SUM_INT. Berbeda dengan fungsi bawaan SUM yang mengembalikan nilai BIGINT, MY_SUM_INT menerima argumen INT dan mengembalikan INT.
package com.starrocks.udf.sample;
public class SumInt {
public static class State {
int counter = 0;
public int serializeLength() { return 4; }
}
public State create() {
return new State();
}
public void destroy(State state) {
}
public final void update(State state, Integer val) {
if (val != null) {
state.counter+= val;
}
}
public void serialize(State state, java.nio.ByteBuffer buff) {
buff.putInt(state.counter);
}
public void merge(State state, java.nio.ByteBuffer buffer) {
int val = buffer.getInt();
state.counter += val;
}
public Integer finalize(State state) {
return state.counter;
}
}
Kelas kustom Anda harus mengimplementasikan metode berikut.
Tipe data parameter permintaan dan nilai kembali dalam metode tersebut harus sesuai dengan yang dideklarasikan dalam pernyataan CREATE FUNCTION pada Langkah 6. Pemetaan tipe data harus mengikuti aturan dalam Pemetaan tipe data.
|
Metode |
Deskripsi |
|
State create() |
Membuat objek state. |
|
void destroy(State) |
Menghapus objek state. |
|
void update(State, ...) |
Memperbarui state. Argumen pertama adalah objek state, dan argumen berikutnya adalah parameter permintaan yang ditentukan dalam deklarasi fungsi. Satu atau lebih parameter permintaan didukung. |
|
void serialize(State, ByteBuffer) |
Menyerialisasi state. |
|
void merge(State, ByteBuffer) |
Menggabungkan state yang telah diserialisasi ke dalam state saat ini. |
|
TYPE finalize(State) |
Menghitung dan mengembalikan hasil akhir dari state. |
Saat mengembangkan UDAF, Anda harus menggunakan kelas buffer java.nio.ByteBuffer untuk menyimpan hasil antara dan metode serializeLength untuk menentukan panjang hasil antara yang telah diserialisasi.
|
Kelas dan metode |
Deskripsi |
|
java.nio.ByteBuffer() |
Kelas buffer ini menyimpan hasil antara. Karena hasil tersebut diserialisasi untuk transfer antar node, Anda harus menggunakan |
|
serializeLength() |
Panjang dalam byte dari hasil antara yang telah diserialisasi. Metode |
Perhatikan persyaratan berikut untuk serialisasi java.nio.ByteBuffer:
-
Jangan mengandalkan metode
remaining()dariByteBufferuntuk mendeserialisasi state. -
Jangan memanggil metode
clear()padaByteBuffer. -
Nilai yang dikembalikan oleh
serializeLength()harus sesuai dengan panjang aktual data yang ditulis ke buffer. Ketidaksesuaian menyebabkan error serialisasi dan deserialisasi.
Kembangkan UDWF
UDWF adalah jenis khusus dari fungsi agregat. Berbeda dengan fungsi agregat biasa, fungsi jendela menghitung nilai pada sekelompok baris (jendela) dan mengembalikan hasil terpisah untuk setiap baris. Biasanya, fungsi jendela mencakup klausa OVER yang membagi baris ke dalam partisi. Fungsi tersebut kemudian menghitung hasil untuk setiap baris berdasarkan jendela tempat baris tersebut berada.
Contoh berikut menggunakan fungsi MY_WINDOW_SUM_INT. Berbeda dengan fungsi bawaan SUM yang mengembalikan nilai BIGINT, MY_WINDOW_SUM_INT menerima argumen INT dan mengembalikan INT.
package com.starrocks.udf.sample;
public class WindowSumInt {
public static class State {
int counter = 0;
public int serializeLength() { return 4; }
@Override
public String toString() {
return "State{" +
"counter=" + counter +
'}';
}
}
public State create() {
return new State();
}
public void destroy(State state) {
}
public void update(State state, Integer val) {
if (val != null) {
state.counter+=val;
}
}
public void serialize(State state, java.nio.ByteBuffer buff) {
buff.putInt(state.counter);
}
public void merge(State state, java.nio.ByteBuffer buffer) {
int val = buffer.getInt();
state.counter += val;
}
public Integer finalize(State state) {
return state.counter;
}
public void reset(State state) {
state.counter = 0;
}
public void windowUpdate(State state,
int peer_group_start, int peer_group_end,
int frame_start, int frame_end,
Integer[] inputs) {
for (int i = (int)frame_start; i < (int)frame_end; ++i) {
state.counter += inputs[i];
}
}
}
Karena fungsi jendela merupakan jenis khusus dari fungsi agregat, kelas kustom Anda harus mengimplementasikan semua metode UDAF yang diperlukan, ditambah metode windowUpdate().
Tipe data parameter permintaan dan nilai kembali dalam metode tersebut harus sesuai dengan yang dideklarasikan dalam pernyataan CREATE FUNCTION pada Langkah 6. Pemetaan tipe data harus mengikuti aturan dalam Pemetaan tipe data.
Metode tambahan yang diperlukan
|
Metode |
Deskripsi |
|
|
Memperbarui data jendela. Untuk informasi lebih lanjut tentang fungsi jendela, lihat Window functions. Untuk setiap baris input, informasi jendela yang sesuai diambil untuk memperbarui hasil antara.
|
Kembangkan UDTF
UDTF menerima satu baris sebagai input dan menghasilkan tabel berisi beberapa baris output. UDTF sering digunakan untuk operasi seperti memisahkan satu kolom menjadi beberapa baris.
Saat ini, UDTF hanya mendukung pengembalian beberapa baris dengan satu kolom.
Contoh berikut menggunakan fungsi MY_UDF_SPLIT. Fungsi ini memisahkan string dengan spasi sebagai pemisah. Baik argumen input maupun nilai kembali bertipe STRING.
package com.starrocks.udf.sample;
public class UDFSplit{
public String[] process(String in) {
if (in == null) return null;
return in.split(" ");
}
}
Kelas kustom Anda harus mengimplementasikan metode berikut.
Tipe data parameter permintaan dan nilai kembali dalam metode tersebut harus sesuai dengan yang dideklarasikan dalam pernyataan CREATE FUNCTION pada Langkah 6. Pemetaan tipe data harus mengikuti aturan dalam Pemetaan tipe data.
|
Metode |
Deskripsi |
|
TYPE[] process() |
Metode |
Langkah 4: Paket proyek Java
Jalankan perintah berikut untuk memaketkan proyek Java.
mvn package
Perintah ini menghasilkan dua file JAR di direktori target: udf-1.0-SNAPSHOT.jar dan udf-1.0-SNAPSHOT-jar-with-dependencies.jar.
Langkah 5: Unggah proyek
Unggah file udf-1.0-SNAPSHOT-jar-with-dependencies.jar ke Bucket OSS dan berikan izin baca-publik pada file JAR tersebut. Untuk informasi lebih lanjut, lihat Simple upload dan Bucket ACL.
Pada Langkah 6, FE memvalidasi file JAR UDF dan menghitung checksum. BE kemudian mengunduh dan mengeksekusi file JAR tersebut.
Langkah 6: Buat UDF di StarRocks
StarRocks menyediakan dua namespace untuk UDF: tingkat database dan tingkat global.
-
Jika Anda tidak memerlukan isolasi visibilitas untuk UDF Anda, Anda dapat membuat UDF global. Saat mereferensikan UDF global, Anda dapat memanggilnya langsung dengan nama fungsinya tanpa awalan katalog atau database, yang menyederhanakan akses.
-
Jika Anda memerlukan isolasi visibilitas atau perlu membuat UDF dengan nama yang sama di database berbeda, Anda dapat membuat UDF tingkat database. Jika sesi Anda saat ini berada dalam database tersebut, Anda dapat memanggil fungsi tersebut langsung dengan namanya. Jika sesi Anda berada di katalog atau database berbeda, Anda harus menggunakan nama lengkap, seperti
catalog.database.function.
Membuat UDF global memerlukan izin CREATE GLOBAL FUNCTION tingkat sistem. Membuat UDF tingkat database memerlukan izin CREATE FUNCTION pada database tersebut. Menggunakan UDF apa pun memerlukan izin USAGE pada UDF tersebut. Untuk informasi tentang cara memberikan izin, lihat GRANT.
Setelah mengunggah file JAR, buat UDF yang sesuai di StarRocks. Untuk membuat UDF global, cukup tambahkan kata kunci GLOBAL ke pernyataan SQL.
Sintaks
CREATE [GLOBAL][AGGREGATE | TABLE] FUNCTION function_name(arg_type [, ...])
RETURNS return_type
[PROPERTIES ("key" = "value" [, ...]) ]
Parameter
|
Parameter |
Wajib |
Deskripsi |
|
GLOBAL |
Tidak |
Menentukan bahwa UDF adalah UDF global. Didukung mulai StarRocks v3.0. |
|
AGGREGATE |
Tidak |
Menentukan bahwa fungsi yang akan dibuat adalah UDAF atau UDWF. |
|
TABLE |
Tidak |
Menentukan bahwa fungsi yang akan dibuat adalah UDTF. |
|
function_name |
Ya |
Nama fungsi. Nama tersebut dapat mencakup nama database, seperti |
|
arg_type |
Ya |
Tipe data argumen fungsi. Untuk tipe data yang didukung, lihat Pemetaan tipe data. |
|
return_type |
Ya |
Tipe data nilai kembali fungsi. Untuk tipe data yang didukung, lihat Pemetaan tipe data. |
|
properties |
Ya |
Properti terkait fungsi. Properti yang berbeda diperlukan untuk jenis UDF yang berbeda. Lihat contoh berikut untuk detailnya. |
Buat scalar UDF
Jalankan perintah berikut untuk membuat scalar UDF dari contoh sebelumnya di StarRocks.
CREATE [GLOBAL] FUNCTION MY_UDF_JSON_GET(string, string)
RETURNS string
PROPERTIES (
"symbol" = "com.starrocks.udf.sample.UDFJsonGet",
"type" = "StarrocksJar",
"file" = "http://<YourBucketName>.oss-cn-xxxx-internal.aliyuncs.com/<YourPath>/udf-1.0-SNAPSHOT-jar-with-dependencies.jar"
);
|
Parameter |
Deskripsi |
|
symbol |
Nama kelas lengkap UDF, dalam format |
|
type |
Jenis UDF. Untuk Java UDF, atur nilai ini menjadi |
|
file |
Path HTTP ke file JAR UDF. Ini harus berupa URL HTTP file di OSS, sebaiknya menggunakan Titik akhir internal. Formatnya adalah |
Buat UDAF
Jalankan perintah berikut untuk membuat UDAF dari contoh sebelumnya di StarRocks.
CREATE [GLOBAL] AGGREGATE FUNCTION MY_SUM_INT(INT)
RETURNS INT
PROPERTIES
(
"symbol" = "com.starrocks.udf.sample.SumInt",
"type" = "StarrocksJar",
"file" = "http://<YourBucketName>.oss-cn-xxxx-internal.aliyuncs.com/<YourPath>/udf-1.0-SNAPSHOT-jar-with-dependencies.jar"
);
Parameter dalam klausa PROPERTIES sama dengan yang dijelaskan dalam Buat scalar UDF.
Buat UDWF
Jalankan perintah berikut untuk membuat UDWF dari contoh sebelumnya di StarRocks.
CREATE [GLOBAL] AGGREGATE FUNCTION MY_WINDOW_SUM_INT(Int)
RETURNS Int
PROPERTIES
(
"analytic" = "true",
"symbol" = "com.starrocks.udf.sample.WindowSumInt",
"type" = "StarrocksJar",
"file" = "http://<YourBucketName>.oss-cn-xxxx-internal.aliyuncs.com/<YourPath>/udf-1.0-SNAPSHOT-jar-with-dependencies.jar"
);
Properti analytic mengidentifikasi fungsi sebagai fungsi jendela. Untuk UDWF, atur nilai ini menjadi true. Parameter lainnya sama dengan yang dijelaskan dalam Buat scalar UDF.
Buat UDTF
Jalankan perintah berikut untuk membuat UDTF dari contoh sebelumnya di StarRocks.
CREATE [GLOBAL] TABLE FUNCTION MY_UDF_SPLIT(string)
RETURNS string
PROPERTIES
(
"symbol" = "com.starrocks.udf.sample.UDFSplit",
"type" = "StarrocksJar",
"file" = "http://<YourBucketName>.oss-cn-xxxx-internal.aliyuncs.com/<YourPath>/udf-1.0-SNAPSHOT-jar-with-dependencies.jar"
);
Parameter dalam klausa PROPERTIES sama dengan yang dijelaskan dalam Buat scalar UDF.
Langkah 7: Gunakan UDF
Setelah membuat UDF, Anda dapat mengujinya dan menggunakannya.
Gunakan scalar UDF
Jalankan perintah berikut untuk menggunakan scalar UDF yang dibuat pada Langkah 6.
SELECT MY_UDF_JSON_GET('{"key":"{\\"in\\":2}"}', '$.key.in');
Gunakan UDAF
Jalankan perintah berikut untuk menggunakan UDAF yang dibuat pada Langkah 6.
SELECT MY_SUM_INT(col1);
Gunakan UDWF
Jalankan perintah berikut untuk menggunakan UDWF yang dibuat pada Langkah 6.
SELECT MY_WINDOW_SUM_INT(intcol)
OVER (PARTITION BY intcol2
ORDER BY intcol3
ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING)
FROM test_basic;
Gunakan UDTF
Jalankan perintah berikut untuk menggunakan UDTF dari contoh sebelumnya.
-- Asumsikan tabel t1 ada dengan kolom a, b, dan c1.
SELECT t1.a,t1.b,t1.c1 FROM t1;
> output:
1,2.1,"hello world"
2,2.2,"hello UDTF."
-- Gunakan fungsi MY_UDF_SPLIT().
SELECT t1.a,t1.b, MY_UDF_SPLIT FROM t1, MY_UDF_SPLIT(t1.c1);
> output:
1,2.1,"hello"
1,2.1,"world"
2,2.2,"hello"
2,2.2,"UDTF."
-
Dalam daftar SELECT,
MY_UDF_SPLITpertama adalah alias kolom default untuk output fungsiMY_UDF_SPLIT. -
Saat ini, Anda tidak dapat menentukan alias tabel atau alias kolom untuk output UDTF, seperti dengan
AS t2(f1).
Lihat UDF
Jalankan perintah berikut untuk melihat informasi tentang UDF.
SHOW [GLOBAL] FUNCTIONS;
Hapus UDF
Jalankan perintah berikut untuk menghapus UDF tertentu.
DROP [GLOBAL] FUNCTION <function_name>(arg_type [, ...]);
FAQ
Q: Apakah saya dapat menggunakan variabel static saat mengembangkan UDF? Apakah variabel static dari UDF yang berbeda saling mengganggu?
A: Ya, Anda dapat menggunakan variabel static. Variabel static diisolasi antar UDF yang berbeda dan tidak saling mengganggu, bahkan jika UDF tersebut memiliki nama kelas yang sama.