All Products
Search
Document Center

Realtime Compute for Apache Flink:StarRocks

Last Updated:Aug 22, 2026

Pelajari cara menggunakan konektor StarRocks.

Latar Belakang

StarRocks adalah gudang data Pemrosesan Paralel Masif (MPP) generasi berikutnya yang menawarkan kinerja sangat cepat di semua skenario serta pengalaman analitik terpadu. StarRocks memiliki keunggulan berikut:

  • Kompatibel dengan protokol MySQL, sehingga Anda dapat menggunakan klien MySQL dan alat Business Intelligence (BI) umum untuk menghubungkan dan menganalisis data.

  • Menggunakan arsitektur terdistribusi:

    • Tabel data dipartisi secara horizontal dan disimpan dalam beberapa replika.

    • Kluster dapat diskalakan secara fleksibel dan mampu menganalisis hingga 10 petabyte (PB) data.

    • Memanfaatkan framework MPP untuk mempercepat komputasi paralel.

    • Mendukung beberapa replika guna menyediakan toleransi kesalahan.

Konektor Flink menyimpan data dalam cache dan menggunakan Stream Load untuk menulisnya secara batch ke tabel sink. Konektor ini membaca dari tabel sumber dengan mengambil data secara batch. Tabel berikut mencantumkan kemampuan konektor StarRocks.

Category

Description

Supported types

tabel sumber, tabel dimensi, tabel sink, dan target ingesti data

Execution mode

streaming mode dan batch mode

Data format

CSV

Connector-specific metrics

None

API types

DataStream, SQL, dan YAML untuk ingesti data

Support for updates/deletions in sink tables

Yes

Prasyarat

Anda harus memiliki kluster StarRocks yang dideploy di EMR atau yang dikelola sendiri di ECS.

Batasan

  • Hanya Ververica Runtime (VVR) 11.1 atau yang lebih baru yang mendukung join dengan tabel dimensi.

  • Untuk menghindari pembatasan akses jaringan, tambahkan port kluster StarRocks berikut ke daftar putih security group atau firewall: 9030, 8030, 8040, 9060, 8060, 9020.

  • Jika tabel StarRocks target berisi kolom generated yang tersembunyi, deklarasikan hanya kolom fisik aktual yang akan ditulis dalam pekerjaan Flink dan kecualikan kolom generated tersebut. Kolom generated tersembunyi yang dibuat secara otomatis oleh ekspresi partisi StarRocks (seperti __generated_partition_column_0) tidak menerima penulisan eksternal. Secara default, konektor membuat permintaan tulis berdasarkan skema lengkap, yang menyebabkan pekerjaan gagal. Untuk informasi lebih lanjut tentang kolom generated di StarRocks, lihat Generated columns.

SQL

Fitur

StarRocks pada E-MapReduce mendukung pernyataan CREATE TABLE AS SELECT (CTAS) dan CREATE DATABASE AS SELECT (CDAS). CTAS menyinkronkan skema dan data dari satu tabel, sedangkan CDAS menyinkronkan seluruh database atau beberapa tabel dalam database yang sama. Untuk informasi lebih lanjut, lihat Gunakan pernyataan CTAS dan CDAS di Realtime Compute for Apache Flink untuk menyinkronkan data dari database MySQL ke StarRocks.

Sintaks

CREATE TABLE USER_RESULT(
 name VARCHAR,
 score BIGINT
 ) WITH (
 'connector' = 'starrocks',
 'jdbc-url'='jdbc:mysql://fe1_ip:query_port,fe2_ip:query_port,fe3_ip:query_port?xxxxx',
 'load-url'='fe1_ip:http_port;fe2_ip:http_port;fe3_ip:http_port',
 'database-name' = 'xxx',
 'table-name' = 'xxx',
 'username' = 'xxx',
 'password' = 'xxx'
 );

Parameter

Type

Parameter

Description

Type

Required

Default

Remarks

General

connector

Menentukan konektor yang akan digunakan.

String

Yes

Nilainya harus starrocks.

jdbc-url

URL Java Database Connectivity (JDBC).

String

Yes

Tentukan alamat IP dan port JDBC FE dalam format jdbc:mysql://ip:port.

database-name

Nama database StarRocks.

String

Yes

table-name

Nama tabel StarRocks.

String

Yes

username

Username untuk menghubungkan ke StarRocks.

String

Yes

password

Password untuk menghubungkan ke StarRocks.

String

Yes

starrocks.create.table.properties

Menentukan properti untuk pembuatan tabel otomatis.

String

No

Menentukan properti awal tabel, seperti engine dan jumlah replika. Contoh: 'starrocks.create.table.properties' = 'buckets 8' atau 'starrocks.create.table.properties' = 'replication_num=1'.

Source-specific

scan-url

URL pemindaian data.

String

No

Tentukan alamat IP dan port HTTP FE. Format: fe_ip:http_port;fe_ip:http_port.

