Tous les produits
Search
Centre de documentation

ApsaraDB for SelectDB:Importer des données avec Flink

Dernière mise à jour :Aug 11, 2026

ApsaraDB for SelectDB est entièrement compatible avec Apache Doris. Utilisez le connecteur Flink Doris pour importer des données historiques depuis des sources telles que MySQL, Oracle, PostgreSQL, SQL Server et Kafka vers SelectDB. Une fois une tâche de capture de changement de données (CDC) lancée dans Flink, celle-ci synchronise également les données incrémentielles de la source vers SelectDB.

Vue d'ensemble

Remarque

Le connecteur Flink Doris prend actuellement en charge uniquement l'écriture de données vers SelectDB. Pour utiliser le Flink Doris Connector afin de vous connecter directement aux nœuds backend de SelectDB et lire des données efficacement, contactez l'équipe de support technique SelectDB pour demander un accès.

Vous pouvez également utiliser le connecteur Flink JDBC pour lire des données depuis SelectDB.

Le connecteur Flink Doris permet à Flink de lire et d'écrire dans Apache Doris pour le traitement et l'analyse de données en temps réel. Comme SelectDB est entièrement compatible avec Apache Doris, ce connecteur constitue une méthode courante pour ingérer des données en streaming dans SelectDB.

Chaque composant fonctionne comme suit :

  • Source

    • Objectif : une source lit les données provenant de systèmes externes pour alimenter un flux de données Flink. Ces systèmes peuvent inclure une file d'attente de messages (comme Apache Kafka), une base de données ou un système de fichiers.

    • Exemple : utilisez Kafka comme source pour lire des messages en temps réel ou lisez des données depuis un fichier.

  • Transformation

    • Objectif : l'étape de transformation traite le flux de données entrant. Ces opérations peuvent inclure le filtrage, le mappage, l'agrégation et le fenêtrage.

    • Exemple : mappez un flux d'entrée pour convertir sa structure de données ou agrégez des données pour calculer une métrique par minute.

  • Sink

    • Objectif : un sink écrit les données traitées d'un flux de données Flink vers un système externe, tel qu'une base de données, un fichier ou une file d'attente de messages.

    • Exemple : écrivez les résultats traités dans une base de données MySQL ou envoyez les données vers un autre topic Kafka.

La figure suivante illustre l'importation de données dans SelectDB à l'aide du connecteur Flink Doris.

image

Prérequis

  • Assurez-vous de la connectivité réseau entre votre source de données, Flink et SelectDB.

    1. Demandez un endpoint public pour votre instance ApsaraDB for SelectDB. Pour plus d'informations, consultez Demander ou libérer un endpoint public.

      Ignorez cette étape si votre environnement Flink et votre source de données se trouvent dans le même Virtual Private Cloud (VPC) que votre instance ApsaraDB for SelectDB. C'est généralement le cas lorsqu'il s'agit de produits Alibaba Cloud ou d'instances Elastic Compute Service (ECS) déployées dans le même VPC.

    2. Ajoutez les adresses IP de votre environnement Flink et de votre source de données à la liste d'autorisation de votre instance ApsaraDB for SelectDB. Pour plus d'informations, consultez Configurer une liste d'autorisation d'adresses IP.

  • Vérifiez que le connecteur Flink Doris est installé.

    Le tableau suivant liste les exigences de version pour Flink et le connecteur Flink Doris.

    Version de Flink

    Version du connecteur Flink Doris

    Lien de téléchargement

    Realtime Compute for Apache Flink : 1.17 ou ultérieure

    Flink open source : 1.15 ou ultérieure

    1.5.2 ou ultérieure. Nous recommandons de télécharger la dernière version.

    Connecteur Flink Doris

    Pour les instructions d'installation, consultez Installer le connecteur Flink Doris.

Ajouter le connecteur Flink Doris

Ajoutez le Flink Doris Connector en fonction de votre environnement.

  • Si vous utilisez Realtime Compute for Apache Flink pour importer des données dans SelectDB, gérez le Flink Doris Connector comme un connecteur personnalisé. Pour plus de détails, consultez Gérer les connecteurs personnalisés.

  • Si vous utilisez un cluster Flink auto-géré, téléchargez le package JAR du Flink Doris Connector correspondant et placez-le dans le répertoire lib de votre installation Flink. Pour le lien de téléchargement, consultez Package JAR.

  • Pour ajouter le Flink Doris Connector en tant que dépendance Maven, ajoutez le code suivant au fichier de configuration des dépendances de votre projet. Pour consulter d'autres versions, reportez-vous au Dépôt Maven.

    <!-- flink-doris-connector -->
    <dependency>
      <groupId>org.apache.doris</groupId>
      <artifactId>flink-doris-connector-1.16</artifactId>
      <version>1.5.2</version>
    </dependency>  

