All Products
Search
Document Center

Realtime Compute for Apache Flink:Konektor MaxCompute

Last Updated:Jun 04, 2026

Konektor MaxCompute memungkinkan Anda membaca dari dan menulis ke MaxCompute (sebelumnya ODPS)—platform gudang data skala eksabita yang sepenuhnya dikelola oleh Alibaba Cloud—secara langsung dari pekerjaan Flink SQL dan DataStream.

Kemampuan

Item Deskripsi
Tipe tabel Tabel sumber, tabel dimensi, tabel sink, dan sink ingesti data
Mode operasi Mode streaming dan mode batch
Tipe API DataStream API, SQL API, dan pekerjaan YAML ingesti data
Semantik At-least-once
Pembaruan atau penghapusan data di tabel sink Batch Tunnel atau Streaming Tunnel: hanya insert. Upsert Tunnel: insert, update, dan delete.

Metrik

Tipe tabel Metrik
Sumber numRecordsIn, numRecordsInPerSecond, numBytesIn, numBytesInPerSecond
Sink numRecordsOut, numRecordsOutPerSecond, numBytesOut, numBytesOutPerSecond
Tabel dimensi dim.odps.cacheSize
Untuk informasi selengkapnya, lihat Metrik pemantauan.

Prasyarat

Sebelum memulai, pastikan Anda telah:

Batasan

  • Konektor hanya mendukung semantik at-least-once. Catatan duplikat mungkin muncul di MaxCompute tergantung pada saluran data (tunnel) yang digunakan. Untuk panduan pemilihan tunnel, lihat bagian "Bagaimana cara memilih saluran data?" dalam FAQ tentang penyimpanan hulu dan hilir.

  • Secara default, sumber beroperasi dalam mode penuh (full mode): hanya membaca dari partisi yang ditentukan oleh opsi partition. Setelah semua data dibaca, pekerjaan selesai dan tidak memantau partisi baru. Untuk memantau partisi baru secara berkelanjutan, konfigurasikan sumber inkremental menggunakan startPartition.

  • Setiap kali cache tabel dimensi diperbarui, tabel tersebut memeriksa partisi terbaru. Setelah sumber dimulai, data yang baru ditambahkan ke partisi yang sedang dibaca tidak akan dibaca—jalankan penerapan hanya setelah partisi berisi data lengkap.

Pilih tunnel

MaxCompute menyediakan tiga tunnel untuk menulis data dari Flink. Pilih berdasarkan kasus penggunaan Anda:

Tunnel Bawaan Kapan digunakan
MaxCompute Batch Tunnel Ya (useStreamTunnel=false, enableUpsert=false) Muatan batch; data tersedia hanya setelah checkpointing. Atur flushIntervalMs=0 untuk menonaktifkan flushing terjadwal.
MaxCompute Streaming Tunnel Tidak (useStreamTunnel=true) Ingesti hampir real-time; data yang diflush langsung tersedia di MaxCompute.
MaxCompute Upsert Tunnel Tidak (enableUpsert=true) Operasi INSERT, UPDATE, dan DELETE pada tabel Delta MaxCompute. Memerlukan VVR 8.0.6+.

Untuk perbandingan detail, lihat bagian "Bagaimana cara memilih saluran data?" dalam FAQ tentang penyimpanan hulu dan hilir.

SQL

Konektor MaxCompute dapat digunakan sebagai tabel sumber, dimensi, atau sink dalam pekerjaan berbasis SQL.

Sintaksis

CREATE TEMPORARY TABLE odps_source(
  id INT,
  user_name VARCHAR,
  content VARCHAR
) WITH (
  'connector' = 'odps',
  'endpoint' = '<yourEndpoint>',
  'project' = '<yourProjectName>',
  'schemaName' = '<yourSchemaName>',
  'tableName' = '<yourTableName>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}',
  'partition' = 'ds=2018****'
);

Opsi konektor

Opsi umum

