Topik ini menjelaskan cara menulis fungsi agregat yang didefinisikan pengguna (UDAF) dalam Java.
Struktur kode UDAF
Anda dapat menulis Java UDAF di IntelliJ IDEA dengan Maven atau MaxCompute Studio. Kode tersebut harus mencakup komponen-komponen berikut:
-
Java package: Opsional.
Anda dapat mengemas kelas Java yang Anda definisikan agar lebih mudah ditemukan dan digunakan kembali.
-
Kelas dan anotasi yang diperlukan: Wajib.
Anda harus mengimpor kelas
com.aliyun.odps.udf.Aggregatordan menggunakan anotasi@Resolve(com.aliyun.odps.udf.annotation.Resolve). Kelascom.aliyun.odps.udf.UDFExceptionbersifat opsional dan dapat digunakan untuk penanganan error. Jika Anda perlu menggunakan kelas terkait UDAF lainnya atau tipe data kompleks, impor kelas yang diperlukan sebagaimana dijelaskan dalam Ikhtisar MaxCompute UDF. -
Anotasi
@Resolve: Wajib.Formatnya adalah
@Resolve(<signature>), di manasignatureadalah signature fungsi yang mendefinisikan tipe data parameter input dan nilai kembali. Signature fungsi UDAF tidak dapat ditentukan melalui refleksi dan hanya dapat diperoleh dengan menggunakan anotasi@Resolve, seperti@Resolve("smallint->varchar(10)"). Untuk informasi selengkapnya tentang anotasi@Resolve, lihat Anotasi @Resolve. -
Kelas Java kustom: Wajib.
Kelas ini merupakan unit organisasi dari kode UDAF. Kelas ini mendefinisikan variabel dan metode yang mengimplementasikan logika bisnis Anda.
-
Metode untuk kelas Java: Wajib.
Kelas Java Anda harus memperluas kelas
com.aliyun.odps.udf.Aggregatordan mengimplementasikan metode-metode berikut.import com.aliyun.odps.udf.ContextFunction; import com.aliyun.odps.udf.ExecutionContext; import com.aliyun.odps.udf.UDFException; public abstract class Aggregator implements ContextFunction { // Metode inisialisasi. @Override public void setup(ExecutionContext ctx) throws UDFException { } // Metode terminasi. @Override public void close() throws UDFException { } // Membuat buffer agregasi. abstract public Writable newBuffer(); // Metode iterate. // buffer adalah buffer agregasi yang menyimpan data antara yang telah dirangkum. Dalam tugas map, buffer ini mengagregasi data untuk satu kelompok, dan metode ini dieksekusi sekali untuk setiap baris. // Writable[] merepresentasikan satu baris data, yang merujuk pada kolom input dalam kode. Misalnya, writable[0] merujuk pada kolom pertama, dan writable[1] merujuk pada kolom kedua. // args adalah parameter yang ditentukan saat memanggil UDAF dalam SQL. Array args itu sendiri tidak boleh null, tetapi elemennya bisa null, yang menunjukkan bahwa data input yang sesuai bernilai null. abstract public void iterate(Writable buffer, Writable[] args) throws UDFException; // Metode terminate. abstract public Writable terminate(Writable buffer) throws UDFException; // Metode merge. abstract public void merge(Writable buffer, Writable partial) throws UDFException; }Metode
iterate,merge, danterminateadalah tiga metode inti yang mengimplementasikan logika utama UDAF. Anda juga harus mengimplementasikan buffer writable kustom.Buffer writable mengonversi objek dalam memori menjadi urutan byte (atau protokol transfer data lainnya) untuk memfasilitasi persistensi ke disk dan transmisi jaringan. Karena MaxCompute menggunakan komputasi terdistribusi untuk memproses fungsi agregat, data harus diserialisasi dan dideserialisasi agar dapat ditransfer antar worker.
Saat menulis Java UDAF, Anda dapat menggunakan tipe Java atau tipe Java Writable. Untuk informasi selengkapnya tentang pemetaan antara tipe data yang didukung MaxCompute dan tipe data Java, lihat Tipe data.
Kode berikut memberikan contoh UDAF.
// Mengemas kelas Java yang didefinisikan dalam org.alidata.odps.udaf.examples.
package org.alidata.odps.udaf.examples;
// Mengimpor kelas dasar yang diperlukan.
import java.io.DataInput;
import java.io.DataOutput;
import java.io.IOException;
import com.aliyun.odps.io.DoubleWritable;
import com.aliyun.odps.io.Writable;
import com.aliyun.odps.udf.Aggregator;
import com.aliyun.odps.udf.UDFException;
import com.aliyun.odps.udf.annotation.Resolve;
// Mendefinisikan kelas Java kustom.
// Menentukan anotasi @Resolve.
@Resolve("double->double")
public class AggrAvg extends Aggregator {
// Mengimplementasikan metode untuk kelas Java.
private static class AvgBuffer implements Writable {
private double sum = 0;
private long count = 0;
@Override
public void write(DataOutput out) throws IOException {
out.writeDouble(sum);
out.writeLong(count);
}
@Override
public void readFields(DataInput in) throws IOException {
sum = in.readDouble();
count = in.readLong();
}
}
private DoubleWritable ret = new DoubleWritable();
@Override
public Writable newBuffer() {
return new AvgBuffer();
}
@Override
public void iterate(Writable buffer, Writable[] args) throws UDFException {
DoubleWritable arg = (DoubleWritable) args[0];
AvgBuffer buf = (AvgBuffer) buffer;
if (arg != null) {
buf.count += 1;
buf.sum += arg.get();
}
}
@Override
public Writable terminate(Writable buffer) throws UDFException {
AvgBuffer buf = (AvgBuffer) buffer;
if (buf.count == 0) {
ret.set(0);
} else {
ret.set(buf.sum / buf.count);
}
return ret;
}
@Override
public void merge(Writable buffer, Writable partial) throws UDFException {
AvgBuffer buf = (AvgBuffer) buffer;
AvgBuffer p = (AvgBuffer) partial;
buf.sum += p.sum;
buf.count += p.count;
}
}
Dalam kode UDAF di atas, buffer pada metode iterate dan merge dapat digunakan kembali. Buffer ini mengagregasi baris input ke dalam buffer sesuai dengan implementasi Anda.
Batasan
Mengakses Internet menggunakan UDF
Secara default, MaxCompute tidak mengizinkan Anda mengakses Internet menggunakan UDF. Jika Anda ingin mengakses Internet menggunakan UDF, isi formulir permohonan koneksi jaringan berdasarkan kebutuhan bisnis Anda dan kirimkan permohonan tersebut. Tim dukungan teknis MaxCompute akan segera menghubungi Anda untuk mengaktifkan konektivitas jaringan. Untuk informasi selengkapnya tentang cara mengisi formulir permohonan koneksi jaringan, lihat Proses koneksi jaringan.
Mengakses VPC menggunakan UDF
Secara default, MaxCompute tidak mengizinkan Anda mengakses resource di VPC menggunakan UDF. Untuk menggunakan UDF guna mengakses resource di VPC, Anda harus membuat koneksi jaringan antara MaxCompute dan VPC tersebut. Untuk informasi selengkapnya tentang operasi terkait, lihat Mengakses resource VPC dari UDF.
Membaca data tabel menggunakan UDF, UDAF, atau UDTF
Anda tidak dapat menggunakan UDF, UDAF, atau UDTF untuk membaca data dari jenis tabel berikut:
Tabel yang telah menjalani evolusi skema
Tabel yang berisi tipe data kompleks
Tabel yang berisi tipe data JSON
Tabel transaksional
Catatan penggunaan
Saat menulis Java UDAF, perhatikan hal-hal berikut:
-
Menyertakan kelas dengan nama yang sama tetapi logika berbeda dalam file JAR UDAF yang berbeda dapat menyebabkan hasil yang tidak terduga atau kegagalan kompilasi. Misalnya, asumsikan UDAF1 dan UDAF2 masing-masing berkorespondensi dengan file resource JAR udaf1.jar dan udaf2.jar. Jika kedua file JAR tersebut berisi kelas bernama
com.aliyun.UserFunction.classtetapi dengan implementasi berbeda, MaxCompute akan memuat salah satu kelas tersebut secara tidak terduga ketika UDAF1 dan UDAF2 dipanggil dalam pernyataan SQL yang sama. -
Dalam Java UDAF, parameter input dan nilai kembali harus berupa tipe objek, seperti String dan Long, bukan tipe primitif.
-
Nilai NULL dalam SQL direpresentasikan sebagai NULL Java. Tipe primitif Java tidak dapat merepresentasikan nilai NULL dalam SQL dan tidak diperbolehkan.
Anotasi @Resolve
Format anotasi @Resolve adalah sebagai berikut.
@Resolve(<signature>)
signature adalah string yang mengidentifikasi tipe data parameter input dan nilai kembali. Saat UDAF dieksekusi, tipe parameter input dan nilai kembalinya harus sesuai dengan tipe yang ditentukan dalam signature fungsi. Selama penguraian semantik, sistem akan memeriksa penggunaan yang tidak sesuai dengan signature fungsi dan melaporkan error jika terjadi ketidakcocokan tipe. Format spesifiknya adalah sebagai berikut.
'arg_type_list -> type'
Deskripsi:
-
arg_type_list: Merepresentasikan tipe data parameter input. Beberapa parameter input dapat ditentukan, dipisahkan oleh koma (,). Tipe data yang didukung adalah BIGINT, STRING, DOUBLE, BOOLEAN, DATETIME, DECIMAL, FLOAT, BINARY, DATE, DECIMAL(presisi,skala), CHAR, VARCHAR, tipe data kompleks (ARRAY, MAP, STRUCT), dan tipe data kompleks bersarang.arg_type_listjuga mendukung tanda bintang (*) atau string kosong ('').-
Jika
arg_type_listberupa tanda bintang (*), artinya fungsi menerima sejumlah parameter input apa pun. -
Jika
arg_type_listberupa string kosong (''), artinya fungsi tidak memiliki parameter input.
Untuk informasi selengkapnya tentang sintaksis lanjutan anotasi Resolve, lihat Parameter dinamis untuk UDAF dan UDTF.
-
-
type: Merepresentasikan tipe data nilai kembali. UDAF hanya mengembalikan satu kolom. Tipe data yang didukung mencakup BIGINT, STRING, DOUBLE, BOOLEAN, DATETIME, DECIMAL, FLOAT, BINARY, DATE, DECIMAL(presisi,skala), tipe data kompleks (ARRAY, MAP, STRUCT), dan tipe data kompleks bersarang.
Saat menulis kode UDAF, Anda dapat memilih tipe data yang sesuai berdasarkan edisi tipe data proyek MaxCompute Anda. Untuk informasi selengkapnya tentang edisi tipe data dan tipe data yang didukung masing-masing edisi, lihat Versi tipe data.
Berikut adalah contoh anotasi @Resolve yang valid.
|
Contoh @Resolve |
Deskripsi |
|
|
Tipe parameter input adalah BIGINT dan DOUBLE, dan tipe nilai kembali adalah STRING. |
|
|
Menerima sejumlah parameter input apa pun, dan tipe nilai kembali adalah STRING. |
|
|
Tidak menerima parameter input, dan tipe nilai kembali adalah DOUBLE. |
|
|
Tipe parameter input adalah ARRAY<BIGINT>, dan tipe nilai kembali adalah STRUCT<x:STRING, y:INT>. |
Tipe data
Tipe data yang didukung MaxCompute bervariasi tergantung edisi tipe datanya. Mulai dari MaxCompute 2.0, tersedia tipe data tambahan, termasuk tipe kompleks seperti ARRAY, MAP, dan STRUCT. Untuk informasi selengkapnya, lihat Edisi tipe data.
Java UDAF Anda harus menggunakan tipe data yang dipetakan ke tipe data MaxCompute. Tabel berikut menjelaskan pemetaan tersebut.
Tipe MaxCompute | Tipe Java | Tipe Java Writable |
TINYINT | java.lang.Byte | ByteWritable |
SMALLINT | java.lang.Short | ShortWritable |
INT | java.lang.Integer | IntWritable |
BIGINT | java.lang.Long | LongWritable |
FLOAT | java.lang.Float | FloatWritable |
DOUBLE | java.lang.Double | DoubleWritable |
DECIMAL | java.math.BigDecimal | BigDecimalWritable |
BOOLEAN | java.lang.Boolean | BooleanWritable |
STRING | java.lang.String | Text |
VARCHAR | com.aliyun.odps.data.Varchar | VarcharWritable |
BINARY | com.aliyun.odps.data.Binary | BytesWritable |
DATE | java.sql.Date | DateWritable |
DATETIME | java.util.Date | DatetimeWritable |
TIMESTAMP | java.sql.Timestamp | TimestampWritable |
INTERVAL_YEAR_MONTH | N/A | IntervalYearMonthWritable |
INTERVAL_DAY_TIME | N/A | IntervalDayTimeWritable |
ARRAY | java.util.List | N/A |
MAP | java.util.Map | N/A |
STRUCT | com.aliyun.odps.data.Struct | N/A |
Parameter input atau nilai kembali UDAF hanya dapat menggunakan Tipe Writable Java jika proyek MaxCompute Anda menggunakan edisi tipe data MaxCompute V2.0.
Penggunaan
Setelah Anda mengembangkan Java UDAF mengikuti proses pengembangan, Anda dapat memanggilnya dalam SQL MaxCompute sebagai berikut:
Gunakan UDF dalam proyek MaxCompute: Caranya mirip dengan penggunaan fungsi bawaan. Anda dapat menggunakan fungsi yang didefinisikan pengguna dengan cara yang sama seperti menggunakan fungsi bawaan.
Gunakan UDF lintas proyek: Gunakan UDF dari Proyek B di Proyek A. Pernyataan berikut menunjukkan contohnya:
select B:udf_in_other_project(arg0, arg1) as res from table_t;. Untuk informasi selengkapnya tentang berbagi lintas proyek, lihat Akses resource lintas proyek berbasis paket.
Untuk contoh lengkap cara mengembangkan dan memanggil Java UDAF menggunakan MaxCompute Studio, lihat Contoh.
Contoh
Contoh ini menunjukkan cara menggunakan MaxCompute Studio untuk mengembangkan UDAF bernama AggrAvg yang menghitung nilai rata-rata. Gambar berikut mengilustrasikan logikanya.

-
Pemotongan data input: MaxCompute mengikuti alur kerja pemrosesan MapReduce untuk memotong data input menjadi irisan yang dapat diproses secara efisien oleh worker.
Anda harus mengonfigurasi ukuran irisan menggunakan parameter
odps.stage.mapper.split.size. Untuk informasi selengkapnya tentang logika pemotongan, lihat Alur kerja pemrosesan MapReduce. -
Fase pertama perhitungan rata-rata: Setiap worker menghitung jumlah catatan data dan menghitung jumlahnya dalam irisan tersebut. Jumlah dan total dari setiap irisan dianggap sebagai hasil antara.
-
Fase kedua perhitungan rata-rata: Hasil antara dari setiap irisan pada fase pertama diagregasi.
-
Output akhir:
r.sum/r.countadalah rata-rata dari semua data input.
Langkah-langkah berikut menjelaskan cara mengembangkan dan memanggil Java UDAF:
Persiapkan lingkungan.
Sebelum Anda dapat mengembangkan dan men-debug UDF di MaxCompute Studio, instal MaxCompute Studio dan hubungkan ke proyek MaxCompute. Untuk informasi selengkapnya, lihat topik-topik berikut:
-
Tulis kode UDAF
Di explorer Project, klik kanan direktori kode sumber modul () dan pilih .
-
Di kotak dialog Create new MaxCompute java class, klik UDAF, masukkan nama di bidang Name, lalu tekan Enter. Untuk contoh ini, beri nama kelas Java tersebut AggrAvg.
Name adalah nama kelas Java MaxCompute. Jika Anda belum membuat paket, masukkan nama dalam format packagename.classname. Paket akan dibuat secara otomatis.
-
Di editor kode, tempel kode UDAF berikut.
import java.io.DataInput; import java.io.DataOutput; import java.io.IOException; import com.aliyun.odps.io.DoubleWritable; import com.aliyun.odps.io.Writable; import com.aliyun.odps.udf.Aggregator; import com.aliyun.odps.udf.UDFException; import com.aliyun.odps.udf.annotation.Resolve; @Resolve("double->double") public class AggrAvg extends Aggregator { private static class AvgBuffer implements Writable { private double sum = 0; private long count = 0; @Override public void write(DataOutput out) throws IOException { out.writeDouble(sum); out.writeLong(count); } @Override public void readFields(DataInput in) throws IOException { sum = in.readDouble(); count = in.readLong(); } } private DoubleWritable ret = new DoubleWritable(); @Override public Writable newBuffer() { return new AvgBuffer(); } @Override public void iterate(Writable buffer, Writable[] args) throws UDFException { DoubleWritable arg = (DoubleWritable) args[0]; AvgBuffer buf = (AvgBuffer) buffer; if (arg != null) { buf.count += 1; buf.sum += arg.get(); } } @Override public Writable terminate(Writable buffer) throws UDFException { AvgBuffer buf = (AvgBuffer) buffer; if (buf.count == 0) { ret.set(0); } else { ret.set(buf.sum / buf.count); } return ret; } @Override public void merge(Writable buffer, Writable partial) throws UDFException { AvgBuffer buf = (AvgBuffer) buffer; AvgBuffer p = (AvgBuffer) partial; buf.sum += p.sum; buf.count += p.count; } }
-
Debug UDAF secara lokal
Untuk operasi debugging selengkapnya, lihat Debug UDF dengan menjalankannya secara lokal.
Klik kanan file AggrAvg.java di proyek Anda dan pilih Run 'AggrAvg.main()'. Di kotak dialog Run/Debug Configurations yang muncul, konfigurasikan parameter berikut: MaxCompute project ke
local, MaxCompute table kekmeans_in, Table columns kedim1,dim2, Download Record limit ke100, dan Data Column Separator ke koma. Lalu, klik OK.CatatanAnda dapat menggunakan data yang ditunjukkan pada gambar sebagai referensi untuk parameter eksekusi.
-
Kemas UDAF yang telah Anda buat ke dalam file JAR, unggah file JAR tersebut ke proyek MaxCompute, dan daftarkan fungsinya. Misalnya, fungsi tersebut diberi nama
user_udaf.Untuk informasi selengkapnya tentang operasi pengemasan, lihat Langkah-langkah.
Di IntelliJ IDEA, klik kanan file Java yang berisi UDAF dan pilih Deploy to server.... Di kotak dialog Package a jar, submit resource and register function, konfigurasikan MaxCompute project, Resource name, Main class (masukkan nama kelas UDAF), dan Function name. Centang kotak Force update if already exists lalu klik OK untuk menyelesaikan penerapan.
-
Di bilah navigasi kiri MaxCompute Studio, klik Project Explorer, klik kanan proyek MaxCompute target, mulai klien MaxCompute, dan eksekusi perintah SQL untuk memanggil UDAF yang baru dibuat.
Asumsikan tabel target
my_tablememiliki skema dan data berikut.+------------+------------+ | col0 | col1 | +------------+------------+ | 1.2 | 2.0 | | 1.6 | 2.1 | +------------+------------+Jalankan pernyataan SQL berikut untuk memanggil UDAF.
select user_udaf(col0) as c0 from my_table;Hasil berikut dikembalikan.
+----+ | c0 | +----+ | 1.4| +----+