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.
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:
-
Proyek dan topik DataHub. Lihat Memulai dengan DataHub.
-
Langganan DataHub (diperlukan untuk source). Lihat Membuat langganan.
-
ID AccessKey dan Rahasia AccessKey Akun Alibaba Cloud. Lihat Operasi Konsol.
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.
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.
| 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
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
-
Untuk daftar lengkap konektor yang didukung oleh Realtime Compute for Apache Flink, lihat Konektor yang didukung.
-
Untuk menghubungkan ke DataHub menggunakan konektor Kafka, lihat Message Queue for Apache Kafka.
-
Bagaimana cara melanjutkan penerapan yang gagal setelah topik DataHub di-split atau diskalakan?
-
Dapatkah saya menghapus topik DataHub yang sedang dikonsumsi?