Exemples

Environnement d'exemple

Cet exemple utilise Flink SQL, Flink CDC et l'API DataStream pour migrer des données de la table employees dans la base de données test d'une instance ApsaraDB RDS for MySQL vers la table employees dans la base de données test d'une instance SelectDB. Modifiez les paramètres de ces exemples pour les adapter à votre scénario. L'environnement d'exemple est le suivant :

  • Environnement autonome Flink 1.16

  • Java

  • Base de données cible : test

  • Table cible : employees

  • Base de données source : test

  • Table source : employees

Préparer l'environnement

Environnement Flink

  1. Préparez un environnement Java.

    Flink nécessite un environnement Java pour fonctionner. Installez un Java Development Kit (JDK) et configurez la variable d'environnement JAVA_HOME.

    Pour la liste des versions Java prises en charge, consultez Compatibilité Java. Cet exemple utilise Java 8. Pour les instructions d'installation, consultez Installer le JDK.

  2. Téléchargez le package d'installation Flink flink-1.16.3-bin-scala_2.12.tgz. Si cette version est obsolète, téléchargez-en une autre depuis Apache Flink.

    wget https://www.apache.si/flink/flink-1.16.3/flink-1.16.3-bin-scala_2.12.tgz
  3. Décompressez le package d'installation.

    tar -zxvf flink-1.16.3-bin-scala_2.12.tgz
  4. Accédez au répertoire lib dans le dossier d'installation de Flink et ajoutez les connecteurs nécessaires aux étapes suivantes.

    • Ajoutez le connecteur Flink Doris.

      wget https://repo.maven.apache.org/maven2/org/apache/doris/flink-doris-connector-1.16/1.5.2/flink-doris-connector-1.16-1.5.2.jar
    • Ajoutez le connecteur Flink MySQL.

      wget https://repo1.maven.org/maven2/com/ververica/flink-sql-connector-mysql-cdc/2.4.2/flink-sql-connector-mysql-cdc-2.4.2.jar
  5. Démarrez le cluster Flink.

    Dans le répertoire bin de votre installation Flink, exécutez la commande suivante :

    ./start-cluster.sh 

SelectDB cible

  1. Créez une instance ApsaraDB for SelectDB. Pour plus d'informations, consultez Créer une instance.

  2. Connectez-vous à l'instance. Pour plus d'informations, consultez Se connecter à une instance.

  3. Créez une base de données de test nommée test.

    CREATE DATABASE test;
  4. Créez une table de test nommée employees.

    USE test;
    
    -- Create table
    CREATE TABLE employees (
        emp_no       int NOT NULL,
        birth_date   date,
        first_name   varchar(20),
        last_name    varchar(20),
        gender       char(2),
        hire_date    date
    )
    UNIQUE KEY(`emp_no`)
    DISTRIBUTED BY HASH(`emp_no`) BUCKETS 1;

MySQL source

  1. Créez une instance ApsaraDB RDS for MySQL.

  2. Créez une base de données de test nommée test.

    CREATE DATABASE test;
  3. Créez une table de test nommée employees.

    USE test;
    
    CREATE TABLE employees (
        emp_no INT NOT NULL PRIMARY KEY,
        birth_date DATE,
        first_name VARCHAR(20),
        last_name VARCHAR(20),
        gender CHAR(2),
        hire_date DATE
    );
  4. Insérez des données.

    INSERT INTO employees (emp_no, birth_date, first_name, last_name, gender, hire_date) VALUES
    (1001, '1985-05-15', 'John', 'Doe', 'M', '2010-06-20'),
    (1002, '1990-08-22', 'Jane', 'Smith', 'F', '2012-03-15'),
    (1003, '1987-11-02', 'Robert', 'Johnson', 'M', '2015-07-30'),
    (1004, '1992-01-18', 'Emily', 'Davis', 'F', '2018-01-05'),
    (1005, '1980-12-09', 'Michael', 'Brown', 'M', '2008-11-21');

