Cette rubrique explique comment créer un connecteur de réception AnalyticDB. Ce connecteur permet de diffuser des données d'une rubrique source dans une instance ApsaraMQ for Kafka vers une table d'une base de données AnalyticDB.
Prérequis
Pour plus d'informations, consultez la page Prérequis.
Étape 1 : Créer les ressources de destination
Créez une ressource AnalyticDB for MySQL ou AnalyticDB for PostgreSQL.
AnalyticDB for MySQL : Dans la console AnalyticDB for MySQL, créez un cluster et un compte de base de données, connectez-vous au cluster, puis créez une base de données. Pour plus d'informations, consultez les pages Créer un cluster, Créer un compte de base de données, Se connecter à un cluster et Créer une base de données.
AnalyticDB for PostgreSQL : Dans la console AnalyticDB for PostgreSQL, créez une instance et un compte de base de données, puis connectez-vous à la base de données. Pour plus d'informations, consultez les pages Créer une instance, Créer et gérer des utilisateurs et Connexion client.
Cet exemple utilise une base de données AnalyticDB for MySQL nommée adb_sink_database et une table nommée adb_sink_table.
Étape 2 : Créer et activer le connecteur de réception AnalyticDB
Connectez-vous à la console ApsaraMQ for Kafka. Sur la page Overview, sélectionnez une région dans la section Resource Distribution.
Dans le volet de navigation de gauche, choisissez .
-
Sur la page Tasks, cliquez sur Create Task.
-
Configurez la tâche
-
À l'étape Source, sélectionnez Message Queue for Apache Kafka comme Data Provider. Configurez les paramètres suivants, puis cliquez sur Next.
Parameter
Description
Example
Region
Région de l'instance source Message Queue for Apache Kafka.
China (Beijing)
Kafka instance
Instance source Message Queue for Apache Kafka.
alikafka_post-cn-jte3****
Topic
Rubrique source des messages.
demo-topic
Group ID
Groupe de consommateurs de l'instance source.
-
Quick Create (recommandé) : Crée automatiquement un ID de groupe au format
GID_EVENTBRIDGE_xxx. -
Use Existing : Sélectionnez un ID de groupe existant. Ne partagez pas cet ID de groupe avec d'autres services pour éviter toute perturbation de la consommation des messages.
Quick Create
Consumer offset
Offset de démarrage de la consommation des messages.
-
Latest offset (latest)
-
Earliest offset (earliest)
Latest offset (latest)
Network configuration
Type de réseau pour le routage des messages.
-
Basic Network
-
Self-managed Internet
Basic Network
VPC
Requis uniquement si le paramètre Network configuration est défini sur Self-managed Internet.
vpc-bp17fapfdj0dwzjkd****
vSwitch
Requis uniquement si le paramètre Network configuration est défini sur Self-managed Internet.
vsw-bp1gbjhj53hdjdkg****
Security group
Requis uniquement si le paramètre Network configuration est défini sur Self-managed Internet.
alikafka_pre-cn-7mz2****
Data Format
Format d'encodage du contenu du message. Nous recommandons Json en l'absence d'exigences spécifiques.
-
Json : Encode les données binaires sous forme d'objet JSON dans la charge utile à l'aide de UTF-8.
-
Text : Encode les données binaires sous forme de chaîne UTF-8 dans la charge utile. Il s'agit du format par défaut.
-
Binary : Encode les données binaires sous forme de chaîne encodée en Base64 dans la charge utile.
Json
Messages
Paramètre de Advanced configuration. Nombre maximal de messages par lot. Une requête est envoyée lorsque le nombre cumulé de messages atteint cette valeur. Valeurs possibles : de 1 à 10 000.
100
Interval (Unit: Seconds)
Paramètre de Advanced configuration. Intervalle, en secondes, entre l'agrégation et l'envoi des messages vers la destination. Valeurs possibles : de 0 à 15. Une valeur de 0 signifie que les messages sont livrés immédiatement.
3
-
À l'étape Filtering, définissez le Pattern Content pour filtrer les événements. Pour plus d'informations, consultez la page event pattern.
À l'étape Transformation, configurez la transformation des données pour effectuer des opérations telles que le fractionnement, le mappage, l'enrichissement et le routage dynamique. Pour plus d'informations, consultez la page Use Function Compute to clean message data.
-
À l'étape Sink, définissez le Service Type sur AnalyticDB, configurez les paramètres suivants, puis cliquez sur Save.
Parameter
Description
Example
Instance type
Sélectionnez le type de base de données de votre instance de destination.
-
AnalyticDB for MySQL
-
AnalyticDB for PostgreSQL
MySQL version
AnalyticDB instance ID
Sélectionnez l'instance de destination.
gp-bp10uo5n536wd****
Database name
Sélectionnez la base de données de destination.
adb_sink_database
Table name
Sélectionnez la table de données de destination.
adb_sink_table
Data Mapping
Utilisez des expressions JSONPath pour définir les règles d'extraction des données. Lorsque le paramètre Data Format est défini sur Json à l'étape Source, les données diffusées depuis ApsaraMQ for Kafka sont encapsulées dans une structure CloudEvents, comme illustré ci-dessous :
{ "data": { "topic": "demo-topic", "partition": 0, "offset": 2, "timestamp": 1739756629123, "headers": { "headers": [], "isReadOnly": false }, "key":"adb-sink-k1", "value": { "userid":"xiaoming", "source":"shanghai" } }, "id": "7702ca16-f944-4b08-***-***-0-2", "source": "acs:alikafka", "specversion": "1.0", "type": "alikafka:Topic:Message", "datacontenttype": "application/json; charset=utf-8", "time": "2025-02-17T01:43:49.123Z", "subject": "acs:alikafka:alikafka_serverless-cn-lf6418u6701:topic:demo-topic", "aliyunaccountid": "1******6789" }Mappez chaque colonne de la table de destination à un champ du message source à l'aide d'une expression JSONPath. Par exemple, pour mapper le champ
useriddu message à une colonne de table, utilisez l'expression$.data.value.userid.Database username
Saisissez le nom d'utilisateur du compte de base de données.
user
Database password
Saisissez le mot de passe du compte de base de données.
Network configuration
-
Source : Connectez-vous à AnalyticDB via un VPC.
-
VPC : Connectez-vous à AnalyticDB via Internet public.
VPC
VPC
Sélectionnez l'ID du VPC. Ce paramètre est requis uniquement lorsque le paramètre Network configuration est défini sur Public Network.
vpc-bp17fapfdj0dwzjkd****
vSwitch
Sélectionnez l'ID du vSwitch. Ce paramètre est requis uniquement lorsque le paramètre Network configuration est défini sur VPC.
ImportantAprès avoir sélectionné un vSwitch, vous devez ajouter le bloc CIDR du vSwitch à la liste blanche d'adresses IP de l'instance AnalyticDB for MySQL. Pour plus d'informations, consultez la page Configure an IP address whitelist.
vsw-bp1gbjhj53hdjdkg****
Security group
Sélectionnez le groupe de sécurité. Ce paramètre est requis uniquement lorsque le paramètre Network configuration est défini sur VPC.
test_group
-
-
-
Revenez à la page VPC. Recherchez la tâche créée et cliquez sur Tasks dans la colonne Enable.
-
Dans la boîte de dialogue Actions, lisez le message, puis cliquez sur Note.
Le démarrage de la tâche prend entre 30 et 60 secondes. Vous pouvez surveiller la progression dans la colonne OK de la page Status.
Étape 3 : Vérifier le connecteur de réception AnalyticDB
Sur la page Tasks, recherchez votre tâche et cliquez sur le nom de la rubrique source dans la colonne Tasks.
Sur la page des détails de la rubrique, cliquez sur Event Source.
-
Dans le panneau Send Test Message, configurez le corps du message, puis cliquez sur Start to Send and Consume Message.
RemarqueLe corps du message doit être au format JSON. Les champs spécifiés dans vos règles de mappage des données seront extraits et écrits dans les colonnes correspondantes de la table de destination.
Dans la boîte de dialogue Start to Send and Consume Message, sélectionnez l'onglet Console. Définissez Message key sur
adb-sink-k1et Message body sur{"userid":"xiaoming","source":"shanghai"}. Pour Send to a specific partition, sélectionnez No, puis cliquez sur OK. Sur la page OK, recherchez votre tâche et cliquez sur le nom de l'instance de destination dans la colonne Tasks.
Sur la page Event Target de l'instance, cliquez sur Basic Information dans le coin supérieur droit.
-
Dans la console Data Management Service (DMS), exécutez l'instruction SQL suivante pour interroger toutes les données de la table.
SELECT * FROM adb_sink_table;La requête doit renvoyer un enregistrement avec un
useridégal àxiaominget unsourceégal àshanghai. Cela confirme que les données ont été écrites avec succès dans la table de destination.