All Products
Search
Document Center

Realtime Compute for Apache Flink:Konektor DataHub

Last Updated:Jul 23, 2026

Konektor DataHub memungkinkan Anda membaca data streaming dari Alibaba Cloud DataHub ke dalam pekerjaan Flink dan menulis hasil pemrosesan kembali ke topik DataHub. Konektor ini mendukung baik Flink SQL maupun DataStream API.

Catatan DataHub kompatibel dengan protokol Kafka. Untuk menghubungkan Flink ke DataHub menggunakan protokol Kafka, gunakan konektor Kafka standar — bukan konektor Upsert Kafka. Untuk detailnya, lihat Kompatibilitas dengan Kafka.

Kemampuan

Item Deskripsi
Jenis yang didukung Source dan sink
Mode menjalankan Streaming dan batch
Format data N/A
Metrik N/A
Jenis API DataStream dan SQL
Dukungan pembaruan/penghapusan data di sink Tidak didukung. Sink hanya menulis baris insert-only ke topik target.

Prasyarat

Sebelum memulai, pastikan Anda telah memiliki:

Batasan

  • Membaca tabel sumber DataHub dengan pekerjaan batch tidak disarankan. Dalam mode batch, tabel sumber DataHub tidak dapat mencapai status selesai dan pekerjaan terus menunggu alih-alih berhenti.

Sintaksis

CREATE TEMPORARY TABLE datahub_input (
  `time` BIGINT,
  `sequence`  STRING METADATA VIRTUAL,
  `shard-id` BIGINT METADATA VIRTUAL,
  `system-time` TIMESTAMP METADATA VIRTUAL
) WITH (
  'connector' = 'datahub',
  'subId' = '<yourSubId>',
  'endPoint' = '<yourEndPoint>',
  'project' = '<yourProjectName>',
  'topic' = '<yourTopicName>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}'
);

Opsi konektor

Umum

Opsi Tipe Wajib Bawaan Deskripsi
connector String Ya (none) Jenis konektor. Atur nilai ini ke datahub.
endPoint String Ya (none) Titik akhir proyek DataHub. Nilainya bervariasi berdasarkan wilayah. Lihat Endpoints.
project String Ya (none) Nama proyek DataHub.
topic String Ya (none) Nama topik DataHub. Untuk topik BLOB (data tidak terstruktur dan tidak bertipe), tabel Flink harus berisi tepat satu kolom VARBINARY.
accessId String Ya (none) ID AccessKey Akun Alibaba Cloud Anda. Simpan sebagai variabel alih-alih hardcoding. Lihat Mengelola variabel.
accessKey String Ya (none) Rahasia AccessKey Akun Alibaba Cloud Anda.
retryTimeout Integer Tidak 1800000 Timeout maksimum untuk upaya pengulangan, dalam milidetik.
retryInterval Integer Tidak 1000 Interval antar upaya pengulangan, dalam milidetik.
CompressType String Tidak lz4 Algoritma kompresi untuk baca dan tulis. Nilai yang valid: lz4, deflate, "" (dinonaktifkan). Memerlukan VVR 6.0.5 atau lebih baru.

Khusus source

Opsi Tipe Wajib Bawaan Deskripsi
subId String Ya (none) ID langganan DataHub.
maxFetchSize Integer Tidak 50 Jumlah catatan yang diambil per permintaan. Tingkatkan nilai ini untuk meningkatkan throughput baca.
maxBufferSize Integer Tidak 50 Jumlah maksimum catatan yang di-cache dari pembacaan asinkron. Tingkatkan nilai ini untuk meningkatkan throughput baca.
fetchLatestDelay Integer Tidak 500 Durasi tidur dalam milidetik saat tidak ada data tersedia. Kurangi nilai ini untuk mengurangi latensi baca pada topik dengan traffic rendah.
lengthCheck String Tidak NONE Aturan penanganan baris ketika jumlah bidang yang diurai tidak sesuai dengan jumlah kolom yang didefinisikan. Nilai yang valid: NONE, SKIP, EXCEPTION, PAD. Lihat Aturan validasi jumlah bidang.
columnErrorDebug Boolean Tidak false Menentukan apakah akan mengaktifkan logging debug untuk error parsing bidang. Atur ke true untuk mencetak log exception parsing.
startTime String Tidak (none) Timestamp untuk mulai mengonsumsi data. Format: yyyy-MM-dd hh:mm:ss.
endTime String Tidak (none) Timestamp untuk berhenti mengonsumsi data. Format: yyyy-MM-dd hh:mm:ss.
startTimeMs Long Tidak -1 Timestamp untuk mulai mengonsumsi data, dalam milidetik. Mengambil prioritas atas startTime. Lihat Posisi awal konsumsi.