Importation avec Flink SQL

  1. Démarrez le client Flink SQL.

    Dans le répertoire bin de votre installation Flink, exécutez la commande suivante :

    ./sql-client.sh
  2. Dans le client Flink SQL, soumettez un job Flink.

    1. Créez une table source MySQL.

      La clause WITH dans l'instruction suivante spécifie la configuration de la MySQL CDC Source. Pour plus d'informations sur les paramètres, consultez MySQL | Apache Flink CDC.

      CREATE TABLE employees_source (
          emp_no INT,
          birth_date DATE,
          first_name STRING,
          last_name STRING,
          gender STRING,
          hire_date DATE,
          PRIMARY KEY (`emp_no`) NOT ENFORCED
      ) WITH (
          'connector' = 'mysql-cdc',
          'hostname' = '127.0.0.1', 
          'port' = '3306',
          'username' = 'root',
          'password' = '****',
          'database-name' = 'test',
          'table-name' = 'employees'
      );
    2. Créez une table sink SelectDB.

      La clause WITH dans l'instruction suivante spécifie la configuration pour SelectDB. Pour plus d'informations sur les paramètres, consultez Paramètres du sink.

      CREATE TABLE employees_sink (
          emp_no       INT ,
          birth_date   DATE,
          first_name   STRING,
          last_name    STRING,
          gender       STRING,
          hire_date    DATE
      ) 
      WITH (
        'connector' = 'doris',
        'fenodes' = 'selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080',
        'table.identifier' = 'test.employees',
        'username' = 'admin',
        'password' = '****'
      );
    3. Synchronisez les données de la table source MySQL vers la table sink SelectDB.

      INSERT INTO employees_sink SELECT * FROM employees_source;
  3. Vérifiez l'importation des données.

    Connectez-vous à SelectDB et exécutez l'instruction suivante pour afficher les données importées.

    SELECT * FROM test.employees;

Importation avec Flink CDC

Important

Realtime Compute for Apache Flink ne prend pas en charge les jobs basés sur JAR. Utilisez plutôt des jobs basés sur YAML avec CDC 3.0.

Utilisez Flink CDC pour importer des données dans SelectDB.

Pour exécuter un job Flink CDC, utilisez le programme flink situé dans votre répertoire d'installation Flink. La syntaxe est la suivante :

<FLINK_HOME>/bin/flink run \
    -Dexecution.checkpointing.interval=10s \
    -Dparallelism.default=1 \
    -c org.apache.doris.flink.tools.cdc.CdcTools \
    lib/flink-doris-connector-1.16-1.5.2.jar \
    <mysql-sync-database|oracle-sync-database|postgres-sync-database|sqlserver-sync-database> \
    --database <selectdb-database-name> \
    [--job-name <flink-job-name>] \
    [--table-prefix <selectdb-table-prefix>] \
    [--table-suffix <selectdb-table-suffix>] \
    [--including-tables <mysql-table-name|name-regular-expr>] \
    [--excluding-tables <mysql-table-name|name-regular-expr>] \
    --mysql-conf <mysql-cdc-source-conf> [--mysql-conf <mysql-cdc-source-conf> ...] \
    --oracle-conf <oracle-cdc-source-conf> [--oracle-conf <oracle-cdc-source-conf> ...] \
    --sink-conf <doris-sink-conf> [--table-conf <doris-sink-conf> ...] \
    [--table-conf <selectdb-table-conf> [--table-conf <selectdb-table-conf> ...]]

Paramètres

Paramètre

Description

execution.checkpointing.interval

L'intervalle de checkpoint Flink. Ce paramètre affecte la fréquence de synchronisation des données. Une valeur de 10s est recommandée.

parallelism.default

Le parallélisme du job Flink. Augmenter le parallélisme peut améliorer la vitesse de synchronisation des données.

job-name

Le nom du job Flink.

database

Le nom de la base de données cible dans SelectDB.

table-prefix

Le préfixe du nom de la table cible dans SelectDB. Par exemple, --table-prefix ods_.

table-suffix

Le suffixe du nom de la table cible dans SelectDB.

including-tables

Les tables à synchroniser. Utilisez une barre verticale `

` pour séparer plusieurs tables. Les expressions régulières sont prises en charge. Par exemple, `--including-tables table1

tbl.*` synchronise table1 et toutes les tables commençant par tbl.

excluding-tables

Les tables à exclure de la synchronisation. Le format est identique à celui de including-tables.

mysql-conf

La configuration de la source MySQL CDC. Pour plus d'informations, consultez Connecteur MySQL CDC. Les paramètres hostname, username, password et database-name sont obligatoires.

oracle-conf

La configuration de la source Oracle CDC. Pour plus d'informations, consultez Connecteur Oracle CDC. Les paramètres hostname, username, password, database-name et schema-name sont obligatoires.

sink-conf

Paramètres de configuration pour le sink Doris. Pour plus d'informations, consultez Paramètres du sink.

table-conf

Paramètres de configuration pour la table SelectDB. Ils correspondent au contenu de la clause PROPERTIES lors de la création d'une table dans SelectDB.

Remarque
  1. Pour la synchronisation des données, ajoutez la dépendance Flink CDC requise, telle que flink-sql-connector-mysql-cdc-${version}.jar ou flink-sql-connector-oracle-cdc-${version}.jar, dans le répertoire $FLINK_HOME/lib.

  2. La synchronisation complète de base de données est prise en charge à partir de Flink 1.15. Pour télécharger différentes versions du connecteur Flink Doris, consultez Connecteur Flink Doris.