Opsi Wajib Bawaan Tipe Deskripsi
connector Ya STRING Atur ke odps.
endpoint Ya STRING Titik akhir MaxCompute. Lihat Titik akhir.
tunnelEndpoint Tidak STRING Titik akhir Tunnel MaxCompute. Jika tidak ditentukan, MaxCompute mengalokasikan koneksi tunnel melalui Server Load Balancer (SLB).
project Ya STRING Nama proyek MaxCompute.
schemaName Tidak STRING Hanya diperlukan jika fitur skema MaxCompute diaktifkan. Atur ke nama skema tabel. Lihat Operasi skema. VVR 8.0.6+.
tableName Ya STRING Nama tabel MaxCompute.
accessId Ya STRING ID AccessKey yang digunakan untuk mengakses MaxCompute. Lihat Cara melihat ID AccessKey dan Rahasia AccessKey saya?
Penting

Simpan ID AccessKey sebagai variabel. Lihat Kelola variabel.

accessKey Ya STRING Rahasia AccessKey yang digunakan untuk mengakses MaxCompute.
partition Tidak STRING Nama partisi di tabel MaxCompute. Tidak diperlukan untuk tabel non-partisi atau sumber inkremental. Lihat bagian "Bagaimana cara mengonfigurasi opsi partition?" dalam FAQ tentang penyimpanan hulu dan hilir.
compressAlgorithm Tidak SNAPPY STRING Algoritma kompresi untuk Tunnel MaxCompute. Nilai valid: RAW (tanpa kompresi), ZLIB, SNAPPY. SNAPPY meningkatkan throughput sekitar 50% dibandingkan ZLIB dalam skenario pengujian.
quotaName Tidak STRING Nama kuota untuk kelompok sumber daya Tunnel MaxCompute eksklusif. VVR 8.0.3+. Jika ditentukan, hapus tunnelEndpoint — jika tidak, tunnel yang ditentukan oleh tunnelEndpoint akan diprioritaskan.

Opsi sumber

Opsi Wajib Bawaan Tipe Deskripsi
maxPartitionCount Tidak 100 INTEGER Jumlah maksimum partisi yang akan dibaca. Jika melebihi batas, muncul error "The number of matched partitions exceeds the default limit". Membaca dari terlalu banyak partisi dapat membebani MaxCompute dan memperlambat startup pekerjaan — tingkatkan nilai ini hanya jika workload Anda memerlukannya.
useArrow Tidak false BOOLEAN Baca data menggunakan format Arrow, yang memanggil API penyimpanan MaxCompute. Hanya untuk penerapan batch. VVR 8.0.8+.
splitSize Tidak 256 MB MEMORYSIZE Jumlah data yang ditarik per split saat menggunakan format Arrow. Hanya untuk penerapan batch. VVR 8.0.8+.
compressCodec Tidak "" (none) STRING Algoritma kompresi saat membaca dengan format Arrow. Nilai valid: "" (none), ZSTD, LZ4_FRAME. Menentukan kodek meningkatkan throughput dibandingkan tanpa kompresi. Hanya untuk penerapan batch. VVR 8.0.8+.
dynamicLoadBalance Tidak false BOOLEAN Aktifkan alokasi shard dinamis untuk meningkatkan kinerja pemrosesan dan mengurangi waktu baca keseluruhan. Perhatikan bahwa ini dapat menyebabkan kesenjangan data karena operator berbeda membaca jumlah data yang tidak konsisten. Hanya untuk penerapan batch. VVR 8.0.8+.

Opsi sumber inkremental

Sumber inkremental melakukan polling ke MaxCompute secara berkala untuk menemukan partisi baru. Sebelum membaca partisi baru, semua penulisan data ke partisi tersebut harus selesai. Untuk detailnya, lihat bagian "Apa yang harus saya lakukan jika sumber inkremental mendeteksi partisi baru saat data masih ditulis?" dalam FAQ tentang penyimpanan hulu dan hilir.

Pengurutan partisi: Sumber membaca partisi yang urutan abjadnya lebih besar atau sama dengan nilai startPartition. Misalnya, year=2023,month=10 diurutkan sebelum year=2023,month=9 secara abjad, jadi gunakan padding nol pada nilai bulan (gunakan year=2023,month=09 bukan year=2023,month=9) untuk memastikan pengurutan yang benar.

