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 pusat Maven dan dapat langsung digunakan dalam pengembangan Pekerjaan Anda.
Gunakan hanya konektor yang secara eksplisit tercantum sebagai pendukung DataStream API dalam Konektor yang didukung. Konektor yang tidak tercantum dalam daftar tersebut tidak didukung, dan antarmukanya dapat berubah tanpa pemberitahuan sebelumnya.
Konektor DataStream dienkripsi secara komersial dan tidak dapat dijalankan secara langsung. Untuk debugging dan eksekusi lokal, lihat Menjalankan dan mendebug Pekerjaan yang berisi konektor secara lokal.
Anda dapat menggunakan konektor dengan salah satu cara berikut:
(Direkomendasikan) Atur cakupan dependency konektor ke provided
-
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 merekomendasikan penggunaan 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, deklarasikan 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> -
Tambahkan paket JAR konektor.
Metode 1: Gunakan konektor bawaan platform (direkomendasikan)
-
Konektor bawaan lebih disarankan karena mengurangi dependency tambahan dan pemeliharaan versi. Saat platform dipatch atau ditingkatkan, konektor bawaan diperbarui bersama platform sehingga tidak memerlukan pembaruan terpisah.
-
Hanya didukung di VVR 11.2 dan versi setelahnya.
Terapkan Pekerjaan JAR dan tambahkan konfigurasi ke bagian Other Configuration pada bagian Starting Parameters.
Untuk menggunakan beberapa konektor bawaan, misalnya mysql dan kafka, gunakan konfigurasi berikut. Untuk nama setiap konektor bawaan, lihat dokumentasi konektor tersebut di Konektor.
pipeline.used-builtin-connectors: mysql;kafka
Metode 2: Tambahkan paket JAR konektor sebagai dependency tambahan
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.
Misalnya, tambahkan ververica-connector-mysql-1.17-vvr-8.0.8.jar dan ververica-connector-kafka-1.17-vvr-8.0.8.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 merekomendasikan penggunaan engine terbaru. Untuk informasi lebih lanjut, lihat Engine. -
Karena konektor dikemas langsung ke dalam JAR Pekerjaan sebagai dependency proyek, konektor tersebut 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 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 merekomendasikan penggunaan engine terbaru. Untuk detail versi tertentu, lihat Versi Engine. -
Untuk dependency terkait Flink, atur cakupannya ke
provideddengan menambahkan<scope>provided</scope>pada 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 tersebut.
-
Jika DataStream API didukung oleh konektor Flink bawaan, kami merekomendasikan menggunakan dependency bawaan yang sesuai.
Dokumentasi terkait
-
Untuk contoh referensi pengembangan lengkap, lihat Mengembangkan Pekerjaan JAR.
-
Untuk konektor lain yang mendukung DataStream, lihat Konektor yang didukung.
-
Karena konektor DataStream dienkripsi secara komersial, konektor tersebut tidak dapat dijalankan atau didebug secara lokal tanpa konfigurasi khusus. Untuk petunjuknya, lihat Menjalankan dan mendebug Pekerjaan yang berisi konektor secara lokal.