Posisi awal konsumsi

Opsi startTimeMs mengontrol dari mana source mulai membaca:

  • -1 (bawaan): Dimulai dari offset paling akhir di topik. Jika tidak ada offset yang tersedia, fallback ke offset paling awal.

  • Timestamp tertentu: Dimulai dari catatan pertama pada atau setelah timestamp yang ditentukan.

Penting

Nilai bawaan -1 dapat menyebabkan kehilangan data. Jika pekerjaan Anda gagal sebelum checkpoint pertamanya, offset paling akhir di topik mungkin telah maju, dan catatan yang ditulis selama jendela tersebut dilewati. Atur startTimeMs secara eksplisit ke timestamp tertentu untuk mengontrol posisi awal.

Aturan validasi jumlah bidang

Opsi lengthCheck menentukan tindakan yang diambil ketika jumlah bidang yang diurai dalam suatu baris tidak sesuai dengan jumlah kolom yang didefinisikan:

Nilai Perilaku
NONE (bawaan) Jika bidang yang diurai > kolom yang didefinisikan: baca dari kiri ke kanan hingga jumlah kolom yang didefinisikan. Jika bidang yang diurai < kolom yang didefinisikan: lewati baris tersebut.
SKIP Lewati baris di mana jumlah bidang yang diurai berbeda dari jumlah kolom yang didefinisikan.
EXCEPTION Melempar exception ketika jumlah bidang yang diurai berbeda dari jumlah kolom yang didefinisikan.
PAD Baca dari kiri ke kanan. Jika bidang yang diurai > kolom yang didefinisikan: baca dari kiri ke kanan hingga jumlah kolom yang didefinisikan. Jika bidang yang diurai < kolom yang didefinisikan: isi bidang yang hilang dengan null.

Khusus sink

Opsi Tipe Wajib Bawaan Deskripsi
batchCount Integer Tidak 500 Jumlah maksimum baris per batch tulis.
batchSize Integer Tidak 512000 Ukuran maksimum batch tulis, dalam byte.
flushInterval Integer Tidak 5000 Interval flush, dalam milidetik.
hashFields String Tidak null Daftar nama kolom yang dipisahkan koma yang digunakan untuk merutekan baris ke shard. Baris dengan nilai yang sama pada kolom-kolom ini ditulis ke shard yang sama. Bawaan (null) menggunakan penulisan acak. Contoh: hashFields=a,b.
timeZone String Tidak (none) Zona waktu yang digunakan saat mengonversi bidang TIMESTAMP.
schemaVersion Integer Tidak -1 Versi skema dalam registri skema terdaftar.

Perilaku flush batch tulis

Batch tulis di-flush ke DataHub ketika salah satu kondisi berikut terpenuhi lebih dulu:

  • Jumlah baris yang dibuffer mencapai batchCount.

  • Ukuran total data yang dibuffer mencapai batchSize.

  • Waktu sejak flush terakhir melebihi flushInterval.

Menambah nilai batchCount, batchSize, atau flushInterval meningkatkan throughput tulis dengan biaya latensi yang lebih tinggi.

Pemetaan tipe data

Tipe Flink Tipe DataHub
TINYINT TINYINT
BOOLEAN BOOLEAN
INTEGER INTEGER
BIGINT BIGINT
BIGINT TIMESTAMP
FLOAT FLOAT
DOUBLE DOUBLE
DECIMAL DECIMAL
VARCHAR STRING
SMALLINT SMALLINT
VARBINARY BLOB

Metadata

Bidang metadata bersifat read-only (R). Deklarasikan sebagai METADATA VIRTUAL dalam definisi tabel source untuk menyertakannya dalam kueri tanpa menuliskannya kembali ke DataHub.

Catatan Bidang metadata hanya tersedia saat menggunakan VVR 3.0.1 atau lebih baru.
Kunci Tipe data Deskripsi R/W
shard-id BIGINT METADATA VIRTUAL ID shard dari catatan tersebut. R
sequence STRING METADATA VIRTUAL Nomor urut catatan dalam shard. R
system-time TIMESTAMP METADATA VIRTUAL Waktu DataHub menerima catatan tersebut. R

Contoh

Source

Contoh berikut membaca data dari topik DataHub dan mencetaknya ke konsol.

CREATE TEMPORARY TABLE datahub_input (
  `time` BIGINT,
  `sequence`  STRING METADATA VIRTUAL,
  `shard-id` BIGINT METADATA VIRTUAL,
  `system-time` TIMESTAMP METADATA VIRTUAL
) WITH (
  'connector' = 'datahub',
  'subId' = '<yourSubId>',
  'endPoint' = '<yourEndPoint>',
  'project' = '<yourProjectName>',
  'topic' = '<yourTopicName>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}'
);