Opsi Wajib Bawaan Tipe Deskripsi
startPartition Ya STRING Partisi awal untuk pembacaan inkremental. Saat ditentukan, partition diabaikan. Untuk tabel berpartisi multi-level, konfigurasikan nilai kolom partisi dalam urutan menurun berdasarkan level. Lihat bagian "Bagaimana cara mengonfigurasi startPartition?" dalam FAQ tentang penyimpanan hulu dan hilir.
subscribeIntervalInSec Tidak 30 INTEGER Interval polling dalam detik.
modifiedTableOperation Tidak NONE Enum Aksi saat partisi dimodifikasi selama pembacaan. Sesi unduh disimpan di checkpoint; jika data di partisi berubah setelah sesi dimulai, melanjutkan dari checkpoint gagal dan penerapan restart berulang kali. Nilai valid: NONE — perbarui startPartition untuk melewati partisi yang tidak tersedia dan restart tanpa state; SKIP — otomatis lewati partisi yang tidak tersedia saat melanjutkan. VVR 8.0.3+. Jika diatur ke salah satu nilai tersebut, data yang sudah dibaca dari partisi yang dimodifikasi tetap dipertahankan; data yang belum dibaca dibuang.

Opsi sink

Opsi Wajib Bawaan Tipe Deskripsi
useStreamTunnel Tidak false BOOLEAN Gunakan MaxCompute Streaming Tunnel alih-alih Batch Tunnel. true: Streaming Tunnel; false: Batch Tunnel. Lihat Pilih tunnel.
flushIntervalMs Tidak 30000 (30 dtk) LONG Interval flush buffer writer tunnel, dalam milidetik. Untuk Streaming Tunnel: data yang diflush langsung tersedia. Untuk Batch Tunnel: data tersedia hanya setelah checkpointing — atur ke 0 untuk menonaktifkan flushing terjadwal. Dipicu ketika flushIntervalMs atau batchSize tercapai.
batchSize Tidak 67108864 (64 MB) LONG Ukuran buffer dalam byte. Data diflush saat buffer mencapai ukuran ini. Dipicu ketika batchSize atau flushIntervalMs tercapai.
numFlushThreads Tidak 1 INTEGER Jumlah thread yang digunakan untuk flush buffer writer tunnel. Nilai lebih dari 1 memungkinkan flushing konkuren di beberapa partisi.
slotNum Tidak 0 INTEGER Jumlah slot Tunnel untuk menerima data dari Flink. Lihat Ikhtisar layanan transmisi data untuk batas slot.
dynamicPartitionLimit Tidak 100 INTEGER Jumlah maksimum partisi dinamis yang ditulis antara dua checkpoint. Jika melebihi batas, muncul error "Too many dynamic partitions". Menulis ke banyak partisi meningkatkan beban pada MaxCompute dan memperlambat checkpointing — tingkatkan nilai ini hanya jika workload Anda memerlukannya.
retryTimes Tidak 3 INTEGER Jumlah maksimum percobaan ulang untuk permintaan server MaxCompute (kegagalan pembuatan sesi, pengiriman, atau flush).
sleepMillis Tidak 1000 INTEGER Interval percobaan ulang dalam milidetik.
enableUpsert Tidak false BOOLEAN Gunakan MaxCompute Upsert Tunnel. true: memproses catatan INSERT, UPDATE_AFTER, dan DELETE; false: menggunakan tunnel yang ditentukan oleh useStreamTunnel. VVR 8.0.6+. Jika sink mengalami error atau gangguan berkepanjangan selama commit sesi dalam mode upsert, atur paralelisme operator sink menjadi 10 atau kurang.
upsertAsyncCommit Tidak false BOOLEAN Gunakan mode asinkron saat melakukan commit sesi upsert. Mode asinkron mengurangi waktu commit, tetapi data yang dicommit tidak langsung dapat dikueri. VVR 8.0.6+.
upsertCommitTimeoutMs Tidak 120000 (120 dtk) INTEGER Timeout untuk commit sesi upsert, dalam milidetik. VVR 8.0.6+.
sink.operation Tidak insert STRING Mode tulis untuk tabel Delta. insert: mode append; upsert: mode update. VVR 8.0.10+.
sink.parallelism Tidak INTEGER Paralelisme tulis untuk tabel Delta. Default mengikuti paralelisme hulu. Nilai write.bucket.num harus merupakan kelipatan integral dari sink.parallelism untuk kinerja tulis dan efisiensi memori optimal. VVR 8.0.10+.
sink.file-cached.enable Tidak false BOOLEAN Aktifkan mode cache file saat menulis ke partisi dinamis tabel Delta. Mengurangi jumlah file kecil yang ditulis ke server tetapi meningkatkan latensi tulis. Aktifkan saat sink memiliki paralelisme tinggi. VVR 8.0.10+.
sink.file-cached.writer.num Tidak 16 INTEGER Thread upload konkuren per task dalam mode cache file. Hindari mengatur nilai terlalu tinggi — menulis ke banyak partisi secara simultan dapat menyebabkan error kehabisan memori (OOM). Hanya berlaku saat sink.file-cached.enable=true. VVR 8.0.10+.
sink.bucket.check-interval Tidak 60000 INTEGER Interval pemeriksaan ukuran file dalam mode cache file, dalam milidetik. Hanya berlaku saat sink.file-cached.enable=true. VVR 8.0.10+.
sink.file-cached.rolling.max-size Tidak 16 MB MEMORYSIZE Ukuran maksimum file cache tunggal. Saat melebihi batas, data diupload ke server. Hanya berlaku saat sink.file-cached.enable=true. VVR 8.0.10+.
sink.file-cached.memory Tidak 64 MB MEMORYSIZE Memori off-heap maksimum untuk penulisan file dalam mode cache file. Hanya berlaku saat sink.file-cached.enable=true. VVR 8.0.10+.
sink.file-cached.memory.segment-size Tidak 128 KB MEMORYSIZE Ukuran segmen buffer untuk penulisan file dalam mode cache file. Hanya berlaku saat sink.file-cached.enable=true. VVR 8.0.10+.
sink.file-cached.flush.always Tidak true BOOLEAN Apakah menggunakan cache saat menulis file dalam mode cache file. Hanya berlaku saat sink.file-cached.enable=true. VVR 8.0.10+.
sink.file-cached.write.max-retries Tidak 3 INTEGER Jumlah percobaan ulang untuk upload data dalam mode cache file. Hanya berlaku saat sink.file-cached.enable=true. VVR 8.0.10+.
upsert.writer.max-retries Tidak 3 INTEGER Jumlah maksimum percobaan ulang untuk menulis ke bucket dalam sesi Upsert Writer. VVR 8.0.10+.
upsert.writer.buffer-size Tidak 64 MB MEMORYSIZE Ukuran buffer total di semua bucket dalam sesi Upsert Writer. Data diflush saat total mencapai ambang batas ini. Tingkatkan untuk efisiensi tulis yang lebih baik; turunkan jika menulis ke banyak partisi menyebabkan error OOM. VVR 8.0.10+.
upsert.writer.bucket.buffer-size Tidak 1 MB MEMORYSIZE Ukuran buffer per bucket dalam sesi Upsert Writer. Turunkan jika memori server Flink tidak mencukupi. VVR 8.0.10+.
upsert.write.bucket.num Ya INTEGER Jumlah bucket untuk tabel Delta target. Harus sesuai dengan write.bucket.num yang dikonfigurasi pada tabel Delta. VVR 8.0.10+.
upsert.write.slot-num Tidak 1 INTEGER Slot tunnel per sesi upsert. VVR 8.0.10+.
upsert.commit.max-retries Tidak 3 INTEGER Jumlah maksimum percobaan ulang untuk commit sesi upsert. VVR 8.0.10+.
upsert.commit.thread-num Tidak 16 INTEGER Paralelisme untuk commit sesi upsert. Hindari nilai besar — commit konkuren berlebihan meningkatkan konsumsi sumber daya dan dapat menyebabkan masalah kinerja. VVR 8.0.10+.
upsert.commit.timeout Tidak 600 INTEGER Timeout untuk commit sesi upsert, dalam detik. VVR 8.0.10+.
upsert.flush.concurrent Tidak 2 INTEGER Jumlah maksimum flush bucket konkuren per partisi. Setiap flush bucket menempati satu slot Tunnel. VVR 8.0.10+.
insert.commit.thread-num Tidak 16 INTEGER Paralelisme untuk commit sesi insert. VVR 8.0.10+.
insert.arrow-writer.enable Tidak false BOOLEAN Gunakan format Arrow untuk insert. VVR 8.0.10+.
insert.arrow-writer.batch-size Tidak 512 INTEGER Jumlah maksimum baris per batch format Arrow. VVR 8.0.10+.
insert.arrow-writer.flush-interval Tidak 100000 INTEGER Interval flush writer dalam milidetik. VVR 8.0.10+.
insert.writer.buffer-size Tidak 64 MB MEMORYSIZE Ukuran cache untuk writer buffered. VVR 8.0.10+.
upsert.partial-column.enable Tidak false BOOLEAN Perbarui hanya kolom yang ditentukan (pembaruan kolom parsial). Hanya berlaku untuk sink tabel Delta. Lihat Perbarui data di kolom tertentu. Saat true: jika catatan dengan primary key yang sama ada, field non-null yang ditentukan akan ditimpa; jika tidak ada catatan yang cocok, catatan baru dimasukkan dengan nilai baru untuk kolom yang ditentukan dan null untuk semua kolom yang tidak ditentukan. VVR 8.0.11+.