Catatan

Untuk menentukan beberapa alamat IP dan port, pisahkan dengan titik koma (;).

scan.connect.timeout-ms

Timeout untuk flink-connector-starrocks menghubungkan ke StarRocks.

Konektor melaporkan error jika koneksi tidak berhasil dibuat dalam batas waktu ini.

String

No

1000

Unit: milidetik.

scan.params.keep-alive-min

Durasi keep-alive untuk tugas kueri.

String

No

10

scan.params.query-timeout-s

Timeout untuk tugas kueri.

Jika tidak ada hasil yang dikembalikan dalam periode ini, sistem menghentikan tugas kueri.

String

No

600

Unit: detik.

scan.params.mem-limit-byte

Batas memori untuk satu kueri pada node BE.

String

No

1073741824 (1 GB)

Unit: byte.

scan.max-retries

Jumlah maksimum percobaan ulang untuk kueri yang gagal.

Konektor melaporkan error jika batas ini terlampaui.

String

No

1

Sink-specific

load-url

URL impor data.

String

Yes

Tentukan alamat IP dan port HTTP FE dalam format fe_ip:http_port;fe_ip:http_port.

Catatan

Untuk menentukan beberapa alamat IP dan port, pisahkan dengan titik koma (;).

sink.semantic

Semantik pengiriman untuk penulisan.

String

No

at-least-once

Nilai yang valid:

  • at-least-once (default): Menjamin data dikirim setidaknya sekali.

  • exactly-once: Menjamin data dikirim tepat sekali.

sink.buffer-flush.max-bytes

Jumlah maksimum data yang disimpan dalam buffer sebelum flush.

String

No

94371840 (90 MB)

Rentang valid: 64 MB hingga 10 GB.

sink.buffer-flush.max-rows

Jumlah maksimum baris yang disimpan dalam buffer sebelum flush.

String

No

500000

Rentang valid: 1.000 hingga 5.000.000.

sink.buffer-flush.interval-ms

Interval flush buffer.

String

No

300000

Rentang valid: 1.000 ms hingga 3.600.000 ms.

sink.max-retries

Jumlah maksimum percobaan ulang untuk penulisan yang gagal.

String

No

3

Rentang valid: 0 hingga 1000.

sink.connect.timeout-ms

Timeout untuk menghubungkan ke StarRocks.

String

No

1000

Rentang valid: 100 hingga 60.000. Unit: milidetik.

sink.properties.*

Properti Stream Load tambahan untuk sink.

String

No

Parameter ini mengontrol perilaku Stream Load. Misalnya, sink.properties.format menentukan format data yang diimpor, seperti CSV. Untuk parameter lainnya, lihat Stream Load.

Dimension-specific

lookup.cache.enabled

Menentukan apakah caching untuk tabel dimensi diaktifkan.

Boolean

No

true

Nilai yang valid:

  • true: Mengaktifkan caching. Setelah data tabel dibaca pertama kali, sistem menyimpannya dalam memori. Permintaan berikutnya menggunakan data cache dalam periode validitasnya untuk mengurangi overhead I/O.

  • false: Menonaktifkan caching. Setiap kueri langsung mengakses sumber data.

Penting
  • Fitur ini memerlukan mesin Realtime Compute for Apache Flink VVR 11.1 dan yang lebih baru.

  • Kami menyarankan menonaktifkan fitur ini dalam skenario berikut:

    • Data tabel dimensi sering diperbarui, dan data real-time diperlukan.

    • Tabel berisi banyak data, sehingga berisiko menyebabkan overflow memori.

Pemetaan tipe data

StarRocks data type

Flink data type

NULL

NULL

BOOLEAN

BOOLEAN

TINYINT

TINYINT

SMALLINT

SMALLINT

INT

INT

BIGINT

BIGINT

BIGINT UNSIGNED

Catatan

Memerlukan mesin Realtime Compute for Apache Flink VVR 8.0.10 dan yang lebih baru.

DECIMAL(20,0)

LARGEINT

DECIMAL(20,0)

FLOAT

FLOAT

DOUBLE

DOUBLE

DATE

DATE

DATETIME

TIMESTAMP

DECIMAL

DECIMAL

DECIMALV2

DECIMAL

DECIMAL32

DECIMAL

DECIMAL64

DECIMAL

DECIMAL128

DECIMAL

CHAR(m)

Catatan
  • VVR 8.0.10 secara otomatis memperpanjang panjang CHAR tiga kali lipat (m=n*3, dengan n<=85) untuk mengakomodasi perbedaan encoding antara MySQL dan StarRocks.

  • VVR 8.0.11 dan yang lebih baru secara otomatis memperpanjang panjang CHAR empat kali lipat (m=n*4, dengan n<=63) untuk mengakomodasi perbedaan encoding antara MySQL dan StarRocks.

  • Panjang maksimum untuk tipe CHAR StarRocks adalah 255. Oleh karena itu, Flink memetakan tipe CHAR ke tipe CHAR StarRocks hanya jika panjang yang diperluas secara otomatis tidak melebihi 255.

