Tous les produits
Search
Centre de documentation

ApsaraDB for SelectDB:Importer des données à l'aide de Kafka

Dernière mise à jour :Aug 21, 2026

ApsaraDB for SelectDB prend en charge le connecteur Doris Kafka pour s'abonner automatiquement aux données et les synchroniser depuis Kafka. Cette rubrique explique comment utiliser le connecteur Doris Kafka pour synchroniser les données vers ApsaraDB for SelectDB.

Informations générales

Kafka Connect est un outil qui permet de diffuser des données de manière fiable entre Apache Kafka et d'autres systèmes. Vous pouvez définir des connecteurs pour déplacer de grands volumes de données vers ou depuis Kafka.

Le connecteur Kafka fourni par la communauté Doris s'exécute au sein d'un cluster Kafka Connect. Il lit les données d'une rubrique Kafka et les écrit dans ApsaraDB for SelectDB.

Dans les scénarios professionnels, les utilisateurs utilisent généralement le connecteur Debezium pour pousser les données de modification de base de données vers Kafka, ou appellent une API pour écrire des données au format JSON dans Kafka en temps réel. Le connecteur Doris Kafka s'abonne automatiquement aux données présentes dans Kafka et synchronise ces données vers ApsaraDB for SelectDB.

Modes d'exécution de Kafka Connect

Kafka Connect propose deux modes d'exécution.

Mode autonome

Avertissement

Le mode autonome n'est pas recommandé pour les environnements de production.

Configuration du mode autonome

Configurez le fichier connect-standalone.properties.

# Modify the broker address
bootstrap.servers=127.0.0.1:9092

Dans le répertoire de configuration de Kafka, créez un fichier connect-selectdb-sink.properties et ajoutez le contenu suivant :

name=test-selectdb-sink
connector.class=org.apache.doris.kafka.connector.DorisSinkConnector
topics=topic_test
doris.topic2table.map=topic_test:test_kafka_tbl
buffer.count.records=10000
buffer.flush.time=120
buffer.size.bytes=5000000
doris.urls=selectdb-cn-4xl3jv1****-public.selectdbfe.rds.aliyuncs.com
doris.http.port=8030
doris.query.port=9030
doris.user=admin
doris.password=****
doris.database=test_db
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.json.JsonConverter

Démarrage en mode autonome

$KAFKA_HOME/bin/connect-standalone.sh -daemon $KAFKA_HOME/config/connect-standalone.properties $KAFKA_HOME/config/connect-selectdb-sink.properties

Mode distribué

Configuration du mode distribué

Configurez le fichier connect-distributed.properties.

# Modify the broker address
bootstrap.servers=127.0.0.1:9092
 
# Modify group.id. The ID must be the same for all workers in the same cluster.
group.id=connect-cluster

Démarrage en mode distribué

$KAFKA_HOME/bin/connect-distributed.sh -daemon $KAFKA_HOME/config/connect-distributed.properties

Ajout d'un connecteur

curl -i http://127.0.0.1:8083/connectors -H "Content-Type: application/json" -X POST -d '{
  "name":"test-selectdb-sink-cluster",
  "config":{
    "connector.class":"org.apache.doris.kafka.connector.DorisSinkConnector",
    "topics":"topic_test",
    "doris.topic2table.map": "topic_test:test_kafka_tbl",
    "buffer.count.records":"10000",
    "buffer.flush.time":"120",
    "buffer.size.bytes":"5000000",
    "doris.urls":"selectdb-cn-4xl3jv1****-public.selectdbfe.rds.aliyuncs.com",
    "doris.user":"admin",
    "doris.password":"***",
    "doris.database":"test_db",
    "doris.http.port":"8030",
    "doris.query.port":"9030",
    "key.converter":"org.apache.kafka.connect.storage.StringConverter",
    "value.converter":"org.apache.kafka.connect.json.JsonConverter"
  }
}'

Paramètres

Paramètre

Description

name

Nom du connecteur. Il doit s'agir d'une chaîne ne contenant aucun caractère de contrôle ISO et doit être unique dans l'environnement Kafka Connect.