Opsi tabel dimensi

Saat penerapan dimulai, tabel dimensi memuat semua data dari partisi yang ditentukan oleh partition. Opsi partition mendukung fungsi max_pt(). Saat cache dimuat ulang, partisi terbaru dibaca ulang. Atur partition ke max_two_pt() untuk memuat data dari dua partisi.

Tabel dimensi memerlukan cache=ALL. Tingkatkan memori node join minimal empat kali ukuran data tabel remote. Untuk tabel dimensi besar, gunakan petunjuk SHUFFLE_HASH untuk mendistribusikan data secara merata. Untuk tabel sangat besar yang menyebabkan sering terjadi garbage collection (GC) JVM, ubah menjadi tabel dimensi key-value dengan kebijakan cache LRU (Least Recently Used)—misalnya, tabel dimensi ApsaraDB for HBase.
Opsi Wajib Bawaan Tipe Deskripsi
cache Ya STRING Kebijakan cache. Harus diatur ke ALL dan dideklarasikan secara eksplisit dalam pernyataan DDL. Semua data tabel dimensi dimuat ke cache sebelum penerapan dijalankan. Lookup berikutnya hanya mencari cache. Cache dimuat ulang setelah entri kedaluwarsa.
cacheSize Tidak 100000 LONG Jumlah maksimum baris yang dicache. Jika melebihi batas, muncul error "Row count of table <table-name> partition <partition-name> exceeds maxRowCount limit". Cache besar mengonsumsi memori heap JVM signifikan dan memperlambat startup serta refresh cache — tingkatkan nilai ini hanya jika workload Anda memerlukannya.
cacheTTLMs Tidak Long.MAX_VALUE LONG Timeout cache dalam milidetik.
cacheReloadTimeBlackList Tidak STRING Periode waktu saat cache tidak direfresh. Gunakan selama periode trafik puncak (seperti acara promosi) untuk mencegah ketidakstabilan penerapan akibat refresh cache. Lihat bagian "Bagaimana cara mengonfigurasi cacheReloadTimeBlackList?" dalam FAQ tentang penyimpanan hulu dan hilir.
maxLoadRetries Tidak 10 INTEGER Jumlah maksimum percobaan ulang untuk pemuatan cache awal saat startup penerapan. Jika percobaan ulang habis, penerapan gagal.