CHAR(n)

VARCHAR(m)

Catatan
  • VVR 8.0.10 secara otomatis memperpanjang panjang VARCHAR tiga kali lipat (m=n*3, dengan n>85) untuk mengakomodasi perbedaan encoding antara MySQL dan StarRocks.

  • VVR 8.0.11 dan yang lebih baru secara otomatis memperpanjang panjang VARCHAR empat kali lipat (m=n*4, dengan n>63) untuk mengakomodasi perbedaan encoding antara MySQL dan StarRocks.

  • Panjang maksimum untuk tipe CHAR StarRocks adalah 255. Oleh karena itu, jika panjang yang diperluas secara otomatis dari tipe CHAR Flink melebihi 255, Flink memetakan tipe tersebut ke tipe VARCHAR StarRocks.

CHAR(n)

VARCHAR

STRING

VARBINARY

Catatan

Memerlukan mesin Realtime Compute for Apache Flink VVR 8.0.10 dan yang lebih baru.

VARBINARY

Contoh kode

CREATE TEMPORARY TABLE IF NOT EXISTS `runoob_tbl_source` (
  `runoob_id` BIGINT NOT NULL,
  `runoob_title` STRING NOT NULL,
  `runoob_author` STRING NOT NULL,
  `submission_date` DATE NULL
) WITH (
  'connector' = 'starrocks',
  'jdbc-url' = 'jdbc:mysql://ip:9030',
  'scan-url' = 'ip:18030',
  'database-name' = 'db_name',
  'table-name' = 'table_name',
  'password' = 'xxxxxxx',
  'username' = 'xxxxx'
);
CREATE TEMPORARY TABLE IF NOT EXISTS `runoob_tbl_sink` (
  `runoob_id` BIGINT NOT NULL,
  `runoob_title` STRING NOT NULL,
  `runoob_author` STRING NOT NULL,
  `submission_date` DATE NULL
  PRIMARY KEY(`runoob_id`)
  NOT ENFORCED
) WITH (
  'jdbc-url' = 'jdbc:mysql://ip:9030',
  'connector' = 'starrocks',
  'load-url' = 'ip:18030',
  'database-name' = 'db_name',
  'table-name' = 'table_name',
  'password' = 'xxxxxxx',
  'username' = 'xxxx',
  'sink.buffer-flush.interval-ms' = '5000'
);

INSERT INTO runoob_tbl_sink SELECT * FROM runoob_tbl_source;
Catatan

StarRocks mengizinkan kolom kunci primer menjadi NULLABLE. Namun, Flink tidak mendukung kunci primer yang berisi kolom nullable. Model konsistensi data Flink mensyaratkan bahwa kunci primer bersifat unik dan non-nullable. Jika tidak, Flink melemparkan error Invalid primary key. Column 'xxx' is nullable. Untuk informasi lebih lanjut, lihat Error "Invalid primary key. Column 'xxx' is nullable.".

Ingesti Data

Gunakan konektor StarRocks Pipeline untuk menulis catatan data dan perubahan skema dari sumber data hulu ke database StarRocks eksternal. Konektor StarRocks mendukung edisi komunitas dan EMR Serverless StarRocks yang sepenuhnya dikelola oleh Alibaba Cloud.

Fitur

  • Pembuatan database dan tabel otomatis.

    Jika database atau tabel hulu tidak ada di instans StarRocks hilir, konektor akan membuatnya secara otomatis. Anda dapat menggunakan parameter table.create.properties.* untuk mengonfigurasi opsi pembuatan tabel otomatis.

  • Sinkronisasi perubahan skema.

    Konektor StarRocks secara otomatis menerapkan event CreateTableEvent, AddColumnEvent, dan DropColumnEvent ke database hilir.

  • VVR 11.1 dan yang lebih baru mendukung perubahan tipe kolom yang kompatibel. Untuk informasi lebih lanjut, lihat ALTER TABLE | StarRocks.

