Tous les produits
Search
Centre de documentation

ApsaraMQ for Kafka:Créer un connecteur de sortie MaxCompute

Dernière mise à jour :Aug 11, 2026

Créez un connecteur de sortie MaxCompute pour exporter les données d'une rubrique d'une instance ApsaraMQ for Kafka vers une table MaxCompute.

Prérequis

Pour connaître la procédure détaillée, consultez la rubrique Prérequis pour les connecteurs de sortie.

Notes

Pour utiliser la fonctionnalité de partitionnement de MaxCompute, vous devez créer une colonne de partition supplémentaire nommée time avec le type de données STRING lors de la création de la table.

Étape 1 : Créer la ressource de destination

Créez une table à l'aide du client MaxCompute. Pour plus d'informations, consultez la rubrique Créer une table.

L'exemple suivant crée une table nommée kafka_to_maxcompute comportant trois colonnes et avec la fonctionnalité de partitionnement activée :

CREATE TABLE IF NOT EXISTS kafka_to_maxcompute(topic STRING,valueName STRING,valueAge BIGINT) PARTITIONED by (time STRING);

Si vous n'utilisez pas la fonctionnalité de partitionnement, utilisez l'instruction suivante :

CREATE TABLE IF NOT EXISTS kafka_to_maxcompute(topic STRING,valueName STRING,valueAge BIGINT);

Après l'exécution de l'instruction, le résultat suivant s'affiche :

Sur la page Tables, consultez les informations relatives à la table kafka_to_maxcompute créée. Son type de partition est partitioned table et son type de table est internal table. Le schéma de la table contient trois champs : topic (string), valueName (string) et valueAge (bigint). Aucun d'entre eux n'est une clé primaire. Le champ de partition est time (string).

