All Products
Search
Document Center

Realtime Compute for Apache Flink:Konektor MySQL DataStream

Last Updated:Sep 01, 2026

Topik ini menjelaskan cara melakukan debugging dan menjalankan pekerjaan DataStream yang menggunakan konektor MySQL.

MySQL CDC DataStream API

Penting

Untuk membaca atau menulis data menggunakan DataStream API, Anda harus menggunakan konektor DataStream yang sesuai. Untuk informasi lebih lanjut tentang cara mengonfigurasi konektor DataStream, lihat DataStream connector setup.

Buat program DataStream API dan gunakan MySqlSource. Contoh berikut menunjukkan kode dan dependensi pom:

import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import com.ververica.cdc.debezium.JsonDebeziumDeserializationSchema;
import com.ververica.cdc.connectors.mysql.source.MySqlSource;

public class MySqlSourceExample {
  public static void main(String[] args) throws Exception {
    MySqlSource<String> mySqlSource = MySqlSource.<String>builder()
        .hostname("yourHostname")
        .port(yourPort)
        .databaseList("yourDatabaseName") // set captured database
        .tableList("yourDatabaseName.yourTableName") // set captured table
        .username("yourUsername")
        .password("yourPassword")
        .deserializer(new JsonDebeziumDeserializationSchema()) // converts SourceRecord to JSON String
        .build();
    StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
    // enable checkpoint
    env.enableCheckpointing(3000);
    env
      .fromSource(mySqlSource, WatermarkStrategy.noWatermarks(), "MySQL Source")
      // set 4 parallel source tasks
      .setParallelism(4)
      .print().setParallelism(1); // use parallelism 1 for sink to keep message ordering
    env.execute("Print MySQL Snapshot + Binlog");
  }
}
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-core</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-connector-base</artifactId>
    <version>${flink.version}</version>
    <scope>provided</scope>
</dependency>
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-table-common</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>
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-table-api-java-bridge</artifactId>
    <version>${flink.version}</version>
    <scope>provided</scope>
</dependency>
<dependency>
    <groupId>com.alibaba.ververica</groupId>
    <artifactId>ververica-connector-mysql</artifactId>
    <version>${vvr.version}</version>
</dependency>

Saat membuat MySqlSource, Anda harus menentukan parameter berikut dalam kode:

Parameter

Deskripsi

hostname

Alamat IP atau hostname dari database MySQL.

port

Nomor port layanan database MySQL.

databaseList

Nama database MySQL.

Catatan

Nama database mendukung ekspresi reguler untuk membaca data dari beberapa database. Anda dapat menggunakan .* untuk mencocokkan semua database.

username

Username layanan database MySQL.

password

Password layanan database MySQL.

deserializer

Deserializer yang mengonversi catatan bertipe SourceRecord ke tipe tertentu. Nilai yang valid:

  • RowDataDebeziumDeserializeSchema: mengonversi SourceRecord ke RowData, struktur data internal Flink Table atau SQL.

  • JsonDebeziumDeserializationSchema: mengonversi SourceRecord ke string JSON.

Anda juga harus menentukan parameter berikut dalam dependensi pom:

${vvr.version}

Versi engine Realtime Compute for Apache Flink. Contoh: 1.17-vvr-8.0.4-3.

Catatan

Gunakan nomor versi yang ditampilkan di Maven sebagai acuan resmi, karena versi hotfix dirilis secara berkala dan pembaruan tersebut mungkin tidak diumumkan melalui saluran lain.

${flink.version}

Versi Apache Flink. Contoh: 1.17.2.

Penting

Gunakan versi Apache Flink yang sesuai dengan versi engine Realtime Compute for Apache Flink untuk menghindari masalah ketidakcocokan saat waktu proses. Untuk informasi lebih lanjut tentang pemetaan versi, lihat Engine.

Solusi debugging DataStream

Buat program DataStream API dan gunakan MySqlSource. Contoh berikut menunjukkan kode:

Penting

Untuk debugging lokal, Anda harus mengunduh file JAR yang diperlukan dan mengonfigurasi dependensinya. Untuk informasi lebih lanjut, lihat Run and debug jobs that contain connectors locally. Topik ini menyediakan contoh proyek Maven. Untuk informasi lebih lanjut, lihat MysqlCDCDemo.zip.

package com.alibaba.realtimecompute;

import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import com.ververica.cdc.debezium.JsonDebeziumDeserializationSchema;
import com.ververica.cdc.connectors.mysql.source.MySqlSource;
import org.apache.flink.configuration.Configuration;

public class MysqlCDCDemo {