Catatan penggunaan

  • Setiap tabel yang disinkronkan harus memiliki kunci primer. Untuk tabel tanpa kunci primer, Anda harus menentukannya dalam blok transform agar data dapat ditulis ke hilir. Contoh:

    transform:
      - source-table: ...
        primary-keys: id, ...
  • Untuk tabel yang dibuat secara otomatis, kunci bucket sama dengan kunci primer, dan tabel tidak dapat memiliki kunci partisi.

  • Saat menyinkronkan perubahan skema, kolom baru hanya dapat ditambahkan di akhir daftar kolom yang sudah ada. Dalam mode perubahan skema Lenient (default), penyisipan kolom di posisi lain akan secara otomatis dipindahkan ke akhir.

  • Jika Anda menggunakan versi StarRocks sebelum 2.5.7, Anda harus secara eksplisit menentukan jumlah bucket melalui parameter table.create.num-buckets. StarRocks 2.5.7 dan versi yang lebih baru dapat secara otomatis menentukan jumlah bucket yang sesuai.

  • Jika Anda menggunakan StarRocks 3.2 atau versi yang lebih baru, kami menyarankan untuk mengaktifkan opsi table.create.properties.fast_schema_evolution guna mempercepat perubahan skema.

  • Masalah streaming dapat terjadi saat Anda menggunakan CDC YAML untuk mengingesti data ke EMR Serverless StarRocks. Anda dapat menerapkan salah satu solusi berikut:

    • Gunakan konektor StarRocks Flink SQL dan atur parameter sink.version=V1.

    • Aktifkan parameter FE emr_internal_redirect.

    • Gunakan nama domain StarRocks Private Zone alih-alih SLB.

Sintaks

source:
  ...

sink:
  type: starrocks
  name: StarRocks Sink
  jdbc-url: jdbc:mysql://127.0.0.1:9030
  load-url: 127.0.0.1:8030
  username: root
  password: pass
  sink.buffer-flush.interval-ms: 5000   # Atur interval flush data.

Konfigurasi

Parameter

Description

Type

Required

Default

Remarks

type

Menentukan jenis konektor sink.

String

Yes

Atur ke starrocks.

name

Nama tampilan sink.

String

No

jdbc-url

URL JDBC untuk koneksi database.

String

Yes

Mendukung beberapa alamat yang dipisahkan koma (,). Contoh: jdbc:mysql://fe_host1:fe_query_port1,fe_host2:fe_query_port2,fe_host3:fe_query_port3.

load-url

URL HTTP node FE untuk Stream Load.

String

Yes

Mendukung beberapa alamat yang dipisahkan titik koma (;). Contoh: fe_host1:fe_http_port1;fe_host2:fe_http_port2.

username

Username untuk koneksi StarRocks.

String

Yes

User ini harus memiliki izin SELECT dan INSERT minimal pada tabel target. Anda dapat memberikan izin yang diperlukan dengan perintah StarRocks GRANT.

password

Password untuk koneksi StarRocks.

String

Yes

sink.semantic

Semantik pengiriman untuk penulisan data.

String

No

at-least-once

Hanya at-least-once yang didukung. Menyetel parameter ini secara eksplisit ke exactly-once tidak menyebabkan error, tetapi nilainya secara otomatis diatur ulang ke at-least-once. Untuk menggunakan semantik exactly-once, gunakan konektor StarRocks Flink SQL. Untuk informasi lebih lanjut, lihat bagian SQL topik ini.

sink.label-prefix

Awalan label untuk pekerjaan Stream Load.

String

No

Nilai hanya boleh berisi huruf Inggris, angka, tanda hubung (-), dan garis bawah (_). Karakter lain dapat menyebabkan load gagal.

sink.connect.timeout-ms

Timeout untuk membuat koneksi HTTP.

Integer

No

30000

Unit: milidetik. Nilai harus antara 100 dan 60000.

sink.wait-for-continue.timeout-ms

Timeout untuk menunggu respons 100 Continue dari server.

Integer

No

30000

Unit: milidetik. Nilai harus antara 3000 dan 600000.

sink.buffer-flush.max-bytes

Ukuran maksimum cache dalam memori, dalam byte, sebelum flush dipicu.

Long

No

94371840

Unit: byte. Nilai harus antara 64 MB dan 10 GB.

Catatan
  • Ukuran cache ini dibagi oleh semua tabel. Saat buffer penuh, konektor memilih beberapa tabel untuk di-flush.

  • Mengatur nilai yang lebih besar dapat meningkatkan throughput tetapi dapat meningkatkan latensi ingesti.

sink.buffer-flush.max-rows

Jumlah maksimum baris dalam cache memori sebelum flush dipicu.

Long

No

500000

Nilai harus antara 1.000 dan 5.000.000.

sink.buffer-flush.interval-ms

Interval waktu antar flush untuk buffer setiap tabel.

Long

No

300000

Unit: milidetik.

Catatan

Untuk pekerjaan yang menyinkronkan data dalam jumlah kecil, kurangi nilai ini untuk menghindari penundaan lama sebelum data dipersisten.

sink.max-retries

Jumlah maksimum percobaan ulang.

Long

No

3

Nilai harus antara 0 dan 1000.

sink.scan-frequency.ms

Frekuensi konektor memeriksa apakah buffer perlu di-flush.

Long

No

50

Unit: milidetik.

sink.io.thread-count

Jumlah thread yang digunakan untuk Stream Load.

Integer

No

2

sink.at-least-once.use-transaction-stream-load

Menentukan apakah akan menggunakan antarmuka transaksi Stream Load untuk ingesti data.

Boolean

No