Paramètres du sink

Paramètre

Valeur par défaut

Obligatoire

Description

fenodes

Aucune

Oui

L'endpoint et le port HTTP de votre instance ApsaraDB for SelectDB.

Vous pouvez obtenir l'VPC Endpoint (ou le Public Endpoint) et le HTTP Port depuis la page Instance Details > Network Information dans la console ApsaraDB for SelectDB.

Exemple : selectdb-cn-4xl3jv1****.selectdbfe.rds.aliyuncs.com:8080.

table.identifier

Aucune

Oui

Le nom de la base de données et de la table. Exemple : test_db.test_table.

username

Aucune

Oui

Le nom d'utilisateur de la base de données pour votre instance ApsaraDB for SelectDB.

password

Aucune

Oui

Le mot de passe de l'utilisateur de la base de données pour votre instance ApsaraDB for SelectDB.

jdbc-url

Aucune

Non

Les informations de connexion JDBC pour votre instance ApsaraDB for SelectDB.

Vous pouvez obtenir l'VPC Endpoint (ou le Public Endpoint) et le MySQL Port depuis la page Instance Details > Network Information dans la console ApsaraDB for SelectDB.

Exemple : jdbc:mysql://selectdb-cn-4xl3jv1****.selectdbfe.rds.aliyuncs.com:9030.

auto-redirect

true

Non

Indique s'il faut rediriger les requêtes Stream Load. Si cette option est activée, Stream Load écrit les données via les frontends (FE) et les informations du backend (BE) ne sont pas récupérées.

doris.request.retries

3

Non

Le nombre de tentatives de renvoi d'une requête à SelectDB.

doris.request.connect.timeout

30s

Non

Le délai d'expiration pour la connexion à SelectDB.

doris.request.read.timeout

30s

Non

Le délai d'expiration pour la lecture des données depuis SelectDB.

sink.label-prefix

""

Oui

Le préfixe de libellé pour les importations Stream Load. Dans un scénario de validation en deux phases (2PC), ce préfixe doit être globalement unique pour garantir la sémantique exactly-once (EOS) de Flink.

sink.properties

Aucune

Non

Les paramètres d'importation pour Stream Load. Configurez les propriétés comme suit :

  • Pour le format CSV :

    sink.properties.format='csv' 
    sink.properties.column_separator=','
    sink.properties.line_delimiter='\n' 
  • Pour le format JSON :

    sink.properties.format='json' 

Pour plus de paramètres, consultez Stream Load.

sink.buffer-size

1048576

Non

La taille du tampon d'écriture, en octets. La valeur par défaut de 1 Mo est recommandée.

sink.buffer-count

3

Non

Le nombre de tampons d'écriture. La valeur par défaut est recommandée.

sink.max-retries

3

Non

Le nombre maximal de tentatives après un échec de validation. La valeur par défaut est 3.

sink.use-cache

false

Non

Indique s'il faut utiliser un cache en mémoire pour la reprise après exception. Si cette option est activée, les données de la période de checkpoint sont conservées dans le cache.

sink.enable-delete

true

Non

Indique s'il faut synchroniser les événements de suppression. Cette option n'est prise en charge que pour les tables utilisant le modèle Unique Key.

sink.enable-2pc

true

Non

Indique s'il faut activer la validation en deux phases (2PC). Cette option est activée par défaut (true) pour garantir la sémantique exactly-once (EOS).

sink.enable.batch-mode

false

Non

Indique s'il faut utiliser le mode batch pour écrire des données dans SelectDB. Lorsqu'il est activé, l'opération d'écriture est déclenchée par la taille du tampon ou par le temps, selon les paramètres sink.buffer-flush.max-rows, sink.buffer-flush.max-bytes et sink.buffer-flush.interval, plutôt que par les checkpoints Flink.

Lorsque le mode batch est activé, la sémantique exactly-once (EOS) n'est pas garantie. Vous pouvez utiliser le modèle Unique Key pour assurer l'idempotence.

sink.flush.queue-size

2

Non

La taille de la file d'attente du tampon en mode batch.

sink.buffer-flush.max-rows

50000

Non

Le nombre maximal de lignes par écriture par lot en mode batch.

sink.buffer-flush.max-bytes

10MB

Non

La taille maximale en octets par écriture par lot en mode batch.

sink.buffer-flush.interval

10s

Non

L'intervalle de vidage asynchrone du tampon en mode batch. La valeur minimale est de 1 seconde.

sink.ignore.update-before

true

Non

Indique s'il faut ignorer les événements update-before. Par défaut, ils sont ignorés.