connector.class

Nom de classe ou alias du connecteur. Définissez cette valeur sur org.apache.doris.kafka.connector.DorisSinkConnector.

topics

Liste des rubriques sources séparées par des virgules.

doris.topic2table.map

Mappage entre les rubriques et les tables. Séparez plusieurs mappages par une virgule (,). Exemple : topic1:tb1,topic2:tb2. Si ce paramètre n'est pas spécifié, le connecteur suppose que les noms de rubrique et de table sont identiques.

buffer.count.records

Nombre d'enregistrements mis en mémoire tampon pour chaque partition Kafka avant d'être vidés dans ApsaraDB for SelectDB. La valeur par défaut est de 10 000 enregistrements.

buffer.flush.time

Intervalle en secondes pour vider la mémoire tampon. La valeur par défaut est 120.

buffer.size.bytes

Taille cumulative en octets des enregistrements mis en mémoire tampon pour chaque partition Kafka. La valeur par défaut est 5 000 000.

doris.urls

Point de terminaison de connexion à ApsaraDB for SelectDB.

Vous pouvez obtenir les paramètres pertinents sur la page Instance Details > Network Information dans la console ApsaraDB for SelectDB.

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

doris.http.port

Le port HTTP par défaut pour ApsaraDB for SelectDB est 8080.

doris.query.port

Le port du protocole MySQL de ApsaraDB for SelectDB est 9030 par défaut.

doris.user

Nom d'utilisateur ApsaraDB for SelectDB.

doris.password

Mot de passe ApsaraDB for SelectDB.

doris.database

Base de données ApsaraDB for SelectDB dans laquelle les données sont écrites.

key.converter

Classe de convertisseur JSON pour la clé.

value.converter

Classe de convertisseur JSON pour la valeur.

jmx

Indique si les métriques internes du connecteur doivent être récupérées via JMX. Pour plus d'informations, consultez Doris-Connector-JMX. La valeur par défaut est true.

enable.delete

Indique si les opérations de suppression doivent être synchronisées. La valeur par défaut est false.

label.prefix

Préfixe d'étiquette pour les données importées à l'aide de Stream Load. La valeur par défaut est le nom de l'application du connecteur.

auto.redirect

Lorsque cette option est activée, le connecteur redirige les requêtes Stream Load via le frontend (FE) vers un backend (BE) cible, ce qui évite d'avoir à récupérer les informations BE.

load.model

Méthode d'importation des données. Les méthodes suivantes sont prises en charge :

  • stream_load : importe les données directement dans SelectDB.

  • copy_into : importe les données dans le stockage d'objets, puis les charge dans SelectDB.

La valeur par défaut est stream_load.

sink.properties.*

Paramètres d'importation pour Stream Load.

Exemple : pour spécifier un séparateur de colonne, utilisez sink.properties.column_separator=,.

Pour plus d'informations, consultez Stream Load.

delivery.guarantee

Spécifie la garantie de cohérence des données lors de la consommation des données Kafka et de leur importation dans ApsaraDB for SelectDB. Les niveaux pris en charge sont at_least_once et exactly_once, et la valeur par défaut est at_least_once.

Actuellement, ApsaraDB for SelectDB peut uniquement garantir que les données importées à l'aide de copy into respectent la sémantique exactly_once.

enable.2pc

Indique si la validation en deux phases doit être activée pour garantir la sémantique exactly-once.

Remarque

Pour les autres configurations courantes des récepteurs Kafka Connect, consultez Configuring Connectors.

Exemples

