Créez et exécutez des tâches d'extraction, de transformation et de chargement (ETL) pour nettoyer, transformer et déverser les données dans vos instances ApsaraMQ for Kafka. Cette rubrique explique comment créer une tâche ETL dans la console ApsaraMQ for Kafka afin de transmettre les données traitées depuis un topic source vers un topic de destination.
Prérequis
Avant d'exécuter une tâche ETL, effectuez les opérations suivantes :
-
Créez un topic source et un topic de destination sur vos instances ApsaraMQ for Kafka. Pour plus d'informations, consultez la section Étape 1 : Créer un topic.
RemarqueSi vous créez manuellement les topics auxiliaires requis par une tâche ETL, vous devez également créer les topics destinés à stocker les informations relatives à la tâche. Pour connaître les paramètres de création des topics, reportez-vous à l'étape 3 de la section « Créer une tâche ETL » de cette rubrique.
Activez Function Compute. Pour plus d'informations, consultez la page Activer Function Compute.
-
Obtenez les autorisations requises si vous utilisez un utilisateur RAM (Resource Access Management). Pour plus d'informations, consultez la page Accorder des autorisations aux utilisateurs RAM.
Le code suivant présente un exemple de politique définissant les autorisations nécessaires :
{ "Version": "1", "Statement": [ { "Action": [ "kafka:CreateETLTask", "kafka:ListETLTask", "kafka:DeleteETLTask" ], "Resource": "*", "Effect": "Allow" } ] }
Contexte
Le processus ETL extrait, transforme et charge des données vers une destination. Écrivez une fonction pour implémenter la logique de la tâche ETL. ApsaraMQ for Kafka appelle cette fonction pour traiter les données du topic source et les envoyer au topic de destination.
Lors du traitement des données, Function Compute crée automatiquement le service et la fonction correspondants. Le nom du service créé suit le format
_FC-kafka-nom de la tâche ETL.Le topic source et le topic de destination sur les instances ApsaraMQ for Kafka doivent se trouver dans la même région.
Function Compute vous permet de consulter les journaux des appels de fonction pour résoudre les problèmes. Pour plus d'informations, consultez la page Configurer la fonctionnalité de journalisation.
La fonctionnalité de tâche ETL de ApsaraMQ for Kafka est en aperçu public. Elle est indépendante des instances ApsaraMQ for Kafka. ApsaraMQ for Kafka ne facture pas l'utilisation de cette fonctionnalité. Si la tâche ETL que vous créez dépend d'autres services, référez-vous aux règles de facturation des services correspondants.
Activer ETL
Les tâches ETL sont créées dans le module Intégration de l'écosystème Connecteur de la console ApsaraMQ for Kafka. Ce module offre des capacités de filtrage et de transformation des données. Pour plus d'informations, consultez la page Vue d'ensemble.
Lors de votre première utilisation de la fonctionnalité de tâche ETL, autorisez ApsaraMQ for Kafka à accéder aux services connexes. Après confirmation, le système crée automatiquement le rôle lié au service AliyunServiceRoleForAlikafkaETL. ApsaraMQ for Kafka peut endosser ce rôle pour accéder aux services utilisés lors du traitement d'une tâche ETL. Pour plus d'informations, consultez la page Rôles liés au service.
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 ETL Tasks, cliquez sur Create Task.
Dans le message Service Authorization qui s'affiche, cliquez sur OK.
Créer une tâche ETL
Cette section décrit comment créer et déployer une tâche ETL pour extraire les données d'un topic source, les traiter, puis charger les données traitées vers un topic de destination.
Sur la page ETL Tasks, cliquez sur Create Task.
À l'étape Configure Basic Information, saisissez un nom de tâche et cliquez sur Next.
-
À l'étape Configure Source and Destination, spécifiez la source de données, le topic de destination et les informations de consommation. Cliquez ensuite sur Next.
Paramètre
Description
Exemple
Instance
Les instances auxquelles appartiennent le topic source et le topic de destination.
alikafka_pre-cn-7pp2bz47****
alikafka_post-cn-2r42bvbm****
Topic
Le topic source et le topic de destination.
RemarqueLe topic source et le topic de destination doivent être différents.
topic****
test****
Consumer Offset
Le décalage à partir duquel vous souhaitez consommer les messages. Ce paramètre s'affiche uniquement après avoir cliqué sur Advanced Settings. Valeurs valides :
Earliest Offset : Consomme les messages à partir du décalage le plus ancien.
Latest Offset : Consomme les messages à partir du décalage le plus récent.
Latest Offset
Failure Handling Policy
Indique s'il faut envoyer les messages suivants si l'envoi d'un message échoue. Ce paramètre s'affiche uniquement après avoir cliqué sur Advanced Settings. Valeurs valides :
Continue Subscription : Envoie les messages suivants même si l'envoi d'un message échoue.
Stop Subscription : N'envoie pas les messages suivants si l'envoi d'un message échoue.
Continue Subscription
Resource Creation Method
La méthode utilisée pour créer les topics auxiliaires requis par la tâche ETL. Ce paramètre s'affiche uniquement après avoir cliqué sur Advanced Settings. Valeurs valides :
Auto
Manual
Auto
Consumer Group
Le consumer group utilisé par la tâche ETL. Ce paramètre s'affiche uniquement si vous définissez le paramètre Resource Creation Method sur Manual. Nous vous recommandons d'utiliser etl-cluster comme préfixe du nom du consumer group.
etl-cluster-kafka
Task Offset Topic
Le topic utilisé pour stocker les décalages des consommateurs. Ce paramètre s'affiche uniquement si vous définissez le paramètre Manual sur Resource Creation Method.
Topic : le nom du topic. Nous vous recommandons d'utiliser etl-offset comme préfixe du nom du topic.
Partitions : le nombre de partitions dans le topic. Ce paramètre doit être défini sur une valeur supérieure à 1.
Storage Engine : le moteur de stockage utilisé par le topic. Ce paramètre doit être défini sur Local Storage.
RemarqueVous ne pouvez définir le paramètre Storage Engine sur Local Storage que si vous créez un topic sur une instance Professional Edition.
cleanup.policy : la politique de nettoyage des journaux utilisée par le topic. Ce paramètre doit être défini sur Compact.
etl-offset-kafka
Task Configuration Topic
Le topic utilisé pour stocker les configurations des tâches. Ce paramètre s'affiche uniquement si vous définissez le paramètre Manual sur Resource Creation Method.
Topic : le nom du topic. Nous vous recommandons d'utiliser etl-config comme préfixe du nom du topic.
Partitions : le nombre de partitions dans le topic. Ce paramètre doit être défini sur 1.
Storage Engine : le moteur de stockage utilisé par le topic. Ce paramètre doit être défini sur Local Storage.
RemarqueVous ne pouvez définir le paramètre Storage Engine sur Local Storage que si vous créez un topic sur une instance Professional Edition.
cleanup.policy : la politique de nettoyage des journaux utilisée par le topic. Ce paramètre doit être défini sur Compact.
etl-config-kafka
Task Status Topic
Le topic utilisé pour stocker l'état de la tâche. Ce paramètre s'affiche uniquement si vous définissez le paramètre Manual sur Resource Creation Method.
Topic : le nom du topic. Nous vous recommandons d'utiliser etl-status comme préfixe du nom du topic.
Partitions : le nombre de partitions dans le topic. Nous vous recommandons de définir ce paramètre sur 6.
Storage Engine : le moteur de stockage utilisé par le topic. Ce paramètre doit être défini sur Local Storage.
RemarqueVous ne pouvez définir le paramètre Storage Engine sur Local Storage que si vous créez un topic sur une instance Professional Edition.
cleanup.policy : la politique de nettoyage des journaux utilisée par le topic. Ce paramètre doit être défini sur Compact.
etl-status-kafka
Dead-letter Queue Topic
Le topic utilisé pour stocker les données d'erreur du framework ETL. Ce paramètre s'affiche uniquement si vous définissez le paramètre Manual sur Resource Creation Method. Pour économiser les ressources de topic, ce topic peut être identique au Error Data Topic.
Topic : le nom du topic. Nous vous recommandons d'utiliser etl-error comme préfixe du nom du topic.
Partitions : le nombre de partitions dans le topic. Nous vous recommandons de définir ce paramètre sur 6.
Storage Engine : le moteur de stockage utilisé par le topic. Ce paramètre peut être défini sur Local Storage ou Cloud Storage.
RemarqueVous ne pouvez définir le paramètre Storage Engine sur Local Storage que si vous créez un topic sur une instance Professional Edition.
etl-error-kafka
Error Data Topic
Le topic utilisé pour stocker les données d'erreur du récepteur (sink). Ce paramètre s'affiche uniquement si vous définissez le paramètre Manual sur Resource Creation Method. Pour économiser les ressources de topic, ce topic peut être identique au Dead-letter Queue Topic.
Topic : le nom du topic. Nous vous recommandons d'utiliser etl-error comme préfixe du nom du topic.
Partitions : le nombre de partitions dans le topic. Nous vous recommandons de définir ce paramètre sur 6.
Storage Engine : le moteur de stockage utilisé par le topic. Ce paramètre peut être défini sur Local Storage ou Cloud Storage.
RemarqueVous ne pouvez définir le paramètre Storage Engine sur Local Storage que si vous créez un topic sur une instance Professional Edition.
etl-error-kafka
-
À l'étape Configure Function, configurez les paramètres et cliquez sur Create.
Avant de cliquer sur Create, cliquez sur Test pour vérifier que la fonction fonctionne comme prévu.
Paramètre
Description
Exemple
Programming Language
Le langage dans lequel la fonction est écrite. Définissez ce paramètre sur Python 3.
Python3
Template
Le modèle de fonction fourni par le système. Après avoir sélectionné un modèle de fonction, le système remplit automatiquement le champ Code avec le code du modèle.
Add Prefix/Suffix
Code
Le code utilisé pour traiter le message. ApsaraMQ for Kafka fournit des modèles de fonctions que vous pouvez utiliser pour nettoyer et transformer les données. Vous pouvez modifier le code du modèle de fonction sélectionné en fonction de vos besoins métier.
RemarqueVous pouvez importer des modules Python selon vos besoins.
Le message dans le code est au format dictionnaire. Vous n'avez qu'à modifier la clé et la valeur.
Renvoyez le message traité. Si la fonction sert à filtrer les messages, renvoyez None.
def deal_message(message): for keyItem in message.keys(): if (keyItem == 'key'): message[keyItem] = message[keyItem] + "KeySurfix" continue if (keyItem == 'value'): message[keyItem] = message[keyItem] + "ValueSurfix" continue return messageMessage Key
La clé du message à traiter dans le topic source. Cliquez sur Test Code pour afficher le paramètre.
demo
Message Content
La valeur du message à traiter dans le topic source.
{"key": "test"}
Une fois la tâche ETL créée, vous pouvez la consulter sur la page ETL Tasks. Le système déploie automatiquement la tâche.
Envoyer un message de test
Une fois la tâche ETL déployée, envoyez un message de test au topic ApsaraMQ for Kafka source pour vérifier si les données peuvent être traitées par la fonction configurée et envoyées au topic de destination.
Sur la page ETL Tasks, localisez la tâche ETL que vous avez créée et cliquez sur Test dans la colonne Actions.
-
Dans le panneau Send Message, saisissez les informations du message de test et cliquez sur OK pour envoyer le message de test.
Dans le champ Message Key, saisissez la clé du message de test. Exemple : demo.
Dans le champ Message Content, saisissez le contenu du message de test. Exemple : {"key": "test"}.
-
Configurez le paramètre Send to Specified Partition pour indiquer si le message de test doit être envoyé à une partition spécifique.
Si vous souhaitez envoyer le message de test à une partition spécifique, cliquez sur Partition ID et saisissez un ID de partition dans le champ Yes. Exemple : 0. Pour savoir comment interroger les ID de partition, consultez la page Afficher l'état des partitions.
Si vous ne souhaitez pas envoyer le message de test à une partition spécifique, cliquez sur No.
Afficher les journaux de la fonction
Après avoir extrait et traité les données dans ApsaraMQ for Kafka, consultez les journaux de la fonction pour vérifier si le topic de destination a bien reçu les données traitées. Pour plus d'informations, consultez la page Configurer la fonctionnalité de journalisation.
Figure 1. Afficher les journaux de la fonction 
Afficher les détails d'une tâche ETL
Une fois la tâche ETL créée, affichez ses détails dans la console ApsaraMQ for Kafka.
Sur la page ETL Tasks, localisez la tâche ETL que vous avez créée et cliquez sur Details dans la colonne Actions.
Sur la page Task Details, affichez les détails de la tâche ETL.
Supprimer une tâche ETL
Si vous n'avez plus besoin d'une tâche ETL, supprimez-la dans la console ApsaraMQ for Kafka.
Sur la page ETL Tasks, localisez la tâche ETL à supprimer et cliquez sur Delete dans la colonne Actions.
Dans le message Notes qui s'affiche, cliquez sur OK.