Exemples de synchronisation

Synchronisation MySQL

<FLINK_HOME>/bin/flink run \
    -Dexecution.checkpointing.interval=10s \
    -Dparallelism.default=1 \
    -c org.apache.doris.flink.tools.cdc.CdcTools \
    lib/flink-doris-connector-1.16-1.5.2.jar \
    mysql-sync-database \
    --database test \
    --mysql-conf hostname=127.0.0.1 \
    --mysql-conf port=3306 \
    --mysql-conf username=root \
    --mysql-conf password="password" \
    --mysql-conf database-name=test \
    --including-tables "employees" \
    --sink-conf fenodes=selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080 \
    --sink-conf username=admin \
    --sink-conf password=****

Synchronisation Oracle

<FLINK_HOME>/bin/flink run \
    -Dexecution.checkpointing.interval=10s \
    -Dparallelism.default=1 \
    -c org.apache.doris.flink.tools.cdc.CdcTools \
    lib/flink-doris-connector-1.16-1.5.2.jar \
    oracle-sync-database \
    --database test_db \
    --oracle-conf hostname=127.0.0.1 \
    --oracle-conf port=1521 \
    --oracle-conf username=admin \
    --oracle-conf password="password" \
    --oracle-conf database-name=XE \
    --oracle-conf schema-name=ADMIN \
    --including-tables "tbl1|test.*" \
    --sink-conf fenodes=selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080 \
    --sink-conf username=admin \
    --sink-conf password=****

Synchronisation PostgreSQL

<FLINK_HOME>/bin/flink run \
    -Dexecution.checkpointing.interval=10s \
    -Dparallelism.default=1 \
    -c org.apache.doris.flink.tools.cdc.CdcTools \
    lib/flink-doris-connector-1.16-1.5.2.jar \
    postgres-sync-database \
    --database db1\
    --postgres-conf hostname=127.0.0.1 \
    --postgres-conf port=5432 \
    --postgres-conf username=postgres \
    --postgres-conf password="123456" \
    --postgres-conf database-name=postgres \
    --postgres-conf schema-name=public \
    --postgres-conf slot.name=test \
    --postgres-conf decoding.plugin.name=pgoutput \
    --including-tables "tbl1|test.*" \
    --sink-conf fenodes=selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080 \
    --sink-conf username=admin \
    --sink-conf password=****

Synchronisation SQL Server

<FLINK_HOME>/bin/flink run \
    -Dexecution.checkpointing.interval=10s \
    -Dparallelism.default=1 \
    -c org.apache.doris.flink.tools.cdc.CdcTools \
    lib/flink-doris-connector-1.16-1.5.2.jar \
    sqlserver-sync-database \
    --database db1\
    --sqlserver-conf hostname=127.0.0.1 \
    --sqlserver-conf port=1433 \
    --sqlserver-conf username=sa \
    --sqlserver-conf password="123456" \
    --sqlserver-conf database-name=CDC_DB \
    --sqlserver-conf schema-name=dbo \
    --including-tables "tbl1|test.*" \
    --sink-conf fenodes=selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080 \
    --sink-conf username=admin \
    --sink-conf password=****