Pemetaan tipe data

Untuk daftar lengkap tipe data MaxCompute, lihat Sistem tipe data MaxCompute versi 2.0.

Tipe MaxCompute Tipe Flink
BOOLEAN BOOLEAN
TINYINT TINYINT
SMALLINT SMALLINT
INT INTEGER
BIGINT BIGINT
FLOAT FLOAT
DOUBLE DOUBLE
DECIMAL(precision, scale) DECIMAL(precision, scale)
CHAR(n) CHAR(n)
VARCHAR(n) VARCHAR(n)
STRING STRING
BINARY BYTES
DATE DATE
DATETIME TIMESTAMP(3)
TIMESTAMP TIMESTAMP(9)
TIMESTAMP_NTZ TIMESTAMP(9)
ARRAY ARRAY
MAP MAP
STRUCT ROW
JSON STRING
Penting

Jika tabel fisik MaxCompute berisi field tipe komposit bersarang (ARRAY, MAP, atau STRUCT) atau field JSON, atur tblproperties('columnar.nested.type'='true') saat membuat tabel agar Realtime Compute for Apache Flink dapat membaca dan menulis data dengan benar.

Flink CDC (pratinjau publik)

Konektor MaxCompute dapat mengingesti data Change Data Capture (CDC) sebagai sink dalam pekerjaan berbasis YAML. Memerlukan VVR 11.1+.