true

Opsi ini hanya berlaku jika database mendukungnya.

sink.ignore.update-before

Menentukan apakah akan mengabaikan catatan update-before dalam operasi pembaruan.

Boolean

No

true

Jika kunci primer diubah melalui modul Transform (misalnya, ketika primary-keys menentukan kunci primer yang berbeda dari hulu), atur sink.ignore.update-before ke false. Jika tidak, baris yang sesuai dengan kunci primer lama tidak dihapus, sehingga menghasilkan data usang.

Hanya Ververica Runtime (VVR) 11.8 atau yang lebih baru yang mendukung parameter ini.

sink.ignore.delete

Menentukan apakah akan mengabaikan catatan penghapusan.

Boolean

No

false

Jika Anda mengatur parameter ini ke true, catatan penghapusan difilter dan tidak ditulis ke StarRocks. Gunakan pengaturan ini untuk menyimpan data historis di sink dan hanya menyinkronkan operasi insert dan update.

Hanya Ververica Runtime (VVR) 11.8 atau yang lebih baru yang mendukung parameter ini.

sink.properties.*

Properti tambahan untuk sink.

String

No

Untuk properti yang didukung, lihat STREAM LOAD.

table.create.num-buckets

Jumlah bucket untuk tabel yang dibuat secara otomatis.

Integer

No

table.create.properties.*

Properti tambahan untuk pembuatan tabel otomatis.

String

No

Sebagai contoh, Anda dapat meneruskan 'table.create.properties.fast_schema_evolution' = 'true' untuk mengaktifkan perubahan skema cepat. Untuk detailnya, lihat dokumentasi StarRocks.

table.schema-change.timeout

Timeout untuk operasi perubahan skema.

Duration

No

30 min

Harus merupakan jumlah detik bilangan bulat.

Catatan

Jika operasi perubahan skema melebihi batas ini, pekerjaan gagal.

unicode-char.max-bytes

Jumlah byte yang dialokasikan untuk setiap karakter Unicode.

Integer

No

3

Dalam CDC, panjang tipe VARCHAR diukur dalam karakter, sedangkan di StarRocks, panjang tipe VARCHAR diukur dalam byte.

Dalam kebanyakan kasus, karakter Unicode tidak melebihi 3 byte setelah encoding UTF-8. Namun, beberapa karakter langka dan simbol emoji dapat menempati 4 byte atau lebih.

sink.socket.time

Timeout klien HTTP untuk meng-flush data ke StarRocks.

Long

No

-1

Timeout klien HTTP, dalam milidetik, untuk mengirim permintaan Stream Load saat data di-flush ke StarRocks. Nilai -1 menggunakan default sistem, yang berarti tidak ada timeout.

Hanya Ververica Runtime (VVR) 11.8 atau yang lebih baru yang mendukung parameter ini.

sink.close.eof-timeout-ms

Timeout untuk menutup sink.

Long

No

60000

Timeout, dalam milidetik, untuk menunggu antrian flush selesai saat pekerjaan ditutup. Hanya Ververica Runtime (VVR) 11.8 atau yang lebih baru yang mendukung parameter ini.

Gunakan katalog bawaan yang telah ada

VVR 11.5 dan yang lebih baru memungkinkan Anda mereferensikan katalog StarRocks bawaan yang dibuat di halaman Data Management secara langsung dalam pekerjaan ingesti data Flink CDC. Hal ini menyederhanakan konfigurasi dengan mengurangi jumlah properti yang perlu Anda atur secara manual.

sink:
  type: starrocks
  using.built-in-catalog: starrocks_catalog

Pekerjaan ingesti data dapat secara otomatis menggunakan kembali opsi katalog StarRocks berikut:

  • jdbc-url

  • http-url

  • username

  • password

  • table.num-buckets

Untuk mengganti nilai-nilai ini, Anda dapat secara eksplisit mengatur opsi YAML yang sesuai, yang akan memiliki prioritas lebih tinggi.

Pemetaan tipe

Catatan

StarRocks tidak mendukung semua tipe CDC YAML. Menulis tipe yang tidak didukung ke sink menyebabkan pekerjaan gagal. Anda dapat menggunakan fungsi bawaan CAST dalam transform untuk mengonversi data yang tidak didukung, atau menggunakan pernyataan proyeksi untuk menghapusnya dari tabel hasil. Untuk informasi lebih lanjut, lihat Kembangkan pekerjaan ingesti data Flink CDC.

CDC type

StarRocks type

Remarks

TINYINT

TINYINT

SMALLINT

SMALLINT

INT

INT

BIGINT

BIGINT

FLOAT

FLOAT

DOUBLE

DOUBLE

BOOLEAN

BOOLEAN

DATE

DATE

TIMESTAMP

DATETIME

TIMESTAMP_LTZ

DATETIME

DECIMAL(p, s)

DECIMAL(p, s)