Prérequis

  1. Installez un cluster Apache Kafka ou Confluent Cloud, version 2.4.0 ou ultérieure. Cet exemple utilise un environnement Kafka à nœud unique.

    # Download and decompress the package
    wget https://archive.apache.org/dist/kafka/2.4.0/kafka_2.12-2.4.0.tgz
    tar -zxvf kafka_2.12-2.4.0.tgz
    cd kafka_2.12-2.4.0/
    bin/zookeeper-server-start.sh -daemon config/zookeeper.properties 
    bin/kafka-server-start.sh -daemon config/server.properties
  2. Téléchargez doris-kafka-connector-1.0.0.jar et placez le fichier JAR dans le répertoire KAFKA_HOME/libs.

  3. Créez une instance ApsaraDB for SelectDB. Pour plus d'informations, consultez Create an instance.

  4. Connectez-vous à une instance ApsaraDB for SelectDB à l'aide du protocole MySQL. Pour plus d'informations, consultez Connect to an instance.

  5. Créez une base de données de test et une table de test.

    1. Créez une base de données de test.

      CREATE DATABASE test_db;
    2. Créez une table de test.

      USE test_db;
      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;

Exemple 1 : Synchronisation des données JSON

  1. Configurer le récepteur SelectDB

    En prenant le mode autonome comme exemple, créez un fichier selectdb-sink.properties dans le répertoire de configuration de Kafka et ajoutez le contenu suivant :

    name=selectdb_sink
    connector.class=org.apache.doris.kafka.connector.DorisSinkConnector
    topics=test_topic
    doris.topic2table.map=test_topic:example_tbl
    buffer.count.records=10000
    buffer.flush.time=120
    buffer.size.bytes=5000000
    doris.urls=selectdb-cn-4xl3jv1****-public.selectdbfe.rds.aliyuncs.com
    doris.http.port=8030
    doris.query.port=9030
    doris.user=admin
    doris.password=***
    doris.database=test_db
    key.converter=org.apache.kafka.connect.storage.StringConverter
    value.converter=org.apache.kafka.connect.json.JsonConverter
    
    #Optional: Configure a dead-letter queue
    errors.tolerance=all
    errors.deadletterqueue.topic.name=test_error
    errors.deadletterqueue.context.headers.enable = true
    errors.deadletterqueue.topic.replication.factor=1
  2. Démarrer Kafka Connect

    bin/connect-standalone.sh -daemon config/connect-standalone.properties config/selectdb-sink.properties

Exemple 2 : Utiliser Debezium pour synchroniser les données MySQL vers ApsaraDB for SelectDB

Dans de nombreux scénarios professionnels, il est nécessaire de synchroniser les données d'une base de données opérationnelle en temps réel. Cela nécessite l'utilisation du mécanisme de capture des données modifiées (CDC) de la base de données.

Debezium est un outil CDC basé sur Kafka Connect qui peut se connecter à diverses bases de données telles que MySQL, PostgreSQL, SQL Server, Oracle et MongoDB. Il envoie continuellement les modifications de données vers une rubrique Kafka dans un format unifié pour une consommation en temps réel par les récepteurs en aval. Cet exemple utilise MySQL.

  1. Téléchargez Debezium.

    wget https://repo1.maven.org/maven2/io/debezium/debezium-connector-mysql/1.9.8.Final/debezium-connector-mysql-1.9.8.Final-plugin.tar.gz
  2. Décompressez le fichier téléchargé.

    tar -zxvf debezium-connector-mysql-1.9.8.Final-plugin.tar.gz
  3. Placez tous les fichiers JAR extraits dans le répertoire KAFKA_HOME/libs.

  4. Configurez la source MySQL.

    Créez un fichier mysql-source.properties dans le répertoire de configuration de Kafka et ajoutez le contenu suivant :

    name=mysql-source
    connector.class=io.debezium.connector.mysql.MySqlConnector
    database.hostname=rm-bp17372257wkz****.rwlb.rds.aliyuncs.com
    database.port=3306
    database.user=testuser
    database.password=****
    database.server.id=1
    # A unique identifier for this client in Kafka
    database.server.name=test123
    # The databases and tables to sync. By default, all databases and tables are synced.
    database.include.list=test
    table.include.list=test.test_table
    database.history.kafka.bootstrap.servers=localhost:9092
    # The Kafka topic used to store database schema changes
    database.history.kafka.topic=dbhistory
    transforms=unwrap
    # See https://debezium.io/documentation/reference/stable/transformations/event-flattening.html
    transforms.unwrap.type=io.debezium.transforms.ExtractNewRecordState
    # Record delete events
    transforms.unwrap.delete.handling.mode=rewrite

    Après la configuration, le format de nom de rubrique Kafka par défaut est SERVER_NAME.DATABASE_NAME.TABLE_NAME.

    Remarque

    Pour les configurations Debezium, consultez Debezium connector for MySQL.

  5. Configurez le récepteur ApsaraDB for SelectDB.

    Créez un fichier selectdb-sink.properties dans le répertoire de configuration de Kafka et ajoutez le contenu suivant :

    name=selectdb-sink
    connector.class=org.apache.doris.kafka.connector.DorisSinkConnector
    topics=test123.test.test_table
    doris.topic2table.map=test123.test.test_table:test_table
    buffer.count.records=10000
    buffer.flush.time=120
    buffer.size.bytes=5000000
    doris.urls=selectdb-cn-4xl3jv1****-public.selectdbfe.rds.aliyuncs.com
    doris.http.port=8030
    doris.query.port=9030
    doris.user=admin
    doris.password=****
    doris.database=test
    key.converter=org.apache.kafka.connect.json.JsonConverter
    value.converter=org.apache.kafka.connect.json.JsonConverter
    
    #Optional: Configure a dead-letter queue
    #errors.tolerance=all
    #errors.deadletterqueue.topic.name=test_error
    #errors.deadletterqueue.context.headers.enable = true
    #errors.deadletterqueue.topic.replication.factor=1
    Remarque

    Lorsque vous synchronisez des données vers ApsaraDB for SelectDB, vous devez créer une base de données et une table au préalable.

  6. Démarrez Kafka Connect.

    bin/connect-standalone.sh -daemon config/connect-standalone.properties config/mysql-source.properties config/selectdb-sink.properties
    Remarque

    Après le démarrage, vous pouvez vérifier le fichier logs/connect.log pour confirmer que le service a démarré avec succès.