    public static void main(String[] args) throws Exception {
        Configuration conf = new Configuration();
        conf.setString("pipeline.classpaths", "file://" + "absolute path to the MySQL uber JAR");  // Configure the dependency
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(conf);
        MySqlSource<String> mySqlSource = MySqlSource.<String>builder()
                .hostname("hostname")
                .port(3306)
                .databaseList("test_db") // Set the database to capture
                .tableList("test_db.test_table") // Set the table to capture
                .username("username")
                .password("password")
                .deserializer(new JsonDebeziumDeserializationSchema())
                .build();
        env.enableCheckpointing(3000);
        env.fromSource(mySqlSource, WatermarkStrategy.noWatermarks(), "MySQL Source")
                .setParallelism(4)
                .print().setParallelism(1);
        env.execute("Print MySQL Snapshot + Binlog");
    }
}

Saat membuat MySqlSource, Anda harus menentukan parameter berikut dalam kode:

Parameter

Deskripsi

hostname

Alamat IP atau hostname dari database MySQL.

port

Nomor port layanan database MySQL.

databaseList

Nama database MySQL.

Catatan

Nama database mendukung ekspresi reguler untuk membaca data dari beberapa database. Anda dapat menggunakan .* untuk mencocokkan semua database.

username

Username layanan database MySQL.

password

Password layanan database MySQL.

deserializer

Deserializer yang mengonversi catatan bertipe SourceRecord ke tipe tertentu. Nilai yang valid:

  • RowDataDebeziumDeserializeSchema: mengonversi SourceRecord ke RowData, struktur data internal Flink Table atau SQL.

  • JsonDebeziumDeserializationSchema: mengonversi SourceRecord ke string JSON.

dependensi pom

Debugging lokal

Konektor Realtime Compute for Apache Flink berisi konten komersial tambahan dan berbeda dari Apache Flink dalam banyak hal. Untuk debugging lokal, lakukan perubahan berikut pada file pom:

  1. ${flink.version} harus bernilai 1.19.0.

  2. Versi konektor ${vvr.version} yang direkomendasikan:

    1. Untuk engine VVR 8.x, gunakan versi 1.17-vvr-8.0.11-4.

    2. Untuk engine VVR 11.x, gunakan versi konektor terbaru yang sesuai dengan versi engine Anda. Misalnya, untuk engine VVR 11.8, gunakan versi terbaru 1.20-vvr-11.8.0-1-jdk11. Untuk versi lainnya, lihat Maven Repository.

  3. Anda harus menambahkan dependensi konektor Kafka.

<properties>
    <maven.compiler.source>8</maven.compiler.source>
    <maven.compiler.target>8</maven.compiler.target>
    <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
    <java.version>8</java.version>
    <flink.version>1.19.0</flink.version>
    <vvr.version>1.20-vvr-11.8.0-1-jdk11</vvr.version>
</properties>
<dependencies>
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-core</artifactId>
        <version>${flink.version}</version>
    </dependency>
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-streaming-java</artifactId>
        <version>${flink.version}</version>
    </dependency>
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-clients</artifactId>
        <version>${flink.version}</version>
    </dependency>
    <dependency>
        <!-- Use this group ID for VVR 11.x -->
        <groupId>com.alibaba.ververica</groupId>

       <!-- Use this group ID for VVR 8.x -->
        <!-- <groupId>org.apache.flink</groupId> -->
        <artifactId>flink-table-common</artifactId>
        <version>${flink.version}</version>
    </dependency>
    <dependency>
        <groupId>com.alibaba.ververica</groupId>
        <artifactId>ververica-connector-mysql</artifactId>
        <version>${vvr.version}</version>
    </dependency>
    <dependency>
        <groupId>com.alibaba.ververica</groupId>
        <artifactId>ververica-connector-kafka</artifactId>
        <version>${vvr.version}</version>
    </dependency>
</dependencies>

Penyebaran dan debugging pada Realtime Compute for Apache Flink

Saat melakukan debugging pekerjaan pada Realtime Compute for Apache Flink, tidak ada batasan versi konektor. Pastikan bahwa ${flink.version} dan ${vvr.version} sesuai dengan versi engine pekerjaan. Untuk informasi lebih lanjut tentang pemetaannya, lihat Engine. File pom berikut disediakan sebagai referensi:

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-core</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-connector-base</artifactId>
    <version>${flink.version}</version>
    <scope>provided</scope>
</dependency>
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-table-common</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>
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-table-api-java-bridge</artifactId>
    <version>${flink.version}</version>
    <scope>provided</scope>
</dependency>
<dependency>
    <groupId>com.alibaba.ververica</groupId>
    <artifactId>ververica-connector-mysql</artifactId>
    <version>${vvr.version}</version>
</dependency>