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
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.
Prérequis
-
Assurez-vous de la connectivité réseau entre votre source de données, Flink et SelectDB.
-
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.
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.
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 Connectorcomme 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 Connectorcorrespondant et placez-le dans le répertoirelibde votre installation Flink. Pour le lien de téléchargement, consultez Package JAR.-
Pour ajouter le
Flink Doris Connectoren 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
-
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.
-
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 -
Décompressez le package d'installation.
tar -zxvf flink-1.16.3-bin-scala_2.12.tgz -
Accédez au répertoire
libdans 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
-
-
Démarrez le cluster Flink.
Dans le répertoire
binde votre installation Flink, exécutez la commande suivante :./start-cluster.sh
SelectDB cible
Créez une instance ApsaraDB for SelectDB. Pour plus d'informations, consultez Créer une instance.
Connectez-vous à l'instance. Pour plus d'informations, consultez Se connecter à une instance.
-
Créez une base de données de test nommée
test.CREATE DATABASE test; -
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
Créez une instance ApsaraDB RDS for MySQL.
-
Créez une base de données de test nommée
test.CREATE DATABASE test; -
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 ); -
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
-
Démarrez le client Flink SQL.
Dans le répertoire
binde votre installation Flink, exécutez la commande suivante :./sql-client.sh -
Dans le client Flink SQL, soumettez un job Flink.
-
Créez une table source MySQL.
La clause
WITHdans l'instruction suivante spécifie la configuration de laMySQL 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' ); -
Créez une table sink SelectDB.
La clause
WITHdans 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' = '****' ); -
Synchronisez les données de la table source MySQL vers la table sink SelectDB.
INSERT INTO employees_sink SELECT * FROM employees_source;
-
-
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
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 |
||
|
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-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, |
tbl.*` synchronise |
|
excluding-tables |
Les tables à exclure de la synchronisation. Le format est identique à celui de |
||
|
mysql-conf |
La configuration de la source MySQL CDC. Pour plus d'informations, consultez Connecteur MySQL CDC. Les paramètres |
||
|
oracle-conf |
La configuration de la source Oracle CDC. Pour plus d'informations, consultez Connecteur Oracle CDC. Les paramètres |
||
|
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 |
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.
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 : |
|
table.identifier |
Aucune |
Oui |
Le nom de la base de données et de la table. Exemple : |
|
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 : |
|
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 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 ( |
|
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 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 |
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
-
Ajoutez les dépendances suivantes à votre projet Maven.
Dépendances Maven
-
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=falseou 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_seconddans 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_dbdans 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-bytesetsink.buffer-flush.intervalafin 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 responsedans 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.