Utilisation avancée

Opérations sur les connecteurs

# Check the connector status
curl -i http://127.0.0.1:8083/connectors/test-selectdb-sink-cluster/status -X GET
# Delete the current connector
curl -i http://127.0.0.1:8083/connectors/test-selectdb-sink-cluster -X DELETE
# Pause the current connector
curl -i http://127.0.0.1:8083/connectors/test-selectdb-sink-cluster/pause -X PUT
# Resume the current connector
curl -i http://127.0.0.1:8083/connectors/test-selectdb-sink-cluster/resume -X PUT
# Restart tasks within the connector
curl -i http://127.0.0.1:8083/connectors/test-selectdb-sink-cluster/tasks/0/restart -X POST

Pour plus d'informations, consultez Connect REST Interface.

File d'attente des lettres mortes

Par défaut, une erreur de conversion entraîne l'échec du connecteur. Toutefois, vous pouvez configurer le connecteur pour qu'il tolère ces erreurs en les ignorant. Vous pouvez également écrire les détails de l'erreur, l'opération ayant échoué et l'enregistrement problématique dans une file d'attente des lettres mortes pour une analyse ultérieure.

errors.tolerance=all
errors.deadletterqueue.topic.name=test_error_topic
errors.deadletterqueue.context.headers.enable=true
errors.deadletterqueue.topic.replication.factor=1

Pour plus d'informations, consultez Error Reporting in Connect.

Connexion à un cluster Kafka activé pour SSL

Pour accéder à un cluster Kafka activé pour SSL via Kafka Connect, vous devez fournir un fichier de certificat (client.truststore.jks) pour authentifier la clé publique du courtier Kafka. Vous pouvez ajouter la configuration suivante à votre fichier connect-distributed.properties :

# Connect worker
security.protocol=SSL
ssl.truststore.location=/var/ssl/private/client.truststore.jks
ssl.truststore.password=test1234
 
# Embedded consumer for sink connectors
consumer.security.protocol=SSL
consumer.ssl.truststore.location=/var/ssl/private/client.truststore.jks
consumer.ssl.truststore.password=test1234

Pour plus d'informations sur la configuration de Kafka Connect pour se connecter à un cluster Kafka activé pour SSL, consultez Configure Kafka Connect.