Sintaksis

source:
  type: xxx

sink:
  type: maxcompute
  name: MaxComputeSink
  access-id: ${your_accessId}
  access-key: ${your_accessKey}
  endpoint: ${your_maxcompute_endpoint}
  project: ${your_project}
  buckets-num: 8

Opsi konfigurasi

Opsi Wajib Bawaan Tipe Deskripsi
type Ya String Atur ke maxcompute.
name Tidak String Nama sink.
access-id Ya String ID AccessKey Akun Alibaba Cloud atau Pengguna RAM Anda. Dapatkan dari Konsol Resource Access Management (RAM).
access-key Ya String Rahasia AccessKey.
endpoint Ya String Titik akhir MaxCompute. Konfigurasikan berdasarkan wilayah dan metode koneksi jaringan. Lihat Titik akhir.
project Ya String Nama proyek MaxCompute. Untuk menemukannya: login ke Konsol MaxCompute, buka Workspace > Projects, lalu salin nama proyek.
tunnel.endpoint Tidak String Titik akhir Tunnel MaxCompute. Biasanya disimpulkan secara otomatis. Diperlukan di lingkungan jaringan khusus, seperti dengan server proxy.
quota.name Tidak String Nama kuota untuk kelompok sumber daya eksklusif. Jika tidak ditentukan, kelompok sumber daya bersama digunakan.
sts-token Tidak String Token Layanan Keamanan (STS) untuk otentikasi peran RAM. Diperlukan saat mengakses MaxCompute dengan peran RAM.
buckets-num Tidak 16 Integer Jumlah bucket untuk tabel Delta MaxCompute yang dibuat otomatis. Lihat Gudang data hampir real-time.
compress.algorithm Tidak zlib String Algoritma kompresi data. Nilai valid: raw (tanpa kompresi), zlib, snappy.
total.buffer-size Tidak 64 MB String Ukuran buffer di memori. Untuk tabel berpartisi: berlaku per partisi. Untuk tabel non-partisi: berlaku per tabel. Buffer untuk partisi atau tabel berbeda bersifat independen. Data diflush saat buffer penuh.
bucket.buffer-size Tidak 4 MB String Ukuran buffer per bucket. Hanya berlaku saat menulis ke tabel Delta MaxCompute.
commit.thread-num Tidak 16 Integer Jumlah maksimum partisi atau tabel yang dicommit secara konkuren selama checkpointing.
flush.concurrent-num Tidak 4 Integer Jumlah maksimum bucket yang diflush secara konkuren. Hanya berlaku saat menulis ke tabel Delta MaxCompute.

Pemetaan lokasi tabel

Saat konektor membuat tabel secara otomatis di MaxCompute, lokasi dipetakan sebagai berikut:

Penting

Jika fitur skema dinonaktifkan untuk proyek MaxCompute Anda, konektor mengabaikan tableId.namespace. Dalam kasus ini, hanya satu database (atau ekuivalen logisnya) yang diingesti ke MaxCompute—misalnya, hanya satu database MySQL saat mengingesti dari MySQL.

Lokasi MySQL Abstraksi Flink CDC Lokasi MaxCompute
N/A Proyek (dari konfigurasi) Proyek
Database TableId.namespace Skema (diabaikan jika skema dinonaktifkan)
Tabel TableId.tableName Tabel

Pemetaan tipe data

Tipe Flink CDC Tipe MaxCompute
CHAR STRING
VARCHAR STRING
BOOLEAN BOOLEAN
BINARY/VARBINARY BINARY
DECIMAL DECIMAL
TINYINT TINYINT
SMALLINT SMALLINT
INTEGER INTEGER
BIGINT BIGINT
FLOAT FLOAT
DOUBLE DOUBLE
TIME_WITHOUT_TIME_ZONE STRING
DATE DATE
TIMESTAMP_WITHOUT_TIME_ZONE TIMESTAMP_NTZ
TIMESTAMP_WITH_LOCAL_TIME_ZONE (precision > 3) TIMESTAMP
TIMESTAMP_WITH_LOCAL_TIME_ZONE (precision <= 3) DATETIME
TIMESTAMP_WITH_TIME_ZONE (precision > 3) TIMESTAMP
TIMESTAMP_WITH_TIME_ZONE (precision <= 3) DATETIME
ARRAY ARRAY
MAP MAP
ROW STRUCT

