Tous les produits
Search
Centre de documentation

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

Dernière mise à jour :Aug 11, 2026

Créez un connecteur de réception Tablestore pour exporter les données d'un sujet source ApsaraMQ for Kafka vers une table Tablestore.

Prérequis

Étape 1 : Créer une table Tablestore

Créez une table Tablestore afin de stocker les données diffusées depuis ApsaraMQ for Kafka. Pour plus d'informations, consultez la section Procédure.

Cet exemple utilise une instance nommée ots-sink et une table de données nommée ots_sink_table. Lors de la création de la table, définissez trois clés primaires : topic de type STRING (définie comme clé de partition), partition de type INTEGER et offset de type INTEGER.

Étape 2 : Créer et démarrer le connecteur de réception Tablestore

  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.

  4. Sur la page Create Task, définissez les champs Task Name et Description, configurez les paramètres suivants, puis cliquez sur Save.

    • Create Task

      1. À l'étape Source, sélectionnez Message Queue for Apache Kafka en tant que 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.

        Chine (Pékin)

        Kafka instance

        Instance source Message Queue for Apache Kafka.

        alikafka_post-cn-jte3****

        Topic

        Sujet à partir duquel les messages sont consommés.

        demo-topic

        Group ID

        Groupe de consommateurs de l'instance source.

        • Quick Create (recommandé) : un ID de groupe au format GID_EVENTBRIDGE_xxx est automatiquement créé.

        • Use Existing : sélectionnez un ID de groupe existant. Ne partagez pas cet ID de groupe avec d'autres services afin d'éviter toute perturbation de la consommation des messages.

        Quick Create

        Consumer offset

        Offset à partir duquel la consommation des messages commence.

        • 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 champ Network configuration est défini sur Self-managed Internet.

        vpc-bp17fapfdj0dwzjkd****

        vSwitch

        Requis uniquement si le champ Network configuration est défini sur Self-managed Internet.

        vsw-bp1gbjhj53hdjdkg****

        Security group

        Requis uniquement si le champ Network configuration est défini sur Self-managed Internet.

        alikafka_pre-cn-7mz2****

        Data Format

        Format d'encodage du contenu des messages. Nous vous recommandons d'utiliser Json si vous n'avez pas d'exigences spécifiques en matière d'encodage.

        • Json : encode les données binaires sous forme d'objet JSON dans la charge utile à l'aide d'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 valides : 1 à 10 000.

        100

        Interval (Unit: Seconds)

        Paramètre de Advanced configuration. Intervalle, en secondes, auquel les messages sont agrégés et envoyés vers la destination. Valeurs valides : 0 à 15. Une valeur de 0 signifie que les messages sont livrés immédiatement.

        3

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

      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 rubrique Utiliser Function Compute pour nettoyer les données des messages.

      4. À l'étape Sink, définissez le champ Service Type sur Tablestore et configurez les paramètres suivants.

        Parameter

        Description

        Example

        Instance Name

        Nom de l'instance Tablestore de destination.

        ots-sink

        Destination Table

        Nom de la table de données Tablestore de destination.

        ots_sink_table

        Primary Key

        Définissez les règles d'extraction des valeurs de clé primaire et des colonnes d'attribut à l'aide de la syntaxe JsonPath. Lorsque le champ Data Format de la Source est défini sur JSON, les données de sortie de ApsaraMQ for Kafka se présentent au format suivant :

        {
            "data": {
                "topic": "demo-topic",
                "partition": 0,
                "offset": 2,
                "timestamp": 1739756629123,
                "headers": {
                    "headers": [],
                    "isReadOnly": false
                },
                "key":"ots-sink-k1",
                "value": "ots-sink-v1"
            },
            "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"
        }

        Par exemple, pour une clé primaire nommée topic, définissez la règle d'extraction de la valeur sur $.data.topic.

        Attribute Column

        Par exemple, pour une colonne d'attribut nommée key, définissez la règle d'extraction de la valeur sur $.data.key.

        Operation Mode

        Méthode d'écriture des données dans Tablestore.

        • put : si un enregistrement possédant la même clé primaire existe déjà, les nouvelles données écrasent les données existantes.

        • update : lorsque deux enregistrements partagent la même clé primaire, des colonnes incrémentielles sont ajoutées à la ligne sans supprimer les colonnes existantes.

        • delete : supprime la ligne correspondant à la clé primaire spécifiée.

        put

        Network configuration

        • VPC : utilisez un VPC pour acheminer les messages Kafka vers Tablestore.

        • Public Network : achemine les messages Kafka vers Tablestore via le réseau public.

        VPC

        VPC

        Sélectionnez l'ID du VPC. Ce paramètre est requis uniquement si le champ Network Configuration est défini sur VPC.

        vpc-bp17fapfdj0dwzjkd****

        vSwitch

        Sélectionnez l'ID du vSwitch. Ce paramètre est requis uniquement si le champ Network Configuration est défini sur VPC.

        vsw-bp1gbjhj53hdjdkg****

        Security group

        Sélectionnez un groupe de sécurité. Ce paramètre est requis uniquement lorsque le champ Network Configuration est défini sur VPC.

        test_group

    • Task Properties

      Configurez la politique de nouvelle tentative pour les livraisons d'événements ayant échoué et la méthode de gestion des erreurs. Pour plus d'informations, consultez la rubrique Politiques de nouvelle tentative et files d'attente de lettres mortes.

  5. Revenez à la page Tasks. Recherchez votre tâche et cliquez sur Enable dans la colonne Actions.

  6. Dans la boîte de dialogue Note, lisez les informations affichées, puis cliquez sur OK.

    Une fois la tâche activée, elle démarre dans un délai de 30 à 60 secondes. Vous pouvez surveiller son état dans la colonne Status de la page Tasks.

Étape 3 : Tester le connecteur de réception Tablestore

  1. Sur la page Tasks, dans la colonne Event Source correspondant à votre tâche, cliquez sur le sujet source.

  2. Sur la page des détails du sujet, cliquez sur Send Test Message.

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

    Dans l'onglet Console, définissez Message Key sur ots-sink-k1, Message Content sur ots-sink-v1 et Send to Specified Partition sur No.

  4. Revenez à la page Tasks et, dans la colonne Event Target correspondant à votre tâche, cliquez sur le nom de la table de destination.

  5. Sur la page Manage Table, cliquez sur l'onglet Data Management pour afficher les données de la table Tablestore.

    La table de données contient les colonnes topic (clé primaire), partition (clé primaire), offset (clé primaire), key et value. Un enregistrement contenant des valeurs telles que partition=2, offset=10, key=ots-sink-k1 et value=ots-sink-v1 confirme que le message a bien été écrit dans Tablestore.