Topik ini menjelaskan cara menggunakan konektor DataStream untuk membaca dari dan menulis ke sumber data.
Dependencies dan penggunaan konektor DataStream
Untuk membaca dari atau menulis ke sumber data menggunakan DataStream API, gunakan konektor DataStream yang sesuai untuk terhubung ke Realtime Compute for Apache Flink. Konektor DataStream VVR tersedia di repositori Maven Central dan dapat digunakan langsung dalam pengembangan pekerjaan Anda.
Gunakan hanya konektor yang secara eksplisit tercantum sebagai pendukung DataStream API dalam Konektor yang didukung. Konektor yang tidak ada dalam daftar tersebut tidak didukung, dan antarmukanya dapat berubah sewaktu-waktu tanpa pemberitahuan.
Konektor DataStream dienkripsi secara komersial dan tidak dapat dijalankan secara langsung. Untuk debugging dan eksekusi lokal, lihat Jalankan dan debug pekerjaan yang berisi konektor secara lokal.
Anda dapat menggunakan konektor dengan salah satu cara berikut:
(Direkomendasikan) Dependencies tambahan
-
Dalam file
pom.xmlMaven pekerjaan Anda, tambahkan konektor yang diperlukan sebagai dependency proyek dengan cakupanprovided.Catatan-
${vvr.version}menentukan versi engine untuk lingkungan runtime pekerjaan. Misalnya, jika pekerjaan Anda berjalan pada enginevvr-8.0.9-flink-1.17, versi Flink yang sesuai adalah1.17.2. Kami menyarankan Anda menggunakan engine terbaru. Untuk informasi lebih lanjut mengenai versi tertentu, lihat Versi engine. -
Karena paket JAR konektor diperkenalkan sebagai dependency tambahan, Anda tidak perlu memasukkan dependency ini ke dalam paket JAR. Oleh karena itu, Anda harus mendeklarasikan cakupannya sebagai
provided.
<!-- Dependency konektor MySQL --> <dependency> <groupId>com.alibaba.ververica</groupId> <artifactId>ververica-connector-mysql</artifactId> <version>${vvr.version}</version> <scope>provided</scope> </dependency> -
-
Jika Anda perlu mengembangkan konektor baru atau memperluas fungsionalitas konektor yang sudah ada, proyek Anda juga memerlukan dependency terhadap paket konektor umum
flink-connector-baseatauververica-connector-common.<!-- Dependency dasar untuk antarmuka publik konektor Flink --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-base</artifactId> <version>${flink.version}</version> </dependency> <!-- Dependency dasar untuk antarmuka publik konektor Alibaba Cloud --> <dependency> <groupId>com.alibaba.ververica</groupId> <artifactId>ververica-connector-common</artifactId> <version>${vvr.version}</version> </dependency> -
Terapkan pekerjaan JAR dan tambahkan paket JAR konektor yang sesuai ke bagian Additional Dependencies. Anda dapat mengunggah konektor yang Anda kembangkan sendiri atau konektor yang disediakan oleh Realtime Compute for Apache Flink.
Sebagai contoh, tambahkan
ververica-connector-mysql-1.17-vvr-8.0.9.jardanververica-connector-kafka-1.17-vvr-8.0.9.jar.
Dependency proyek
-
Dalam file
pom.xmlMaven pekerjaan Anda, tambahkan konektor yang Anda perlukan sebagai dependency proyek. Contoh berikut menunjukkan cara menambahkan konektor Kafka dan MySQL.Catatan-
${vvr.version}adalah versi engine lingkungan runtime pekerjaan. Misalnya, jika pekerjaan Anda berjalan pada enginevvr-8.0.9-flink-1.17, versi Flink yang sesuai adalah1.17.2. Kami menyarankan Anda menggunakan engine terbaru. Untuk informasi lebih lanjut, lihat Engine. -
Karena konektor dikemas langsung ke dalam JAR pekerjaan sebagai dependency proyek, mereka harus berada dalam cakupan default
compile.
<!-- Dependency konektor Kafka --> <dependency> <groupId>com.alibaba.ververica</groupId> <artifactId>ververica-connector-kafka</artifactId> <version>${vvr.version}</version> </dependency> <!-- Dependency konektor MySQL --> <dependency> <groupId>com.alibaba.ververica</groupId> <artifactId>ververica-connector-mysql</artifactId> <version>${vvr.version}</version> </dependency> -
-
Jika Anda perlu mengembangkan konektor baru atau memperluas fungsionalitas konektor yang sudah ada, proyek Anda juga memerlukan dependency terhadap paket konektor umum
flink-connector-baseatauververica-connector-common.<!-- Antarmuka publik konektor Flink --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-base</artifactId> <version>${flink.version}</version> </dependency> <!-- Dependency dasar untuk antarmuka publik konektor Alibaba Cloud --> <dependency> <groupId>com.alibaba.ververica</groupId> <artifactId>ververica-connector-common</artifactId> <version>${vvr.version}</version> </dependency>
Untuk mencegah konflik dependency, perhatikan hal berikut:
-
${flink.version}adalah versi Flink tempat pekerjaan dijalankan. Versi ini harus sama dengan versi Flink dari engine VVR yang Anda pilih di halaman penerapan pekerjaan. Misalnya, jika Anda memilih enginevvr-8.0.9-flink-1.17di halaman penerapan, versi Flink yang sesuai adalah1.17.2. Kami menyarankan Anda menggunakan engine terbaru. Untuk detail mengenai versi tertentu, lihat Versi Engine. -
Untuk dependency terkait Flink, atur cakupannya menjadi
provideddengan menambahkan<scope>provided</scope>ke dependency tersebut. Ini terutama mencakup dependency non-Konektor dalam gruporg.apache.flinkyang diawali denganflink-. -
Dalam kode sumber Apache Flink, panggil hanya metode yang dianotasi dengan @Public atau @PublicEvolving. Realtime Compute for Apache Flink hanya menjamin kompatibilitas untuk API publik ini.
-
Jika DataStream API didukung oleh konektor bawaan Flink, kami menyarankan Anda menggunakan dependency bawaan yang sesuai.
Dokumentasi terkait
-
Untuk contoh referensi pengembangan lengkap, lihat Kembangkan pekerjaan JAR.
-
Untuk konektor lain yang mendukung DataStream, lihat Konektor yang didukung.
-
Karena konektor DataStream dienkripsi secara komersial, konektor tersebut tidak dapat dijalankan atau di-debug secara lokal tanpa konfigurasi khusus. Untuk petunjuknya, lihat Jalankan dan debug pekerjaan yang berisi konektor secara lokal.