Karena StarRocks tidak mendukung DECIMAL untuk kunci primer, konektor secara otomatis mengonversi kolom kunci primer DECIMAL hulu menjadi VARCHAR dalam skema StarRocks yang disinkronkan.

CHAR(n)

(n <= 85)

CHAR(n × 3)

CDC mengukur panjang dalam karakter, sedangkan StarRocks menggunakan byte. Konektor mengalikan panjang dengan 3 untuk mengakomodasi karakter UTF-8 multi-byte.

Catatan

Panjang maksimum tipe CHAR StarRocks adalah 255. Oleh karena itu, hanya tipe CHAR CDC dengan panjang hingga 85 yang dipetakan ke tipe CHAR StarRocks.

Catatan

Anda dapat mengatur parameter unicode-char.max-bytes untuk mengalokasikan lebih banyak byte untuk setiap karakter Unicode.

CHAR(n)

(n > 85)

VARCHAR(n × 3)

CDC mengukur panjang dalam karakter, sedangkan StarRocks menggunakan byte. Konektor mengalikan panjang dengan 3 untuk mengakomodasi karakter UTF-8 multi-byte.

Catatan

CDC mengukur panjang dalam karakter, sedangkan StarRocks menggunakan byte. Konektor mengalikan panjang dengan 3. Karena hasilnya melebihi batas 255 byte untuk tipe CHAR StarRocks, maka dipetakan ke VARCHAR.

Catatan

Anda dapat mengatur parameter unicode-char.max-bytes untuk mengalokasikan lebih banyak byte untuk setiap karakter Unicode.

VARCHAR(n)

VARCHAR(n × 3)

CDC mengukur panjang dalam karakter, sedangkan StarRocks menggunakan byte. Konektor mengalikan panjang dengan 3 untuk mengakomodasi karakter UTF-8 multi-byte.

Catatan

Anda dapat mengatur parameter unicode-char.max-bytes untuk mengalokasikan lebih banyak byte untuk setiap karakter Unicode.

BINARY(n)

BINARY(n+2)

Dua byte padding ditambahkan untuk memastikan integritas data.

VARBINARY(n)

VARBINARY(n+1)

Satu byte padding ditambahkan untuk memastikan integritas data.

Perubahan skema

Sebagai sink ingesti data, StarRocks mendukung event perubahan skema berikut:

  • CREATE TABLE EVENT

    Catatan

    Jika tabel StarRocks hilir sudah ada, konektor tidak mencoba membuatnya lagi. Pastikan skema tabel hilir kompatibel dengan skema hulu.

  • ADD COLUMN EVENT

    Catatan

    StarRocks mengharuskan kolom kunci primer muncul di awal tabel. Kolom baru apa pun harus ditambahkan setelahnya.

  • DROP COLUMN EVENT

  • TRUNCATE TABLE EVENT

  • DROP TABLE EVENT

Contoh kode

Contoh berikut menunjukkan konfigurasi untuk skenario umum.

Sinkronisasi satu tabel

Sinkronkan satu tabel MySQL ke StarRocks. Jika database dan tabel tujuan tidak ada, konektor secara otomatis membuat tabel kunci primer.

pipeline:
          name: MySQL to StarRocks Pipeline
      source:
        type: mysql
        name: MySQL Source
        hostname: <yourHostname>
        port: 3306
        username: <yourUsername>
        password: ${secret_values.mysql_password}
        tables: test_db.test_source_table
        server-id: 5401-5499
        # (Opsional) Sinkronkan data dari tabel yang baru ditambahkan selama fase inkremental tanpa me-restart pekerjaan.
        scan.binlog.newly-added-table.enabled: true
        # (Opsional) Sinkronkan komentar tabel dan kolom ke tujuan.
        include-comments.enabled: true
        # (Opsional) Hanya deserialize binlog untuk tabel yang ditangkap untuk meningkatkan kinerja baca.
        scan.only.deserialize.captured.tables.changelog.enabled: true
      
      sink:
        type: starrocks
        name: StarRocks Sink
        jdbc-url: jdbc:mysql://<yourFeHostname>:9030
        load-url: <yourFeHostname>:8030
        username: <yourUsername>
        password: ${secret_values.starrocks_password}
       
        # (Opsional) Untuk pekerjaan dengan data sedikit, kurangi interval flush untuk mencegah penundaan persistensi lama. Default: 300000, atau 5 menit.
        sink.buffer-flush.interval-ms: 5000
        # (Opsional) Jika set karakter hulu adalah utf8mb4, atur parameter ini ke 4 untuk mencegah pemotongan teks. Default: 3.
        unicode-char.max-bytes: 4
        # (Opsional) Jumlah bucket untuk tabel yang dibuat otomatis. Parameter ini wajib untuk versi StarRocks sebelum 2.5.7. Versi yang lebih baru dapat menyimpulkan nilainya secara otomatis.
        table.create.num-buckets: 8
        # (Opsional) Jumlah replika untuk tabel yang dibuat otomatis. Konfigurasikan nilai ini berdasarkan kluster Anda.
        table.create.properties.replication_num: 3
        # (Opsional) Untuk StarRocks 3.2 dan yang lebih baru, aktifkan opsi ini untuk mempercepat perubahan skema.
        table.create.properties.fast_schema_evolution: true
        # Catatan: Jika Anda menggunakan transform untuk mengubah kunci primer, Anda juga harus mengatur sink.ignore.update-before: false.
        # Jika tidak, baris yang terkait dengan kunci primer lama tetap ada di tujuan.
      
      pipeline:
        name: MySQL to StarRocks Pipeline