Importation avec l'API DataStream

  1. Ajoutez les dépendances suivantes à votre projet Maven.

    Dépendances Maven

    <properties>
            <maven.compiler.source>8</maven.compiler.source>
            <maven.compiler.target>8</maven.compiler.target>
            <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
            <scala.version>2.12</scala.version>
            <java.version>1.8</java.version>
            <flink.version>1.16.3</flink.version>
            <fastjson.version>1.2.62</fastjson.version>
            <scope.mode>compile</scope.mode>
        </properties>
        <dependencies>
            <!-- https://mvnrepository.com/artifact/com.google.guava/guava -->
            <dependency>
                <groupId>com.google.guava</groupId>
                <artifactId>guava</artifactId>
                <version>28.1-jre</version>
            </dependency>
    
            <dependency>
                <groupId>org.apache.commons</groupId>
                <artifactId>commons-lang3</artifactId>
                <version>3.14.0</version>
            </dependency>
    
            <dependency>
                <groupId>org.apache.doris</groupId>
                <artifactId>flink-doris-connector-1.16</artifactId>
                <version>1.5.2</version>
            </dependency>
    
            <dependency>
                <groupId>org.apache.flink</groupId>
                <artifactId>flink-table-api-scala-bridge_${scala.version}</artifactId>
                <version>${flink.version}</version>
            </dependency>
            <dependency>
                <groupId>org.apache.flink</groupId>
                <artifactId>flink-table-planner_${scala.version}</artifactId>
                <version>${flink.version}</version>
            </dependency>
            <dependency>
                <groupId>org.apache.flink</groupId>
                <artifactId>flink-streaming-scala_${scala.version}</artifactId>
                <version>${flink.version}</version>
            </dependency>
            <dependency>
                <groupId>org.apache.flink</groupId>
                <artifactId>flink-clients</artifactId>
                <version>${flink.version}</version>
            </dependency>
            <dependency>
                <groupId>org.apache.flink</groupId>
                <artifactId>flink-connector-jdbc</artifactId>
                <version>${flink.version}</version>
            </dependency>
            <dependency>
                <groupId>org.apache.flink</groupId>
                <artifactId>flink-connector-kafka</artifactId>
                <version>${flink.version}</version>
            </dependency>
            <dependency>
                <groupId>org.apache.doris</groupId>
                <artifactId>flink-doris-connector-1.16</artifactId>
                <version>1.5.2</version>
            </dependency>
    
            <dependency>
                <groupId>com.ververica</groupId>
                <artifactId>flink-sql-connector-mysql-cdc</artifactId>
                <version>2.4.2</version>
                <exclusions>
                    <exclusion>
                        <artifactId>flink-shaded-guava</artifactId>
                        <groupId>org.apache.flink</groupId>
                    </exclusion>
                </exclusions>
            </dependency>
    
            <dependency>
                <groupId>org.apache.flink</groupId>
                <artifactId>flink-runtime-web</artifactId>
                <version>${flink.version}</version>
            </dependency>
    
        </dependencies>
  2. Code Java principal.

    Le code suivant configure la table source MySQL et la table sink ApsaraDB for SelectDB. Les paramètres correspondent à ceux utilisés dans la section Importer des données avec Flink SQL. Pour plus d'informations, consultez MySQL | Apache Flink CDC et Paramètres du sink.

    package org.example;
    
    import com.ververica.cdc.connectors.mysql.source.MySqlSource;
    import com.ververica.cdc.connectors.mysql.table.StartupOptions;
    import com.ververica.cdc.connectors.shaded.org.apache.kafka.connect.json.JsonConverterConfig;
    import com.ververica.cdc.debezium.JsonDebeziumDeserializationSchema;
    
    import org.apache.doris.flink.cfg.DorisExecutionOptions;
    import org.apache.doris.flink.cfg.DorisOptions;
    import org.apache.doris.flink.sink.DorisSink;
    import org.apache.doris.flink.sink.writer.serializer.JsonDebeziumSchemaSerializer;
    import org.apache.doris.flink.tools.cdc.mysql.DateToStringConverter;
    import org.apache.flink.api.common.eventtime.WatermarkStrategy;
    import org.apache.flink.streaming.api.datastream.DataStreamSource;
    import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
    
    import java.util.HashMap;
    import java.util.Map;
    import java.util.Properties;
    
    public class Main {
        public static void main(String[] args) throws Exception {
            StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
            env.setParallelism(1);
            env.enableCheckpointing(10000);
    
            Map<String, Object> customConverterConfigs = new HashMap<>();
            customConverterConfigs.put(JsonConverterConfig.DECIMAL_FORMAT_CONFIG, "numeric");
            JsonDebeziumDeserializationSchema schema =
                    new JsonDebeziumDeserializationSchema(false, customConverterConfigs);
            
            // Configure the MySQL source table
            MySqlSource<String> mySqlSource = MySqlSource.<String>builder()
                    .hostname("rm-xxx.mysql.rds.aliyuncs***")
                    .port(3306)
                    .startupOptions(StartupOptions.initial())
                    .databaseList("db_test")
                    .tableList("db_test.employees")
                    .username("root")
                    .password("test_123")
                    .debeziumProperties(DateToStringConverter.DEFAULT_PROPS)
                    .deserializer(schema)
                    .serverTimeZone("Asia/Shanghai")
                    .build();
    
            // Configure the ApsaraDB for SelectDB sink table
            DorisSink.Builder<String> sinkBuilder = DorisSink.builder();
            DorisOptions.Builder dorisBuilder = DorisOptions.builder();
            dorisBuilder.setFenodes("selectdb-cn-xxx-public.selectdbfe.rds.aliyunc****:8080")
                    .setTableIdentifier("db_test.employees")
                    .setUsername("admin")
                    .setPassword("test_123");
            DorisOptions dorisOptions = dorisBuilder.build();
    
            // Configure Stream Load parameters with sink.properties
            Properties properties = new Properties();
            properties.setProperty("format", "json");
            properties.setProperty("read_json_by_line", "true");
            DorisExecutionOptions.Builder executionBuilder = DorisExecutionOptions.builder();
            executionBuilder.setStreamLoadProp(properties);
    
            sinkBuilder.setDorisExecutionOptions(executionBuilder.build())
                    .setSerializer(JsonDebeziumSchemaSerializer.builder().setDorisOptions(dorisOptions).build()) // Serialize the data stream.
                    .setDorisOptions(dorisOptions);
    
            DataStreamSource<String> dataStreamSource = env.fromSource(mySqlSource, WatermarkStrategy.noWatermarks(), "MySQL Source");
            dataStreamSource.sinkTo(sinkBuilder.build());
            env.execute("MySQL to SelectDB");
        }
    }