Étape 2 : Créer et démarrer le connecteur

  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. Dans la page Create Task, configurez les paramètres Task Name et Description. Suivez ensuite les instructions à l'écran pour configurer les autres paramètres.

    • Création de tâche

      1. Dans l'assistant de configuration Source, définissez le paramètre Data Provider sur ApsaraMQ for Kafka, configurez les paramètres suivants, puis cliquez sur Next.

        Paramètre

        Description

        Exemple

        Region

        La région où réside l'instance ApsaraMQ for Kafka.

        China (Hangzhou)

        ApsaraMQ for Kafka Instance

        L'ID de l'instance ApsaraMQ for Kafka dans laquelle les données à acheminer sont produites.

        alikafka_post-cn-9hdsbdhd****

        Topic

        La rubrique de l'instance ApsaraMQ for Kafka dans laquelle les données à acheminer sont produites.

        guide-sink-topic

        Group ID

        L'ID du groupe de l'instance ApsaraMQ for Kafka dans laquelle les données à acheminer sont produites.

        • Quickly Create : Le système crée automatiquement un groupe dont l'ID est au format GID_EVENTBRIDGE_xxx.

        • Use Existing Group : Sélectionnez l'ID d'un groupe existant qui n'est pas utilisé. Si vous sélectionnez un groupe existant en cours d'utilisation, la publication et l'abonnement aux messages existants seront affectés.

        Use Existing Group

        Consumer Offset

        • Latest Offset : Les messages sont consommés à partir du décalage le plus récent.

        • Earliest Offset : Les messages sont consommés à partir du décalage le plus ancien.

        Latest Offset

        Network Configuration

        Si une transmission de données transfrontalière est requise, sélectionnez Self-managed Internet. Dans les autres cas, sélectionnez Basic Network.

        Basic Network

        Data Format

        La fonctionnalité de format de données permet d'encoder les données binaires transmises depuis la source dans un format de données spécifique. Plusieurs formats de données sont pris en charge. Si vous n'avez pas d'exigences particulières en matière d'encodage, spécifiez Json comme valeur.

        • Json : Encode les données binaires au format JSON basé sur UTF-8 et les place dans la charge utile.

        • Text : Format par défaut. Encode les données binaires en une chaîne basée sur UTF-8 et les place dans la charge utile.

        • Binary : Les données binaires sont encodées en chaînes basées sur l'encodage Base64, puis placées dans la charge utile.

        Text

        Messages

        Un paramètre de Advanced Configuration. Le nombre maximal de messages à envoyer par lot. Une requête est envoyée uniquement lorsque le nombre de messages accumulés atteint cette valeur. La valeur doit être un entier compris entre 1 et 10 000.

        2000

        Interval (Unit: Seconds)

        Un paramètre de Advanced Configuration. L'intervalle auquel la fonction est appelée. Le système agrège les messages et les envoie à Function Compute à cet intervalle. La valeur doit être un entier compris entre 0 et 15. L'unité est la seconde. Une valeur de 0 signifie que les messages sont livrés immédiatement.

        3

      2. À l'étape Filtering, définissez un modèle d'événement pour filtrer les requêtes. Pour plus d'informations, consultez la rubrique Modèles d'événements.

      3. À l'étape Transformation, configurez le nettoyage des données pour un traitement complexe tel 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 paramètre Service Type sur acs.maxcompute et configurez les paramètres suivants.

        Paramètre

        Description

        Exemple

        AccessKey ID

        L'AccessKey ID de votre compte Alibaba Cloud pour accéder à MaxCompute.

        yourAccessKeyID

        AccessKey Secret

        L'AccessKey Secret de votre compte Alibaba Cloud.

        yourAccessKeySecret

        MaxCompute Project Name

        Sélectionnez un projet MaxCompute existant.

        test_compute

        MaxCompute Table Name

        Sélectionnez une table MaxCompute existante.

        kafka_to_maxcompute

        MaxCompute Table Input Parameter

        Après avoir sélectionné une table, configurez une value extraction rule pour chaque colonne. Dans l'exemple de message suivant, la valeur de la colonne topic est extraite du champ topic, donc la value extraction rule est $.topic.

        {
              'data': {
                'topic': 't_test',
                'partition': 2,
                'offset': 1,
                'timestamp': 1717048990499,
                'headers': {
                  'headers': [],
                  'isReadOnly': False
                },
                'key': 'MaxCompute-K1',
                'value': 'MaxCompute-V1'
              },
              'id': '9b05fc19-9838-4990-bb49-ddb942307d3f-2-1',
              'source': 'acs:alikafka',
              'specversion': '1.0',
              'type': 'alikafka:Topic:Message',
              'datacontenttype': 'application/json; charset=utf-8',
              'time': '2024-05-30T06:03:10.499Z',
              'aliyunaccountid': '1413397765616316'
            }

        topic: $.data.topic

        valuename: $.data.value

        valueage: $.data.offset

        Partition Dimension

        • Close : Désactive la fonctionnalité de partitionnement.

        • Enable : Active la fonctionnalité de partitionnement.

          Si vous activez le partitionnement, vous devez configurer des paramètres tels que la valeur de partition :

          • La valeur de partition prend en charge les variables temporelles {yyyy}, {MM}, {dd}, {HH} et {mm}, qui représentent respectivement l'année, le mois, le jour, l'heure et la minute. Les variables temporelles sont sensibles à la casse.

          • La valeur de partition peut également être une constante.

        Enable

        {yyyy}-{MM}-{dd}.{HH}:{mm}.suffix

        Network configuration

        • VPC : Livrez les messages Kafka à MaxCompute via un VPC.

        • Public Network : Livrez les messages Kafka à MaxCompute via Internet.

        internet

        VPC

        Sélectionnez un ID VPC. Ce paramètre est requis uniquement si vous définissez Network Configuration sur VPC.

        vpc-bp17fapfdj0dwzjkd****

        vSwitch

        Sélectionnez un ID vSwitch. Ce paramètre est requis uniquement si vous définissez Network Configuration sur VPC.

        vsw-bp1gbjhj53hdjdkg****

        Security Group

        Sélectionnez un groupe de sécurité. Ce paramètre est requis uniquement si vous définissez Network Configuration sur VPC.

        test_group

    • Propriétés de la tâche

      Configurez des politiques de nouvelle tentative et des files d'attente de lettres mortes pour gérer les erreurs de livraison. Pour plus d'informations, consultez la rubrique Politiques de nouvelle tentative et files d'attente de lettres mortes.

  5. Une fois les configurations précédentes terminées, cliquez sur Save. Sur la page Tasks, localisez la tâche de connecteur de sortie MaxCompute que vous avez créée. La colonne Status affiche Starting. Lorsque le statut passe à Running, le processus de création est terminé.

Étape 3 : Tester le connecteur

  1. Sur la page Tasks, localisez le connecteur de sortie MaxCompute et cliquez sur la rubrique source dans la colonne Event Source.

  2. Sur la page des détails de la rubrique, cliquez sur Send Test Message.

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

    Dans l'onglet Console, définissez Message key sur MaxCompute-K1, Message body sur MaxCompute-V1 et Send to specified partition sur No.

  4. Accédez à la console MaxCompute et exécutez l'instruction SQL suivante pour afficher les informations de partition.

    show PARTITIONS kafka_to_maxcompute;

    Le résultat suivant s'affiche :

    OK
    OK
    OK
    OK
    OK
    OK
    OK
    OK
    OK
    OK
    OK
    OK
    time=2024-05-31.16:37.suffix
    OK
    2024-05-31 16:42:49 INFO ==================================================================
    2024-05-31 16:42:49 INFO Exit code of the Shell command 0
    2024-05-31 16:42:49 INFO --- Invocation of Shell command completed ---
    2024-05-31 16:42:49 INFO Shell run successfully!
    2024-05-31 16:42:49 INFO Current task status: FINISH
    2024-05-31 16:42:49 INFO Cost time is: 1.411s
  5. En vous basant sur les informations de partition, exécutez l'instruction suivante pour afficher les données de la partition.

    SELECT * FROM kafka_to_maxcompute WHERE time="2024-05-31.16:37.suffix";

    La requête renvoie un enregistrement : topic est xxx (masqué), valueName est MaxCompute-V1, valueAge est 4 et time est 2024-05-31.16:37.suffix. Cela confirme que les données ont été exportées d'ApsaraMQ for Kafka vers la table MaxCompute partitionnée.