CREATE TEMPORARY TABLE test_out (
  `time` BIGINT,
  `sequence`  STRING,
  `shard-id` BIGINT,
  `system-time` TIMESTAMP
) WITH (
  'connector' = 'print',
  'logger' = 'true'
);

INSERT INTO test_out
SELECT
  `time`,
  `sequence`,
  `shard-id`,
  `system-time`
FROM datahub_input;

Sink

Contoh berikut membaca dari satu topik DataHub, mengonversi bidang name menjadi huruf kecil, dan menulis hasilnya ke topik DataHub lain.

CREATE TEMPORARY TABLE datahub_source (
  name VARCHAR
) WITH (
  'connector' = 'datahub',
  'endPoint' = '<endPoint>',
  'project' = '<yourProjectName>',
  'topic' = '<yourTopicName>',
  'subId' = '<yourSubId>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}',
  'startTime' = '2018-06-01 00:00:00'
);

CREATE TEMPORARY TABLE datahub_sink (
  name VARCHAR
) WITH (
  'connector' = 'datahub',
  'endPoint' = '<endPoint>',
  'project' = '<yourProjectName>',
  'topic' = '<yourTopicName>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}',
  'batchSize' = '512000',
  'batchCount' = '500'
);

INSERT INTO datahub_sink
SELECT
  LOWER(name)
FROM datahub_source;

DataStream API

Penting

Untuk menggunakan DataStream API dengan DataHub, konfigurasikan konektor DataStream untuk Realtime Compute for Apache Flink. Lihat Pengaturan konektor DataStream.

Baca dari DataHub

VVR menyediakan kelas DatahubSourceFunction, yang mengimplementasikan antarmuka SourceFunction milik Flink.

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1);

// Konfigurasikan source DataHub
DatahubSourceFunction datahubSource =
    new DatahubSourceFunction(
        <yourEndPoint>,
        <yourProjectName>,
        <yourTopicName>,
        <yourSubId>,
        <yourAccessId>,
        <yourAccessKey>,
        "public",
        <yourStartTime>,
        <yourEndTime>
    );
datahubSource.setRequestTimeout(30 * 1000);
datahubSource.enableExitAfterReadFinished();

env.addSource(datahubSource)
    .map((MapFunction<RecordEntry, Tuple2<String, Long>>) this::getStringLongTuple2)
    .print();
env.execute();

private Tuple2<String, Long> getStringLongTuple2(RecordEntry recordEntry) {
    Tuple2<String, Long> tuple2 = new Tuple2<>();
    TupleRecordData recordData = (TupleRecordData) (recordEntry.getRecordData());
    tuple2.f0 = (String) recordData.getField(0);
    tuple2.f1 = (Long) recordData.getField(1);
    return tuple2;
}

Tulis ke DataHub

VVR menyediakan kelas OutputFormatSinkFunction, yang mengimplementasikan antarmuka DatahubSinkFunction.

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

// Konfigurasikan sink DataHub
env.generateSequence(0, 100)
    .map((MapFunction<Long, RecordEntry>) aLong -> getRecordEntry(aLong, "default:"))
    .addSink(
        new DatahubSinkFunction<>(
            <yourEndPoint>,
            <yourProjectName>,
            <yourTopicName>,
            <yourSubId>,
            <yourAccessId>,
            <yourAccessKey>,
            "public",
            <schemaVersion> // Jika registri skema diaktifkan, Anda harus menentukan versi skema. Jika tidak, atur ke 0.
        )
    );
env.execute();

private RecordEntry getRecordEntry(Long message, String s) {
    RecordSchema recordSchema = new RecordSchema();
    recordSchema.addField(new Field("f1", FieldType.STRING));
    recordSchema.addField(new Field("f2", FieldType.BIGINT));
    recordSchema.addField(new Field("f3", FieldType.DOUBLE));
    recordSchema.addField(new Field("f4", FieldType.BOOLEAN));
    recordSchema.addField(new Field("f5", FieldType.TIMESTAMP));
    recordSchema.addField(new Field("f6", FieldType.DECIMAL));
    RecordEntry recordEntry = new RecordEntry();
    TupleRecordData recordData = new TupleRecordData(recordSchema);
    recordData.setField(0, s + message);
    recordData.setField(1, message);
    recordEntry.setRecordData(recordData);
    return recordEntry;
}

Dependensi Maven

Tambahkan konektor DataStream DataHub ke proyek Anda. Semua versi yang tersedia tercantum di repositori pusat Maven.

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

Lanjutan