Utilisation avancée

Mettre à jour des colonnes partielles avec Flink SQL

-- enable checkpoint
SET 'execution.checkpointing.interval' = '10s';

CREATE TABLE cdc_mysql_source (
   id INT
  ,name STRING
  ,bank STRING
  ,age INT
  ,PRIMARY KEY (id) NOT ENFORCED
) WITH (
 'connector' = 'mysql-cdc',
 'hostname' = '127.0.0.1',
 'port' = '3306',
 'username' = 'root',
 'password' = 'password',
 'database-name' = 'database',
 'table-name' = 'table'
);

CREATE TABLE selectdb_sink (
    id INT,
    name STRING,
    bank STRING,
    age INT
) 
WITH (
  'connector' = 'doris',
  'fenodes' = 'selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080',
  'table.identifier' = 'database.table',
  'username' = 'admin',
  'password' = '****',
  'sink.properties.format' = 'json',
  'sink.properties.read_json_by_line' = 'true',
  'sink.properties.columns' = 'id,name,bank,age',
  'sink.properties.partial_columns' = 'true' -- Enable partial column updates.
);

INSERT INTO selectdb_sink SELECT id,name,bank,age FROM cdc_mysql_source;

Utiliser Flink SQL pour supprimer des données par colonne

Dans les scénarios CDC, le sink Doris identifie le type d'événement à partir de RowKind et attribue une valeur à la colonne masquée __DORIS_DELETE_SIGN__ pour effectuer des suppressions. Lorsque la source de données est constituée de messages Kafka, le sink ne peut pas utiliser RowKind pour déterminer le type d'opération. Il doit alors s'appuyer sur un champ spécifique au sein du message, tel que {"op_type":"delete",data:{...}}. Pour supprimer des données où op_type vaut 'delete', vous devez explicitement passer une valeur à la colonne masquée en fonction de votre logique métier. L'exemple Flink SQL suivant montre comment supprimer des données dans Alibaba Cloud SelectDB en se basant sur un champ spécifique dans les données Kafka.

-- Example message: {"op_type":"delete",data:{"id":1,"name":"zhangsan"}}
CREATE TABLE KAFKA_SOURCE(
  data STRING,
  op_type STRING
) WITH (
  'connector' = 'kafka',
  ...
);

CREATE TABLE SELECTDB_SINK(
  id INT,
  name STRING,
  __DORIS_DELETE_SIGN__ INT
) WITH (
  'connector' = 'doris',
  'fenodes' = 'selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080',
  'table.identifier' = 'db.table',
  'username' = 'admin',
  'password' = '****',
  'sink.enable-delete' = 'false',        -- A value of false indicates that the event type is not inferred from RowKind.
  'sink.properties.columns' = 'id, name, __DORIS_DELETE_SIGN__'  -- Explicitly specify the columns for the Stream Load import.
);

INSERT INTO SELECTDB_SINK
SELECT json_value(data,'$.id') as id,
json_value(data,'$.name') as name, 
if(op_type='delete',1,0) as __DORIS_DELETE_SIGN__ 
FROM KAFKA_SOURCE;

