API Flink DataStream menyediakan model pemrograman fleksibel untuk membuat transformasi data, operasi, dan operator kustom guna menangani logika bisnis serta pemrosesan data yang kompleks.
Kompatibilitas Apache Flink
API DataStream yang didukung oleh Realtime Compute for Apache Flink sepenuhnya kompatibel dengan versi open source Apache Flink. Untuk informasi selengkapnya, lihat Apa itu Apache Flink? dan Panduan Pemrograman API Flink DataStream.
Prasyarat
-
Lingkungan pengembangan terintegrasi (IDE), seperti IntelliJ IDEA, telah diinstal.
-
Maven 3.6.3 atau versi yang lebih baru telah diinstal.
-
Pengembangan pekerjaan hanya mendukung JDK 8 dan JDK 11.
-
Kembangkan pekerjaan JAR Anda secara lokal sebelum menerapkan dan menjalankannya di Konsol Realtime Compute for Apache Flink.
Sebelum memulai
Siapkan sumber data yang diperlukan terlebih dahulu karena contoh ini menggunakan konektor sumber data.
-
Contoh ini menggunakan ApsaraMQ for Kafka (2.6.2) dan ApsaraDB RDS for MySQL (8.0) sebagai sumber data.
-
Jika Anda memiliki sumber data yang dikelola sendiri yang memerlukan akses jaringan publik atau akses lintas-VPC, lihat Opsi konektivitas jaringan.
-
Jika Anda belum memiliki sumber data ApsaraMQ for Kafka, beli dan terapkan instans tersebut. Untuk informasi selengkapnya, lihat Langkah 2: Membeli dan menerapkan instans. Saat menerapkan instans, pastikan instans tersebut berada dalam VPC yang sama dengan ruang kerja Realtime Compute for Apache Flink Anda.
-
Jika Anda belum memiliki sumber data ApsaraDB RDS for MySQL, beli instans ApsaraDB RDS for MySQL. Untuk informasi selengkapnya, lihat Langkah 1: Membuat instans ApsaraDB RDS for MySQL dan mengonfigurasi database. Saat membeli instans, pastikan instans tersebut berada dalam wilayah dan VPC yang sama dengan ruang kerja Realtime Compute for Apache Flink Anda.
Mengembangkan pekerjaan
Konfigurasi dependensi lingkungan Flink
Untuk menghindari konflik dependensi JAR, ikuti panduan berikut:
-
${flink.version}menentukan versi Flink untuk eksekusi pekerjaan. Versi ini harus konsisten dengan versi Flink mesin VVR yang Anda pilih pada halaman penerapan pekerjaan. Misalnya, jika Anda memilih mesinvvr-8.0.9-flink-1.17pada halaman penerapan pekerjaan, versi Flink yang sesuai adalah1.17.2. Untuk informasi selengkapnya tentang versi mesin VVR, lihat Bagaimana cara memeriksa versi Flink pekerjaan saat ini?. -
Untuk dependensi Flink, atur cakupan ke
provideddengan menambahkan<scope>provided</scope>. Ini terutama mencakup dependensi non-Konektor dalam gruporg.apache.flinkyang diawali denganflink-. -
Dalam kode sumber Flink, hanya metode yang secara eksplisit dianotasi dengan @Public atau @PublicEvolving yang merupakan API publik. Realtime Compute for Apache Flink hanya menjamin kompatibilitas untuk metode-metode tersebut.
-
Jika Anda menggunakan API DataStream dari konektor Flink bawaan, gunakan dependensi yang disediakan.
Berikut adalah dependensi dasar Flink. Anda mungkin juga perlu menambahkan dependensi logging. Untuk daftar lengkap dependensi, lihat Contoh kode lengkap di akhir topik ini.
Dependensi Flink
Dependensi dan penggunaan konektor
Untuk membaca dan menulis data dengan API DataStream, gunakan konektor DataStream untuk terhubung ke Realtime Compute for Apache Flink. Repositori Maven Central menyediakan konektor DataStream VVR yang dapat Anda gunakan selama pengembangan.
Gunakan hanya konektor yang secara eksplisit ditentukan mendukung API DataStream dalam Konektor yang didukung. Jangan gunakan konektor jika tidak secara eksplisit ditandai mendukung API DataStream, karena antarmuka dan parameternya mungkin berubah di masa depan.
Anda dapat menggunakan konektor dengan salah satu cara berikut:
(Direkomendasikan) Unggah sebagai dependensi tambahan
-
Dalam file pom.xml pekerjaan Anda, tambahkan konektor yang diperlukan sebagai dependensi proyek dengan cakupan "provided". Untuk file dependensi lengkap, lihat Contoh kode lengkap di akhir topik ini.
Catatan-
${vvr.version}adalah versi mesin runtime pekerjaan. Misalnya, jika pekerjaan Anda berjalan pada mesinvvr-8.0.9-flink-1.17, versi Flink yang sesuai adalah1.17.2. Kami merekomendasikan agar Anda menggunakan mesin terbaru. Untuk informasi tentang versi spesifik, lihat Mesin. -
Karena paket JAR konektor ditambahkan sebagai dependensi tambahan, paket tersebut tidak perlu dimasukkan ke dalam JAR aplikasi. Oleh karena itu, Anda harus mendeklarasikan cakupannya sebagai
provided.
<!-- Dependensi konektor Kafka --> <dependency> <groupId>com.alibaba.ververica</groupId> <artifactId>ververica-connector-kafka</artifactId> <version>${vvr.version}</version> <scope>provided</scope> </dependency> <!-- Dependensi 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 ada, proyek Anda juga perlu bergantung pada paket konektor umum
flink-connector-baseatauververica-connector-common.<!-- Dependensi dasar untuk antarmuka publik konektor Flink --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-base</artifactId> <version>${flink.version}</version> </dependency> <!-- Dependensi dasar untuk antarmuka publik konektor Alibaba Cloud --> <dependency> <groupId>com.alibaba.ververica</groupId> <artifactId>ververica-connector-common</artifactId> <version>${vvr.version}</version> </dependency> -
Untuk konfigurasi koneksi DataStream dan contoh kode, lihat dokumentasi untuk konektor DataStream yang sesuai.
Untuk daftar konektor yang mendukung API DataStream, lihat Konektor yang didukung.
-
Terapkan pekerjaan dan tambahkan paket JAR konektor yang sesuai di bidang Additional Dependencies. Untuk informasi selengkapnya, lihat Menerapkan pekerjaan JAR. Anda dapat mengunggah konektor kustom atau konektor yang disediakan oleh Realtime Compute for Apache Flink. Untuk mengunduh konektor, lihat Konektor.
Sebagai contoh, unggah paket JAR konektor
ververica-connector-mysql-1.17-vvr-8.0.4-1.jardanververica-connector-kafka-1.17-vvr-8.0.4-1.jar, serta file konfigurasiconfig.properties.
Masukkan ke dalam JAR pekerjaan
-
Tambahkan konektor yang diperlukan sebagai dependensi proyek dalam file pom.xml pekerjaan Anda. Kode berikut menunjukkan contoh untuk konektor Kafka dan MySQL.
Catatan-
${vvr.version}adalah versi mesin lingkungan runtime pekerjaan. Misalnya, jika pekerjaan Anda berjalan pada versi mesinvvr-8.0.9-flink-1.17, versi Flink yang sesuai adalah1.17.2. Kami merekomendasikan agar Anda menggunakan mesin terbaru. Untuk informasi selengkapnya, lihat Mesin. -
Jika Anda memasukkan konektor ke dalam JAR pekerjaan sebagai dependensi proyek, konektor tersebut harus berada dalam cakupan default (compile).
<!-- Dependensi konektor Kafka --> <dependency> <groupId>com.alibaba.ververica</groupId> <artifactId>ververica-connector-kafka</artifactId> <version>${vvr.version}</version> </dependency> <!-- Dependensi 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 ada, proyek Anda juga memerlukan paket konektor umum
flink-connector-baseatauververica-connector-common.<!-- Dependensi dasar untuk antarmuka publik konektor Flink --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-base</artifactId> <version>${flink.version}</version> </dependency> <!-- Dependensi dasar untuk antarmuka publik konektor Alibaba Cloud --> <dependency> <groupId>com.alibaba.ververica</groupId> <artifactId>ververica-connector-common</artifactId> <version>${vvr.version}</version> </dependency> -
Untuk konfigurasi koneksi DataStream dan contoh kode, lihat dokumentasi untuk konektor DataStream yang sesuai.
Untuk daftar konektor yang mendukung API DataStream, lihat Konektor yang didukung.
Membaca dependensi tambahan dari OSS
Pekerjaan JAR Flink tidak dapat membaca file konfigurasi lokal dari metode main. Sebagai gantinya, unggah file konfigurasi ke bucket OSS ruang kerja Anda, tambahkan sebagai dependensi tambahan selama penerapan, dan baca saat runtime. Bagian berikut memberikan contohnya.
-
Buat file konfigurasi bernama config.properties untuk menghindari hardcoding kredensial dalam kode Anda.
# Kafka bootstrapServers=host1:9092,host2:9092,host3:9092 inputTopic=topic groupId=groupId # MySQL database.url=jdbc:mysql://localhost:3306/my_database database.username=username database.password=password -
Dalam pekerjaan JAR Anda, gunakan kode untuk membaca file config.properties yang disimpan di bucket OSS.
Metode 1: Baca dari bucket ruang kerja
-
Di panel navigasi kiri Konsol Realtime Compute for Apache Flink, buka halaman Artifacts dan unggah file tersebut.
-
Saat runtime, file yang ditambahkan di bidang Additional Dependencies dimuat ke direktori /flink/usrlib pod tempat pekerjaan berjalan.
-
Kode berikut memberikan contoh cara membaca file konfigurasi ini.
Properties properties = new Properties(); Map<String,String> configMap = new HashMap<>(); try (InputStream input = new FileInputStream("/flink/usrlib/config.properties")) { // Muat file properti. properties.load(input); // Dapatkan nilai properti. configMap.put("bootstrapServers",properties.getProperty("bootstrapServers")) ; configMap.put("inputTopic",properties.getProperty("inputTopic")); configMap.put("groupId",properties.getProperty("groupId")); configMap.put("url",properties.getProperty("database.url")) ; configMap.put("username",properties.getProperty("database.username")); configMap.put("password",properties.getProperty("database.password")); } catch (IOException ex) { ex.printStackTrace(); }
Metode 2: Baca dari bucket yang diotorisasi
-
Unggah file konfigurasi ke bucket OSS target.
-
Gunakan OSSClient untuk membaca file langsung dari OSS. Untuk informasi selengkapnya, lihat Unduhan streaming dan Mengelola kredensial akses. Kode berikut memberikan contohnya.
OSS ossClient = new OSSClientBuilder().build("Endpoint", "AccessKeyId", "AccessKeySecret"); try (OSSObject ossObject = ossClient.getObject("examplebucket", "exampledir/config.properties"); BufferedReader reader = new BufferedReader(new InputStreamReader(ossObject.getObjectContent()))) { // baca file dan proses ... } finally { if (ossClient != null) { ossClient.shutdown(); } }
-
Tulis logika bisnis
-
Anda dapat mengintegrasikan sumber data eksternal ke dalam program aliran data Flink.
Watermarkadalah strategi komputasi Flink yang berbasis semantik waktu dan sering digunakan bersama timestamp. Oleh karena itu, contoh ini tidak menggunakan strategi watermark. Untuk informasi selengkapnya, lihat Strategi Watermark.// Integrasikan sumber data eksternal ke dalam program aliran data Flink. // WatermarkStrategy.noWatermarks() menunjukkan bahwa tidak ada strategi watermark yang digunakan. DataStreamSource<String> stream = env.fromSource(kafkaSource, WatermarkStrategy.noWatermarks(), "kafka Source"); -
Transformasi operator mengonversi
DataStream<String>menjadiDataStream<Student>dalam contoh ini. Untuk transformasi operator dan metode pemrosesan yang lebih kompleks, lihat Operator Flink.// Operator yang mengonversi struktur data menjadi Student. DataStream<Student> source = stream .map(new MapFunction<String, Student>() { @Override public Student map(String s) throws Exception { // Data dipisahkan oleh koma. String[] data = s.split(","); return new Student(Integer.parseInt(data[0]), data[1], Integer.parseInt(data[2])); } }).filter(student -> student.score >=60); // Filter data dengan skor 60 atau lebih tinggi.
Mengemas pekerjaan
Kemas pekerjaan dengan menggunakan maven-shade-plugin.
-
Jika Anda menambahkan konektor sebagai dependensi tambahan, pastikan cakupan dependensi konektor diatur ke
providedsaat mengemas pekerjaan. -
Jika Anda memasukkan konektor ke dalam JAR pekerjaan, gunakan cakupan default (compile).
Menguji dan menerapkan pekerjaan
-
Secara default, Realtime Compute for Apache Flink tidak dapat mengakses internet publik, sehingga mencegah pengujian langsung di lingkungan lokal. Kami merekomendasikan agar Anda melakukan pengujian unit secara terpisah. Untuk informasi selengkapnya, lihat Menjalankan dan men-debug pekerjaan dengan konektor secara lokal.
-
Untuk menerapkan pekerjaan JAR, lihat Menerapkan pekerjaan JAR.
Catatan-
Saat penerapan, jika Anda menggunakan konektor sebagai dependensi tambahan, pastikan Anda mengunggah paket JAR terkait.
-
Jika Anda perlu membaca file konfigurasi, Anda juga harus mengunggahnya sebagai dependensi tambahan.
-
Contoh kode lengkap
Contoh ini memproses data dari sumber ApsaraMQ for Kafka dan menulis hasilnya ke sink ApsaraDB RDS for MySQL. Contoh ini hanya untuk referensi. Untuk panduan gaya kode dan kualitas lebih lanjut, lihat Panduan gaya kode dan kualitas.
Contoh ini menghilangkan konfigurasi untuk parameter runtime seperti checkpoint, TTL, dan strategi restart. Anda dapat mengonfigurasi pengaturan ini di halaman deployment details setelah pekerjaan diterapkan. Karena pengaturan dalam kode memiliki prioritas lebih tinggi, kami merekomendasikan mengonfigurasinya setelah penerapan untuk menyederhanakan pembaruan di masa depan. Untuk informasi selengkapnya, lihat Mengonfigurasi informasi penerapan pekerjaan.
FlinkDemo.java
package com.aliyun;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.connector.jdbc.JdbcConnectionOptions;
import org.apache.flink.connector.jdbc.JdbcExecutionOptions;
import org.apache.flink.connector.jdbc.JdbcSink;
import org.apache.flink.connector.jdbc.JdbcStatementBuilder;
import org.apache.flink.connector.kafka.source.KafkaSource;
import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
import org.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDeserializationSchema;
import org.apache.flink.kafka.shaded.org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import java.io.FileInputStream;
import java.io.IOException;
import java.io.InputStream;
import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.util.HashMap;
import java.util.Map;
import java.util.Properties;
public class FlinkDemo {
// Definisikan struktur data.
public static class Student {
public int id;
public String name;
public int score;
public Student(int id, String name, int score) {
this.id = id;
this.name = name;
this.score = score;
}
}
public static void main(String[] args) throws Exception {
// Buat lingkungan eksekusi Flink.
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
Properties properties = new Properties();
Map<String,String> configMap = new HashMap<>();
try (InputStream input = new FileInputStream("/flink/usrlib/config.properties")) {
// Muat file properti.
properties.load(input);
// Dapatkan nilai properti.
configMap.put("bootstrapServers",properties.getProperty("bootstrapServers")) ;
configMap.put("inputTopic",properties.getProperty("inputTopic"));
configMap.put("groupId",properties.getProperty("groupId"));
configMap.put("url",properties.getProperty("database.url")) ;
configMap.put("username",properties.getProperty("database.username"));
configMap.put("password",properties.getProperty("database.password"));
} catch (IOException ex) {
ex.printStackTrace();
}
// Bangun sumber Kafka
KafkaSource<String> kafkaSource = KafkaSource.<String>builder()
.setBootstrapServers(configMap.get("bootstrapServers"))
.setTopics(configMap.get("inputTopic"))
.setStartingOffsets(OffsetsInitializer.latest())
.setGroupId(configMap.get("groupId"))
.setDeserializer(KafkaRecordDeserializationSchema.valueOnly(StringDeserializer.class))
.build();
// Integrasikan sumber data eksternal ke dalam program aliran data Flink.
// WatermarkStrategy.noWatermarks() menunjukkan bahwa tidak ada strategi watermark yang digunakan.
DataStreamSource<String> stream = env.fromSource(kafkaSource, WatermarkStrategy.noWatermarks(), "kafka Source");
// Filter skor 60 atau lebih tinggi.
DataStream<Student> source = stream
.map(new MapFunction<String, Student>() {
@Override
public Student map(String s) throws Exception {
String[] data = s.split(",");
return new Student(Integer.parseInt(data[0]), data[1], Integer.parseInt(data[2]));
}
}).filter(Student -> Student.score >=60);
source.addSink(JdbcSink.sink("INSERT IGNORE INTO student (id, username, score) VALUES (?, ?, ?)",
new JdbcStatementBuilder<Student>() {
public void accept(PreparedStatement ps, Student data) {
try {
ps.setInt(1, data.id);
ps.setString(2, data.name);
ps.setInt(3, data.score);
} catch (SQLException e) {
throw new RuntimeException(e);
}
}
},
new JdbcExecutionOptions.Builder()
.withBatchSize(5) // Jumlah catatan per penulisan batch.
.withBatchIntervalMs(2000) // Interval batch dalam milidetik.
.build(),
new JdbcConnectionOptions.JdbcConnectionOptionsBuilder()
.withUrl(configMap.get("url"))
.withDriverName("com.mysql.cj.jdbc.Driver")
.withUsername(configMap.get("username"))
.withPassword(configMap.get("password"))
.build()
)).name("Sink MySQL");
env.execute("Flink Demo");
}
}
pom.xml
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/maven-v4_0_0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>com.aliyun</groupId>
<artifactId>FlinkDemo</artifactId>
<version>1.0-SNAPSHOT</version>
<name>FlinkDemo</name>
<packaging>jar</packaging>
<properties>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<flink.version>1.17.1</flink.version>
<vvr.version>1.17-vvr-8.0.4-1</vvr.version>
<target.java.version>1.8</target.java.version>
<maven.compiler.source>${target.java.version}</maven.compiler.source>
<maven.compiler.target>${target.java.version}</maven.compiler.target>
<log4j.version>2.14.1</log4j.version>
</properties>
<dependencies>
<!-- Dependensi Apache Flink -->
<!-- Dependensi ini diatur ke 'provided' karena tidak boleh dikemas ke dalam file JAR. -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-java</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-clients</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<!-- Tambahkan dependensi konektor di sini. Mereka harus berada dalam cakupan default (compile). -->
<dependency>
<groupId>com.alibaba.ververica</groupId>
<artifactId>ververica-connector-kafka</artifactId>
<version>${vvr.version}</version>
</dependency>
<dependency>
<groupId>com.alibaba.ververica</groupId>
<artifactId>ververica-connector-mysql</artifactId>
<version>${vvr.version}</version>
</dependency>
<!-- Tambahkan framework logging untuk menghasilkan output konsol saat runtime. -->
<!-- Secara default, dependensi ini dikecualikan dari JAR aplikasi. -->
<dependency>
<groupId>org.apache.logging.log4j</groupId>
<artifactId>log4j-slf4j-impl</artifactId>
<version>${log4j.version}</version>
<scope>runtime</scope>
</dependency>
<dependency>
<groupId>org.apache.logging.log4j</groupId>
<artifactId>log4j-api</artifactId>
<version>${log4j.version}</version>
<scope>runtime</scope>
</dependency>
<dependency>
<groupId>org.apache.logging.log4j</groupId>
<artifactId>log4j-core</artifactId>
<version>${log4j.version}</version>
<scope>runtime</scope>
</dependency>
</dependencies>
<build>
<plugins>
<!-- Java Compiler -->
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>
<version>3.11.0</version>
<configuration>
<source>${target.java.version}</source>
<target>${target.java.version}</target>
</configuration>
</plugin>
<!-- Kami menggunakan maven-shade-plugin untuk membuat fat JAR yang berisi semua dependensi yang diperlukan. -->
<!-- Ubah nilai <mainClass> jika titik masuk program Anda berubah. -->
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-shade-plugin</artifactId>
<version>3.2.0</version>
<executions>
<execution>
<phase>package</phase>
<goals>
<goal>shade</goal>
</goals>
<!-- Kecualikan dependensi yang tidak diperlukan. -->
<configuration>
<artifactSet>
<excludes>
<exclude>org.apache.flink:force-shading</exclude>
<exclude>com.google.code.findbugs:jsr305</exclude>
<exclude>org.slf4j:*</exclude>
<exclude>org.apache.logging.log4j:*</exclude>
</excludes>
</artifactSet>
<filters>
<filter>
<!-- Jangan salin tanda tangan dari direktori META-INF.
Jika tidak, pengecualian keamanan mungkin terjadi saat Anda menggunakan file JAR. -->
<artifact>*:*</artifact>
<excludes>
<exclude>META-INF/*.SF</exclude>
<exclude>META-INF/*.DSA</exclude>
<exclude>META-INF/*.RSA</exclude>
</excludes>
</filter>
</filters>
<transformers>
<transformer
implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer">
<mainClass>com.aliyun.FlinkDemo</mainClass>
</transformer>
</transformers>
</configuration>
</execution>
</executions>
</plugin>
</plugins>
</build>
</project>
Dokumentasi terkait
-
Untuk daftar konektor yang mendukung API DataStream, lihat Konektor yang didukung.
-
Untuk contoh lengkap alur kerja pengembangan pekerjaan JAR Flink dari awal hingga akhir, lihat Pekerjaan JAR Flink.
-
Realtime Compute for Apache Flink juga mendukung pekerjaan SQL dan Python. Untuk panduan pengembangan, lihat Ikhtisar pengembangan pekerjaan dan Mengembangkan pekerjaan Python.