All Products
Search
Document Center

E-MapReduce:Java UDF

Last Updated:Aug 23, 2026

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_udf menjadi TRUE pada 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.

Catatan

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 evaluate adalah titik masuk untuk UDF dan harus merupakan metode anggota publik.

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.

Catatan

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() untuk menentukan ukuran serialisasinya.

serializeLength()

Panjang dalam byte dari hasil antara yang telah diserialisasi. Metode serializeLength() harus mengembalikan INT. Misalnya, State { int counter = 0; public int serializeLength() { return 4; }} menunjukkan bahwa hasil antara adalah INT dan panjang serialisasinya adalah 4 byte. Anda dapat menyesuaikannya sesuai kebutuhan. Misalnya, jika hasil antara adalah LONG dengan panjang 8 byte, Anda dapat menggunakan State { long counter = 0; public int serializeLength() { return 8; }}.

Catatan

Perhatikan persyaratan berikut untuk serialisasi java.nio.ByteBuffer:

  • Jangan mengandalkan metode remaining() dari ByteBuffer untuk mendeserialisasi state.

  • Jangan memanggil metode clear() pada ByteBuffer.

  • 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().

Catatan

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

void windowUpdate(State state, int, int, int , int, ...)

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.

  • peer_group_start: Posisi awal partisi saat ini.

    Partisi adalah sekelompok baris yang memiliki nilai sama untuk kolom yang ditentukan dalam klausa PARTITION BY.

  • peer_group_end: Posisi akhir partisi saat ini.

  • frame_start: Posisi awal frame jendela saat ini.

    Frame jendela adalah subset baris dalam partisi yang menentukan cakupan komputasi untuk baris saat ini. Misalnya, ROWS BETWEEN 1 PRECEDING AND 1 FOLLOWING mendefinisikan frame yang mencakup baris saat ini, baris sebelumnya, dan baris setelahnya.

  • frame_end: Posisi akhir frame jendela saat ini.

  • inputs: Data masukan untuk jendela, disediakan sebagai array kelas wrapper. Kelas wrapper harus sesuai dengan tipe data masukan. Dalam contoh ini, tipe data masukan adalah INT, sehingga array kelas wrapper adalah Integer[].

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.

Catatan

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.

Catatan

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 process() adalah titik masuk untuk UDTF dan harus mengembalikan array.

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.

Catatan

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.

Catatan

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 db1.my_func. Jika function_name mencakup nama database, UDF dibuat di database tersebut. Jika tidak, UDF dibuat di database saat ini. Kombinasi nama fungsi dan argumen harus unik dalam database. Namun, Anda dapat melakukan overload fungsi dengan membuat fungsi lain dengan nama yang sama tetapi tipe argumen berbeda.

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 <package_name>.<class_name>.

type

Jenis UDF. Untuk Java UDF, atur nilai ini menjadi StarrocksJar.

file

Path HTTP ke file JAR UDF. Ini harus berupa URL HTTP file di OSS, sebaiknya menggunakan Titik akhir internal. Formatnya adalah http://<YourBucketName>.oss-cn-xxxx-internal.aliyuncs.com/<YourPath>/<jar_package_name>.

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."
Catatan
  • Dalam daftar SELECT, MY_UDF_SPLIT pertama adalah alias kolom default untuk output fungsi MY_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.