FAQ

  • Q : Comment écrire des données BITMAP ?

    R : Consultez l'exemple ci-dessous :

    CREATE TABLE bitmap_sink (
      dt INT,
      page STRING,
      user_id INT 
    )
    WITH ( 
      'connector' = 'doris', 
      'fenodes' = 'selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080',
      'table.identifier' = 'test.bitmap_test', 
      'username' = 'admin', 
      'password' = '****', 
      'sink.label-prefix' = 'selectdb_label', 
      'sink.properties.columns' = 'dt,page,user_id,user_id=to_bitmap(user_id)'
    );
  • Q : Comment résoudre l'erreur errCode = 2, detailMessage = Label[label_0_1]has already been used, relate to txn[19650] ?

    R : Dans un scénario exactly-once, un job Flink doit être redémarré à partir du dernier checkpoint ou savepoint. Cette erreur survient si vous redémarrez le job à partir d'un état plus ancien. Si la sémantique exactly-once n'est pas requise, vous pouvez désactiver la validation en deux phases (2PC) en définissant sink.enable-2pc=false ou utiliser un sink.label-prefix différent.

  • Q : Comment résoudre l'erreur errCode = 2, detailMessage = transaction[19650]not found ?

    R : Cette erreur se produit lors de la phase de validation. Elle indique que l'ID de transaction enregistré dans le checkpoint a expiré dans ApsaraDB for SelectDB. Lorsque le connecteur tente de valider cette transaction expirée, le serveur signale que la transaction est introuvable. Dans ce cas, vous ne pouvez pas redémarrer le job à partir du checkpoint. Pour éviter ce problème, augmentez le paramètre streaming_label_keep_max_second dans ApsaraDB for SelectDB. La valeur par défaut est de 12 heures.

  • Q : Comment résoudre l'erreur errCode = 2, detailMessage = current running txns on db 10006 is 100, larger than limit 100 ?

    R : Cette erreur indique que le nombre de transactions d'importation simultanées pour une seule base de données a dépassé la limite système de 100. Pour résoudre ce problème, augmentez le paramètre max_running_txn_num_per_db dans ApsaraDB for SelectDB. Pour plus d'informations, consultez max_running_txn_num_per_db.

    Cette erreur peut également survenir si vous modifiez fréquemment le libellé et redémarrez le job. Dans les scénarios de validation en deux phases (2PC) (applicables aux modèles Duplicate Key et Aggregate Key), chaque job nécessite un libellé unique. Lorsqu'un job redémarre à partir d'un checkpoint, Flink annule uniquement les transactions qui ont été pré-validées mais pas encore validées. Si vous modifiez fréquemment le libellé avant de redémarrer, de nombreuses transactions pré-validées ne sont pas annulées et continuent de consommer le quota de transactions. Pour le modèle Unique Key, vous pouvez désactiver le 2PC et concevoir l'opérateur sink pour des écritures idempotentes.

  • Q : Comment garantir l'ordre des données au sein d'un lot lors de l'écriture dans une table utilisant le modèle Unique Key ?

    R : Ajoutez une configuration de colonne de séquence pour garantir l'ordre des données. Pour plus d'informations, consultez SEQUENCE.

  • Q : Pourquoi aucune donnée n'est-elle synchronisée alors que le job Flink ne signale aucune erreur ?

    R : Ce comportement dépend de la version du connecteur. Dans les versions antérieures à 1.1.0, les écritures sont effectuées par lots et pilotées par les données ; vous devez donc vérifier que la source en amont produit des données. À partir de la version 1.1.0, les écritures sont déclenchées par les checkpoints, que vous devez activer pour écrire des données.

  • Q : Comment résoudre l'erreur tablet writer write failed, tablet_id=190958, txn_id=3505530, err=-235 ?

    R : Cette erreur survient généralement dans les versions du connecteur antérieures à 1.1.0. Elle est causée par une fréquence d'écriture trop élevée, ce qui crée trop de versions sur le tablet. Pour résoudre ce problème, augmentez les paramètres sink.buffer-flush.max-bytes et sink.buffer-flush.interval afin de réduire la fréquence des Stream Load.

  • Q : Comment ignorer les données incorrectes lors d'une importation Flink ?

    R : Si les données sources contiennent des enregistrements qui ne correspondent pas au schéma de la table de destination (par exemple, type de données ou longueur incorrects), le job Stream Load échoue et Flink réessaie en continu. Pour ignorer ces données incorrectes, vous pouvez soit désactiver le mode strict de Stream Load en définissant strict_mode=false,max_filter_ratio=1, soit ajouter une étape de transformation pour filtrer les données invalides avant qu'elles n'atteignent l'opérateur sink.

  • Q : Comment mapper la table source à la table ApsaraDB for SelectDB ?

    R : Lorsque vous importez des données à l'aide du connecteur Flink Doris, assurez-vous que les mappages suivants sont corrects : (1) Les colonnes et les types de la table source doivent correspondre à ceux définis dans Flink SQL. (2) Les colonnes et les types dans Flink SQL doivent correspondre à ceux de la table ApsaraDB for SelectDB.

  • Q : Comment résoudre l'erreur TApplicationException: get_next failed: out of sequence response: expected 4 but got 3 ?

    R : Cette erreur indique un bug de concurrence au sein du framework Thrift sous-jacent. Pour la résoudre, mettez à niveau vers la dernière version du connecteur Flink Doris et utilisez une version compatible de Flink.

  • Q : Comment résoudre l'erreur DorisRuntimeException: Fail to abort transaction 26153 with urlhttp://192.168.XX.XX ?

    R : Pour diagnostiquer ce problème, recherchez l'expression abort transaction response dans les journaux du TaskManager. Le code d'état HTTP dans l'entrée de journal indique si le problème provient du client ou du serveur.