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
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 |
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 : |
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 :
La valeur par défaut est |
sink.properties.* | Paramètres d'importation pour Stream Load. Exemple : pour spécifier un séparateur de colonne, utilisez 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 Actuellement, ApsaraDB for SelectDB peut uniquement garantir que les données importées à l'aide de copy into respectent la sémantique |
enable.2pc | Indique si la validation en deux phases doit être activée pour garantir la sémantique exactly-once. |
Pour les autres configurations courantes des récepteurs Kafka Connect, consultez Configuring Connectors.
Exemples
Prérequis
-
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 Téléchargez doris-kafka-connector-1.0.0.jar et placez le fichier JAR dans le répertoire KAFKA_HOME/libs.
Créez une instance ApsaraDB for SelectDB. Pour plus d'informations, consultez Create an instance.
Connectez-vous à une instance ApsaraDB for SelectDB à l'aide du protocole MySQL. Pour plus d'informations, consultez Connect to an instance.
-
Créez une base de données de test et une table de test.
-
Créez une base de données de test.
CREATE DATABASE test_db; -
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
-
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 -
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.
-
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 -
Décompressez le fichier téléchargé.
tar -zxvf debezium-connector-mysql-1.9.8.Final-plugin.tar.gz Placez tous les fichiers JAR extraits dans le répertoire KAFKA_HOME/libs.
-
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=rewriteAprès la configuration, le format de nom de rubrique Kafka par défaut est
SERVER_NAME.DATABASE_NAME.TABLE_NAME.RemarquePour les configurations Debezium, consultez Debezium connector for MySQL.
-
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=1RemarqueLorsque vous synchronisez des données vers ApsaraDB for SelectDB, vous devez créer une base de données et une table au préalable.
-
Démarrez Kafka Connect.
bin/connect-standalone.sh -daemon config/connect-standalone.properties config/mysql-source.properties config/selectdb-sink.propertiesRemarqueAprè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.