Tous les produits
Search
Centre de documentation

ApsaraMQ for Kafka:Créer un connecteur de réception AnalyticDB

Dernière mise à jour :Aug 11, 2026

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.

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

  1. Connectez-vous à la console ApsaraMQ for Kafka. Sur la page Overview, sélectionnez une région dans la section Resource Distribution.

  2. Dans le volet de navigation de gauche, choisissez Connector Ecosystem Integration > Tasks.

  3. Sur la page Tasks, cliquez sur Create Task.

    • Configurez la tâche

      1. À 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

      2. À l'étape Filtering, définissez le Pattern Content pour filtrer les événements. Pour plus d'informations, consultez la page event pattern.

      3. À 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.

      4. À 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 userid du 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.

        Important

        Aprè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

  4. Revenez à la page VPC. Recherchez la tâche créée et cliquez sur Tasks dans la colonne Enable.

  5. 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

  1. Sur la page Tasks, recherchez votre tâche et cliquez sur le nom de la rubrique source dans la colonne Tasks.

  2. Sur la page des détails de la rubrique, cliquez sur Event Source.

  3. Dans le panneau Send Test Message, configurez le corps du message, puis cliquez sur Start to Send and Consume Message.

    Remarque

    Le 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-k1 et Message body sur {"userid":"xiaoming","source":"shanghai"}. Pour Send to a specific partition, sélectionnez No, puis cliquez sur OK.

  4. Sur la page OK, recherchez votre tâche et cliquez sur le nom de l'instance de destination dans la colonne Tasks.

  5. Sur la page Event Target de l'instance, cliquez sur Basic Information dans le coin supérieur droit.

  6. 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 à xiaoming et un source égal à shanghai. Cela confirme que les données ont été écrites avec succès dans la table de destination.