All Products
Search
Document Center

MaxCompute:Java UDAF

Last Updated:Aug 22, 2026

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.Aggregator dan menggunakan anotasi @Resolve (com.aliyun.odps.udf.annotation.Resolve). Kelas com.aliyun.odps.udf.UDFException bersifat 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 mana signature adalah 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.Aggregator dan 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, dan terminate adalah 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;
  }
}
Catatan

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.class tetapi 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_list juga mendukung tanda bintang (*) atau string kosong ('').

    • Jika arg_type_list berupa tanda bintang (*), artinya fungsi menerima sejumlah parameter input apa pun.

    • Jika arg_type_list berupa 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.

Catatan

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

@Resolve('bigint,double->string')

Tipe parameter input adalah BIGINT dan DOUBLE, dan tipe nilai kembali adalah STRING.

@Resolve('*->string')

Menerima sejumlah parameter input apa pun, dan tipe nilai kembali adalah STRING.

@Resolve('->double')

Tidak menerima parameter input, dan tipe nilai kembali adalah DOUBLE.

@Resolve('array<bigint>->struct<x:string, y:int>')

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

Catatan

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.

求平均值逻辑

  1. 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.

  2. 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.

  3. Fase kedua perhitungan rata-rata: Hasil antara dari setiap irisan pada fase pertama diagregasi.

  4. Output akhir: r.sum/r.count adalah rata-rata dari semua data input.

Langkah-langkah berikut menjelaskan cara mengembangkan dan memanggil Java UDAF:

  1. 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:

    1. Instal MaxCompute Studio

    2. Hubungkan ke proyek MaxCompute

    3. Buat modul Java MaxCompute

  2. Tulis kode UDAF

    1. Di explorer Project, klik kanan direktori kode sumber modul (src > main > java) dan pilih New > MaxCompute Java.

    2. 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.

    3. 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;
        }
      }
  3. 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 ke kmeans_in, Table columns ke dim1,dim2, Download Record limit ke 100, dan Data Column Separator ke koma. Lalu, klik OK.

    Catatan

    Anda dapat menggunakan data yang ditunjukkan pada gambar sebagai referensi untuk parameter eksekusi.

  4. 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.

  5. 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_table memiliki 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|
    +----+