Sinkronisasi seluruh database

Sinkronkan semua tabel dalam database MySQL ke StarRocks secara bersamaan. Konektor secara otomatis membuat database tujuan dan tabel kunci primer, sehingga Anda tidak perlu membuat setiap tabel terlebih dahulu.

source:
        type: mysql
        name: MySQL Source
        hostname: <yourHostname>
        port: 3306
        username: <yourUsername>
        password: ${secret_values.mysql_password}
        # Gunakan ekspresi reguler untuk mencocokkan semua tabel dalam database. Untuk mencocokkan beberapa database, pisahkan pola dengan koma.
        tables: test_db.\.*
        server-id: 5401-5499
        # (Opsional) Sinkronkan data dari tabel yang baru ditambahkan selama fase inkremental tanpa me-restart pekerjaan.
        scan.binlog.newly-added-table.enabled: true
        # (Opsional) Sinkronkan komentar tabel dan kolom ke tujuan.
        include-comments.enabled: true
      
      sink:
        type: starrocks
        name: StarRocks Sink
        jdbc-url: jdbc:mysql://<yourFeHostname>:9030
        load-url: <yourFeHostname>:8030
        username: <yourUsername>
        password: ${secret_values.starrocks_password}
        # (Opsional) Untuk pekerjaan dengan data sedikit, kurangi interval flush untuk mencegah penundaan persistensi lama. Default: 300000, atau 5 menit.
        sink.buffer-flush.interval-ms: 5000
        # (Opsional) Jika set karakter hulu adalah utf8mb4, atur parameter ini ke 4 untuk mencegah pemotongan teks. Default: 3.
        unicode-char.max-bytes: 4
        # (Opsional) Jumlah bucket untuk tabel yang dibuat otomatis. Parameter ini wajib untuk versi StarRocks sebelum 2.5.7. Versi yang lebih baru dapat menyimpulkan nilainya secara otomatis.
        table.create.num-buckets: 8
        # (Opsional) Jumlah replika untuk tabel yang dibuat otomatis. Konfigurasikan nilai ini berdasarkan kluster Anda.
        table.create.properties.replication_num: 3
        # (Opsional) Untuk StarRocks 3.2 dan yang lebih baru, aktifkan opsi ini untuk mempercepat perubahan skema.
        table.create.properties.fast_schema_evolution: true
      
      pipeline:
        name: MySQL to StarRocks Pipeline

Kecualikan tabel tertentu saat sinkronisasi database penuh

Saat menyinkronkan seluruh database, gunakan ekspresi reguler untuk melewati tabel yang tidak ingin Anda sinkronkan ke tujuan, seperti tabel sementara atau sensitif.

source:
        type: mysql
        name: MySQL Source
        hostname: <yourHostname>
        port: 3306
        username: <yourUsername>
        password: ${secret_values.mysql_password}
        tables: test_db.\.*
        # Tabel yang cocok dengan ekspresi reguler ini tidak disinkronkan.
        tables.exclude: test_db.tmp_.\*
        server-id: 5401-5499
      
      sink:
        type: starrocks
        name: StarRocks Sink
        jdbc-url: jdbc:mysql://<yourFeHostname>:9030
        load-url: <yourFeHostname>:8030
        username: <yourUsername>
        password: ${secret_values.starrocks_password}
        # (Opsional) Versi antarmuka load. V2 memerlukan StarRocks 2.4 atau yang lebih baru. Untuk EMR Serverless StarRocks, gunakan V1 jika terjadi masalah streaming.
        sink.version: V2
        # (Opsional) Untuk pekerjaan dengan data sedikit, kurangi interval flush untuk mencegah penundaan persistensi lama. Default: 300000, atau 5 menit.
        sink.buffer-flush.interval-ms: 5000
        # (Opsional) Jumlah bucket untuk tabel yang dibuat otomatis. Parameter ini wajib untuk versi StarRocks sebelum 2.5.7.
        table.create.num-buckets: 8
        # (Opsional) Untuk StarRocks 3.2 dan yang lebih baru, aktifkan opsi ini untuk mempercepat perubahan skema.
        table.create.properties.fast_schema_evolution: true
      
      pipeline:
        name: MySQL to StarRocks Pipeline

