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
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.
-
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
-
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
À 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.
À 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.
-
À 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
topicest extraite du champtopic, 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.topicvaluename:
$.data.valuevalueage:
$.data.offsetPartition 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.
-
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
Sur la page Tasks, localisez le connecteur de sortie MaxCompute et cliquez sur la rubrique source dans la colonne Event Source.
Sur la page des détails de la rubrique, cliquez sur Send Test Message.
-
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 surMaxCompute-V1et Send to specified partition sur No. -
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 -
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 :
topicestxxx(masqué),valueNameestMaxCompute-V1,valueAgeest4ettimeest2024-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.