Contoh

SQL API

Tabel sumber

Baca semua data dari partisi

Baca semua data dari partisi yang ditentukan oleh partition:

CREATE TEMPORARY TABLE odps_source (
  cid VARCHAR,
  rt DOUBLE
) WITH (
  'connector' = 'odps',
  'endpoint' = '<yourEndpointName>',
  'project' = '<yourProjectName>',
  'tableName' = '<yourTableName>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}',
  'partition' = 'ds=201809*'
);

CREATE TEMPORARY TABLE blackhole_sink (
  cid VARCHAR,
  invoke_count BIGINT
) WITH (
  'connector' = 'blackhole'
);

INSERT INTO blackhole_sink
SELECT
   cid,
   COUNT(*) AS invoke_count
FROM odps_source GROUP BY cid;
Baca data inkremental

Baca data mulai dari partisi yang ditentukan oleh startPartition dan pantau partisi baru secara berkelanjutan:

CREATE TEMPORARY TABLE odps_source (
  cid VARCHAR,
  rt DOUBLE
) WITH (
  'connector' = 'odps',
  'endpoint' = '<yourEndpointName>',
  'project' = '<yourProjectName>',
  'tableName' = '<yourTableName>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}',
  'startPartition' = 'yyyy=2018,MM=09,dd=05' -- Mulai membaca dari partisi 20180905.
);

CREATE TEMPORARY TABLE blackhole_sink (
  cid VARCHAR,
  invoke_count BIGINT
) WITH (
  'connector' = 'blackhole'
);

INSERT INTO blackhole_sink
SELECT cid, COUNT(*) AS invoke_count
FROM odps_source GROUP BY cid;

Tabel sink

Tulis ke partisi statis

Tulis ke partisi yang ditentukan oleh partition:

CREATE TEMPORARY TABLE datagen_source (
  id INT,
  len INT,
  content VARCHAR
) WITH (
  'connector' = 'datagen'
);

CREATE TEMPORARY TABLE odps_sink (
  id INT,
  len INT,
  content VARCHAR
) WITH (
  'connector' = 'odps',
  'endpoint' = '<yourEndpoint>',
  'project' = '<yourProjectName>',
  'tableName' = '<yourTableName>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}',
  'partition' = 'ds=20180905' -- Tulis ke partisi 20180905.
);

INSERT INTO odps_sink
SELECT
  id, len, content
FROM datagen_source;
Tulis ke partisi dinamis

Tulis data ke partisi yang ditentukan saat runtime berdasarkan nilai di kolom ds:

CREATE TEMPORARY TABLE datagen_source (
  id INT,
  len INT,
  content VARCHAR,
  c TIMESTAMP
) WITH (
  'connector' = 'datagen'
);

CREATE TEMPORARY TABLE odps_sink (
  id  INT,
  len INT,
  content VARCHAR,
  ds VARCHAR -- Kolom partisi dinamis.
) WITH (
  'connector' = 'odps',
  'endpoint' = '<yourEndpoint>',
  'project' = '<yourProjectName>',
  'tableName' = '<yourTableName>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}',
  'partition' = 'ds' -- Abaikan nilai; data diarahkan ke partisi berdasarkan field ds.
);

INSERT INTO odps_sink
SELECT
   id,
   len,
   content,
   DATE_FORMAT(c, 'yyMMdd') as ds
FROM datagen_source;

Tabel dimensi

Kunci bernilai tunggal

Tentukan primary key saat setiap kunci memetakan tepat ke satu baris:

CREATE TEMPORARY TABLE datagen_source (
  k INT,
  v VARCHAR
) WITH (
  'connector' = 'datagen'
);

CREATE TEMPORARY TABLE odps_dim (
  k INT,
  v VARCHAR,
  PRIMARY KEY (k) NOT ENFORCED  -- Tentukan primary key.
) WITH (
  'connector' = 'odps',
  'endpoint' = '<yourEndpoint>',
  'project' = '<yourProjectName>',
  'tableName' = '<yourTableName>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}',
  'partition' = 'ds=20180905',
  'cache' = 'ALL'
);