Sinkronisasi ke database dan tabel yang ditentukan

Jika nama database atau tabel StarRocks tujuan harus berbeda dari nama hulu, seperti saat menulis ke database lapisan ODS, gunakan route untuk mengganti namanya.

source:
        type: mysql
        name: MySQL Source
        hostname: <yourHostname>
        port: 3306
        username: <yourUsername>
        password: ${secret_values.mysql_password}
        tables: test_db.\.*
        server-id: 5401-5499
      
      sink:
        type: starrocks
        name: StarRocks Sink
        jdbc-url: jdbc:mysql://<yourFeHostname>:9030
        load-url: <yourFeHostname>:8030
        username: <yourUsername>
        password: ${secret_values.starrocks_password}
        # (Opsional) Untuk pekerjaan dengan data sedikit, kurangi interval flush untuk mencegah penundaan persistensi lama. Default: 300000, atau 5 menit.
        sink.buffer-flush.interval-ms: 5000
        # (Opsional) Jumlah bucket untuk tabel yang dibuat otomatis. Parameter ini wajib untuk versi StarRocks sebelum 2.5.7.
        table.create.num-buckets: 8
      
      route:
        # Sinkronkan semua tabel dalam database test_db MySQL ke database test_db2 StarRocks tanpa mengubah nama tabel.
        # <> adalah placeholder yang diganti dengan nama tabel sumber yang cocok.
        - source-table: test_db.\.*
          sink-table: test_db2.<>
          replace-symbol: <>
      
      pipeline:
        name: MySQL to StarRocks Pipeline

Gabungkan tabel terpartisi

Gabungkan beberapa tabel terpartisi dengan skema identik menjadi satu tabel StarRocks untuk kueri dan analitik terpadu.

source:
        type: mysql
        name: MySQL Source
        hostname: <yourHostname>
        port: 3306
        username: <yourUsername>
        password: ${secret_values.mysql_password}
        # Cocokkan semua tabel terpartisi, seperti user_0 dan user_1.
        tables: test_db.user\.*
        server-id: 5401-5499
      
      sink:
        type: starrocks
        name: StarRocks Sink
        jdbc-url: jdbc:mysql://<yourFeHostname>:9030
        load-url: <yourFeHostname>:8030
        username: <yourUsername>
        password: ${secret_values.starrocks_password}
        # (Opsional) Untuk pekerjaan dengan data sedikit, kurangi interval flush untuk mencegah penundaan persistensi lama. Default: 300000, atau 5 menit.
        sink.buffer-flush.interval-ms: 5000
        # (Opsional) Untuk tabel gabungan, tentukan secara eksplisit jumlah bucket berdasarkan volume data total.
        table.create.num-buckets: 8
      
      route:
        # Gabungkan semua tabel terpartisi ke tabel StarRocks test_db.user.
        - source-table: test_db.user\.*
          sink-table: test_db.user
      
      pipeline:
        name: MySQL to StarRocks Pipeline

Aktifkan mode EVOLVE

Secara default, mode LENIENT tidak menyinkronkan perubahan skema seperti menghapus kolom, menghapus tabel, atau memotong tabel ke tujuan. Jika Anda memerlukan sinkronisasi skema yang ketat, aktifkan mode EVOLVE. Mode ini memiliki batasan signifikan, jadi tinjau catatan berikut sebelum menggunakannya.

Batasan dan catatan penggunaan

  • Penggantian nama kolom tidak didukung. Pekerjaan gagal jika terjadi event penggantian nama kolom hulu.

  • Menghapus kolom atau tabel, atau memotong tabel, diterapkan ke tujuan. Operasi hulu yang tidak disengaja secara langsung memengaruhi tabel tujuan. Mode LENIENT default lebih aman karena tidak menyinkronkan event drop-table atau truncate-table.

  • Jika Anda me-restart pekerjaan tanpa state dan tidak menghapus tabel sink, ketidakcocokan skema hulu dan sink dapat menyebabkan pekerjaan gagal. Anda harus menyesuaikan skema tabel hilir secara manual.

source:
        type: mysql
        name: MySQL Source
        hostname: <yourHostname>
        port: 3306
        username: <yourUsername>
        password: ${secret_values.mysql_password}
        tables: test_db.test_source_table
        server-id: 5401-5499
      
      sink:
        type: starrocks
        name: StarRocks Sink
        jdbc-url: jdbc:mysql://<yourFeHostname>:9030
        load-url: <yourFeHostname>:8030
        username: <yourUsername>
        password: ${secret_values.starrocks_password}
      
      pipeline:
        name: MySQL to StarRocks Pipeline
        # Aktifkan mode EVOLVE untuk menyinkronkan perubahan skema secara ketat. Pekerjaan gagal pada perubahan yang tidak didukung, seperti penggantian nama kolom.
        schema.change.behavior: evolve