CREATE TEMPORARY TABLE blackhole_sink (
  k VARCHAR,
  v1 VARCHAR,
  v2 VARCHAR
) WITH (
  'connector' = 'blackhole'
);

INSERT INTO blackhole_sink
SELECT k, s.v, d.v
FROM datagen_source AS s
INNER JOIN odps_dim FOR SYSTEM_TIME AS OF PROCTIME() AS d ON s.k = d.k;
Kunci bernilai ganda

Abaikan primary key saat kunci dapat memetakan ke beberapa baris:

CREATE TEMPORARY TABLE datagen_source (
  k INT,
  v VARCHAR
) WITH (
  'connector' = 'datagen'
);

CREATE TEMPORARY TABLE odps_dim (
  k INT,
  v VARCHAR
  -- Primary key tidak diperlukan untuk lookup bernilai ganda.
) WITH (
  'connector' = 'odps',
  'endpoint' = '<yourEndpoint>',
  'project' = '<yourProjectName>',
  'tableName' = '<yourTableName>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}',
  'partition' = 'ds=20180905',
  'cache' = 'ALL'
);

CREATE TEMPORARY TABLE blackhole_sink (
  k VARCHAR,
  v1 VARCHAR,
  v2 VARCHAR
) WITH (
  'connector' = 'blackhole'
);

INSERT INTO blackhole_sink
SELECT k, s.v, d.v
FROM datagen_source AS s
INNER JOIN odps_dim FOR SYSTEM_TIME AS OF PROCTIME() AS d ON s.k = d.k;

DataStream API

Penting
  • Untuk menggunakan DataStream API dengan MaxCompute, konfigurasikan konektor DataStream. Lihat Integrasikan konektor DataStream.

  • VVR 6.0.6+ mendukung debugging lokal program DataStream dengan konektor MaxCompute hingga 30 menit. Sesi yang melebihi 30 menit akan dihentikan dengan error. Lihat Debug konektor secara lokal.

  • Membaca dari tabel Delta MaxCompute (tabel yang dibuat dengan primary key dan transactional=true) tidak didukung.

Deklarasikan tabel MaxCompute menggunakan SQL, lalu akses melalui Table API atau DataStream API.

Koneksi ke tabel sumber

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tEnv = StreamTableEnvironment.create(env);
tEnv.executeSql(String.join(
    "\n",
    "CREATE TEMPORARY TABLE IF NOT EXISTS odps_source (",
    "  cid VARCHAR,",
    "  rt DOUBLE",
    ") WITH (",
    "  'connector' = 'odps',",
    "  'endpoint' = '<yourEndpointName>',",
    "  'project' = '<yourProjectName>',",
    "  'tableName' = '<yourTableName>',",
    "  'accessId' = '<yourAccessId>',",
    "  'accessKey' = '<yourAccessPassword>',",
    "  'partition' = 'ds=201809*'",
    ")");
DataStream<Row> source = tEnv.toDataStream(tEnv.from("odps_source"));
source.print();
env.execute("odps source");

Hubungkan ke sink

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tEnv = StreamTableEnvironment.create(env);
tEnv.executeSql(String.join(
    "\n",
    "CREATE TEMPORARY TABLE IF NOT EXISTS odps_sink (",
    "  cid VARCHAR,",
    "  rt DOUBLE",
    ") WITH (",
    "  'connector' = 'odps',",
    "  'endpoint' = '<yourEndpointName>',",
    "  'project' = '<yourProjectName>',",
    "  'tableName' = '<yourTableName>',",
    "  'accessId' = '<yourAccessId>',",
    "  'accessKey' = '<yourAccessPassword>',",
    "  'partition' = 'ds=20180905'",
    ")");
DataStream<Row> data = env.fromElements(
    Row.of("id0", 3.),
    Row.of("id1", 4.));
tEnv.fromDataStream(data).insertInto("odps_sink").execute();

Dependensi Maven

Tambahkan konektor DataStream MaxCompute ke proyek Anda. Semua versi tersedia di repositori pusat Maven.

<dependency>
    <groupId>com.alibaba.ververica</groupId>
    <artifactId>ververica-connector-odps</artifactId>
    <version>${vvr-version}</version>
</dependency>

Langkah berikutnya