Créez un connecteur sink Function Compute pour exporter les données d'un topic source de votre instance ApsaraMQ for Kafka vers une fonction Function Compute.
Prérequis
Avant de créer un connecteur sink Function Compute, assurez-vous que les conditions préalables suivantes sont remplies :
-
ApsaraMQ for Kafka
Activez la fonctionnalité de connecteur pour l'instance ApsaraMQ for Kafka. Pour plus d'informations, consultez la rubrique Activer un connecteur.
-
Créez un topic source pour l'instance ApsaraMQ for Kafka. Pour plus d'informations, consultez la rubrique Étape 1 : Créer un topic.
Cet exemple utilise un topic nommé fc-test-input.
-
Function Compute
-
Créez une fonction dans Function Compute. Pour plus d'informations, consultez la rubrique Créer une fonction.
ImportantLa fonction doit être une fonction événementielle.
Ce guide utilise une fonction événementielle nommée hello_world à titre d'exemple. La fonction se trouve dans le service guide-hello_world et s'exécute dans un environnement d'exécution Python. Le code exemple de la fonction est le suivant :
# -*- coding: utf-8 -*- import logging # To enable the initializer feature # Implement the initializer function as follows: # def initializer(context): # logger = logging.getLogger() # logger.info('initializing') def handler(event, context): logger = logging.getLogger() logger.info('hello world:' + bytes.decode(event)) return 'hello world:' + bytes.decode(event)
-
-
Facultatif : EventBridge
RemarqueCette étape n'est requise que si votre instance ApsaraMQ for Kafka se trouve dans la région Chine (Hangzhou) ou Chine (Chengdu).
Notes
Vous ne pouvez exporter les données d'un topic source d'une instance ApsaraMQ for Kafka vers une fonction Function Compute que si elles se trouvent dans la même région. Pour plus d'informations sur les limites des connecteurs, consultez la rubrique Limites.
-
Si votre instance ApsaraMQ for Kafka se trouve dans la région Chine (Hangzhou) ou Chine (Chengdu), le connecteur est déployé sur EventBridge.
EventBridge est actuellement gratuit. Pour plus d'informations, consultez la rubrique Facturation.
-
Lorsque vous créez un connecteur, EventBridge crée automatiquement les rôles liés au service suivants : AliyunServiceRoleForEventBridgeSourceKafka et AliyunServiceRoleForEventBridgeConnectVPC.
Si un rôle lié au service n'a pas été créé, EventBridge en créera automatiquement un correspondant afin qu'EventBridge puisse utiliser ce rôle pour accéder à ApsaraMQ for Kafka et au VPC.
Si ces rôles liés au service existent déjà, EventBridge ne les recrée pas.
Pour plus d'informations sur les rôles liés au service, consultez la rubrique Rôles liés au service.
Vous ne pouvez pas actuellement afficher les journaux d'exécution des tâches déployées sur EventBridge. Une fois la tâche du connecteur terminée, vérifiez sa progression en consultant l'état de consommation du consumer group du topic source. Pour plus de détails, consultez la rubrique Afficher l'état du consommateur.
Procédure
Utilisez un connecteur sink Function Compute pour exporter les données d'un topic source d'une instance ApsaraMQ for Kafka vers une fonction Function Compute :
-
Facultatif : Activez l'accès inter-régions pour le connecteur sink Function Compute
ImportantSi vous n'avez pas besoin d'accès inter-régions, ignorez cette étape.
Activer l'accès Internet pour le connecteur sink Function Compute
-
Facultatif : Activez l'accès inter-comptes pour le connecteur sink Function Compute
ImportantSi vous n'avez pas besoin d'accès inter-comptes, ignorez cette étape.
-
Facultatif : Créez les topics et le consumer group requis par le connecteur sink Function Compute
ImportantSi vous n'avez pas besoin de personnaliser les noms des topics et du consumer group, vous pouvez ignorer cette étape.
Certains topics requis par un connecteur sink Function Compute doivent utiliser le moteur de stockage local. Si votre instance ApsaraMQ for Kafka est en version majeure 0.10.2, vous ne pouvez pas créer manuellement des topics avec un stockage local. Ces topics doivent être créés automatiquement.
-
Vérifiez les résultats
Activer l'accès Internet pour un connecteur sink FC
Si un connecteur sink Function Compute doit accéder à des services Alibaba Cloud dans d'autres régions, vous devez activer l'accès Internet pour celui-ci. Pour plus de détails, consultez la rubrique Activer l'accès Internet pour un connecteur.
Créer une stratégie personnalisée
Dans le compte cible, créez une stratégie personnalisée pour accorder l'accès à Function Compute.
Connectez-vous à la console RAM.
Dans le volet de navigation de gauche, sélectionnez Permissions > Policies.
Sur la page Policies, cliquez sur Create Policy.
-
Sur la page Create Policy, configurez la stratégie.
-
Sous l'onglet JSON, saisissez le script de stratégie et cliquez sur Next.
Le script de stratégie suivant accorde les autorisations d'accès à Function Compute :
{ "Version": "1", "Statement": [ { "Action": [ "fc:InvokeFunction", "fc:GetFunction" ], "Resource": "*", "Effect": "Allow" } ] } Sous Basic Information, saisissez KafkaConnectorFcAccess pour le champ Name.
Cliquez sur OK.
-
Créer un rôle RAM
Créez un rôle RAM dans le compte cible. Vous ne pouvez pas sélectionner ApsaraMQ for Kafka comme service de confiance lors de la création d'un rôle RAM. Par conséquent, vous devez sélectionner un autre service pris en charge, puis modifier manuellement la stratégie d'approbation après la création du rôle.
Dans le volet de navigation de gauche, sélectionnez Identity Management > Roles.
Sur la page Roles, cliquez sur Create Role.
-
Dans le panneau Create Role, configurez le rôle.
Sélectionnez Alibaba Cloud Service pour le type d'entité de confiance, puis cliquez sur Next.
Dans la section Role Type, sélectionnez Normal Service Role. Pour Role Name, saisissez AliyunKafkaConnectorRole. Dans la liste déroulante Select Trusted Service, sélectionnez Function Compute, puis cliquez sur Complete.
Sur la page Roles, recherchez et cliquez sur AliyunKafkaConnectorRole.
Sur la page de détails AliyunKafkaConnectorRole, cliquez sur l'onglet Trust Policy Management, puis sur Edit Trust Policy.
-
Dans le panneau Edit Trust Policy, remplacez fc par alikafka dans le script, puis cliquez sur OK.
Après avoir enregistré les modifications, sous l'onglet Trust Policy Management pour AliyunKafkaConnectorRole, confirmez que le paramètre
Servicede la stratégie d'approbation est mis à jour versalikafka.aliyuncs.comet que le paramètreActionest défini sursts:AssumeRole.
Ajouter des autorisations
Dans le compte cible, accordez au rôle RAM les autorisations d'accès à Function Compute.
Dans le volet de navigation de gauche, accédez à Identity Management > Roles.
Sur la page Roles, recherchez AliyunKafkaConnectorRole et cliquez sur Add Permissions dans la colonne Actions.
-
Dans le panneau Add Permissions, ajoutez la stratégie KafkaConnectorFcAccess.
Dans la section Select Policy, sélectionnez Custom Policy.
Dans la liste Authorization Policy Name, recherchez et cliquez sur KafkaConnectorFcAccess.
Cliquez sur OK.
Cliquez sur Complete.
Créer des topics pour le connecteur sink Function Compute
Dans la console ApsaraMQ for Kafka, vous pouvez créer manuellement les cinq topics requis par un connecteur sink Function Compute : un topic d'offset de tâche, un topic de configuration de tâche, un topic d'état de tâche, un topic de file d'attente de lettres mortes et un topic de données d'erreur. Ces topics ont des exigences différentes concernant le nombre de partitions et le moteur de stockage. Pour plus d'informations, consultez la rubrique Paramètres de l'étape Configurer le service source.
Connectez-vous à la console ApsaraMQ for Kafka.
-
Sur la page Overview, sélectionnez une région dans la section Resource Distribution.
ImportantVous devez créer des topics dans la même région que votre application, c'est-à-dire la région où l'instance ECS est déployée. Les topics ne peuvent pas être utilisés entre différentes régions. Par exemple, si un topic est créé dans la région Chine (Pékin), le producteur et le consommateur de messages doivent également s'exécuter sur une instance ECS dans la région Chine (Pékin).
Sur la page Instances, cliquez sur le nom de l'instance cible.
Dans le volet de navigation de gauche, cliquez sur Topics.
Sur la page Topics, cliquez sur Create Topic.
-
Dans le panneau Create Topic, configurez les paramètres du topic et cliquez sur OK.
Paramètre
Description
Exemple
Name
Le nom du topic.
RemarqueKafka considère les noms de topics contenant des traits de soulignement (
xxx_xxx) et des points (xxx.xxx) comme identiques. Si vous tentez de créer un topic en double, le système renvoie une erreur.demo
Description
Une brève description du topic.
test demo
Partitions
Le nombre de partitions du topic.
12
Storage Engine
RemarqueActuellement, vous pouvez sélectionner un moteur de stockage uniquement pour les instances Professional Edition non serverless. Pour les autres types d'instances, cette option n'est pas disponible et le Stockage cloud est utilisé par défaut.
Le moteur de stockage pour les messages du topic.
ApsaraMQ for Kafka prend en charge les deux moteurs de stockage suivants :
-
Cloud Storage : Ce moteur utilise des disques Alibaba Cloud pour le stockage sous-jacent et offre des performances élevées, une faible latence et une grande fiabilité grâce à un mécanisme distribué à trois réplicas. Si l'Instance Edition de l'instance est Standard (High Write), seul le Cloud Storage peut être utilisé.
-
Local Storage : Ce moteur utilise l'algorithme de réplication ISR (In-Sync Replica) natif de Kafka et un mécanisme distribué à trois réplicas.
Cloud Storage
Message Type
Le type de messages dans le topic.
-
Normal Message : Par défaut, Kafka distribue les messages ayant la même clé vers la même partition et les stocke dans l'ordre d'envoi. Si un nœud du cluster tombe en panne, les messages peuvent être désordonnés. Si vous définissez le Storage Engine sur Cloud Storage, le système sélectionne Normal Message par défaut.
-
Partitionally Ordered Message : Par défaut, Kafka distribue les messages ayant la même clé vers la même partition et les stocke dans l'ordre d'envoi. Même si un nœud du cluster tombe en panne, l'ordre des messages au sein de la partition est garanti. Cependant, l'envoi de messages vers certaines partitions peut échouer jusqu'à leur récupération. Si vous définissez le Storage Engine sur Local Storage, le système sélectionne Partitionally Ordered Message par défaut.
Normal Message
Log Cleanup Policy
La politique de nettoyage des journaux pour le topic.
Lorsque vous sélectionnez Local Storage comme Storage Engine (actuellement, seules les instances Professional Edition prennent en charge le stockage local et cette option n'est pas disponible pour les instances Standard Edition), vous devez configurer la Log Cleanup Policy.
ApsaraMQ for Kafka prend en charge les deux politiques de nettoyage des journaux suivantes.
-
Delete : Politique de nettoyage des messages par défaut. Si l'espace disque est suffisant, les messages sont conservés pendant la période de rétention spécifiée. Si l'espace disque est insuffisant (généralement lorsque l'utilisation du disque dépasse 85 %), le système supprime les messages plus anciens prématurément pour garantir la disponibilité du service.
-
Compact : Utilise la politique de nettoyage Kafka Log Compaction. La compaction des journaux garantit que le système conserve la dernière valeur pour chaque clé de message. Elle est principalement utilisée pour des scénarios tels que la restauration de l'état après un crash système ou le rechargement d'un cache après un redémarrage du système. Par exemple, Kafka Connect et Confluent Schema Registry utilisent des topics compactés pour stocker l'état du système et les données de configuration.
ImportantLes topics compactés sont généralement utilisés uniquement pour des composants spécifiques de l'écosystème, tels que Kafka Connect ou Confluent Schema Registry. Ne définissez pas cette propriété pour les topics utilisés pour la production et la consommation générales de messages. Pour plus d'informations, consultez la Bibliothèque de démonstration ApsaraMQ for Kafka.
Compact
Tag
Les tags pour le topic.
demo
Après avoir créé le topic, il apparaît dans la liste des topics sur la page Topics.
-
Créer un consumer group pour le connecteur sink FC
Vous pouvez créer manuellement le consumer group pour une tâche de synchronisation de données de connecteur sink Function Compute dans la console ApsaraMQ for Kafka. Le nom du consumer group doit être connect-nom-tâche. Pour plus d'informations, consultez la rubrique Paramètres de l'étape Configurer le service source.
Connectez-vous à la console ApsaraMQ for Kafka.
Sur la page Overview, sélectionnez une région dans la section Resource Distribution.
Sur la page Instances, cliquez sur le nom de l'instance cible.
Dans le volet de navigation de gauche, cliquez sur Groups.
Sur la page Groups, cliquez sur Create Group.
-
Dans le panneau Create Group, saisissez un nom pour le consumer group dans la zone de texte Group ID, saisissez une brève description dans la zone de texte Description, ajoutez des tags au consumer group, puis cliquez sur OK.
Une fois créé, le consumer group apparaît dans la liste sur la page Groups.
Création et déploiement d'un connecteur FC Sink
Créez et déployez un connecteur FC Sink pour synchroniser les données depuis ApsaraMQ for Kafka vers Function Compute.
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, cliquez sur Connectors.
Sur la page Connectors, sélectionnez l'instance à laquelle le connecteur est associé dans la liste déroulante Select Instance et cliquez sur Create Connector.
-
Dans l'assistant Create Connector, effectuez les étapes suivantes.
-
Sous l'onglet Configure Basic Information, configurez les paramètres ci-dessous selon vos besoins, puis cliquez sur Next.
Parameter
Description
Example
Name
Le nom du connecteur. Le nom doit respecter les exigences suivantes :
-
La longueur maximale est de 48 caractères. Seuls les chiffres, les lettres minuscules et les traits d'union (-) sont autorisés. Le nom ne peut pas commencer par un trait d'union.
-
Le nom doit être unique au sein d'une instance ApsaraMQ for Kafka.
La tâche de synchronisation des données du connecteur utilise un consumer group nommé
connect-task-name. Si vous ne créez pas ce consumer group manuellement, le système le génère automatiquement.kafka-fc-sink
Instance
Par défaut, le nom et l'ID de l'instance s'affichent.
demo alikafka_post-cn-st21p8vj****
-
-
Sous l'onglet Configure Source Service, définissez le paramètre Data Source sur Message Queue for Apache Kafka, configurez les paramètres suivants, puis cliquez sur Next.
RemarqueSi vous avez déjà créé un topic et un consumer group, choisissez la création manuelle des ressources et saisissez les informations relatives à vos ressources existantes. Sinon, optez pour la création automatique des ressources.
Tableau 1. Paramètres de configuration du service source
Parameter
Description
Example
Data Source Topic
Le topic source à partir duquel les données sont synchronisées.
fc-test-input
Consumer Thread Concurrency
Le nombre de threads de consommateur simultanés pour le topic source. Valeur par défaut : 6. Valeurs possibles :
-
1
-
2
-
3
-
6
-
12
6
Consumer Offset
L'offset à partir duquel la consommation démarre. Valeurs possibles :
-
Earliest Offset : démarrer la consommation à partir de l'offset le plus ancien.
-
Latest Offset : démarrer la consommation à partir de l'offset le plus récent.
Earliest Offset
VPC ID
L'ID du VPC dans lequel la tâche de synchronisation des données s'exécute. Ce paramètre s'affiche lorsque vous cliquez sur Configure Runtime Environment. Par défaut, il s'agit du VPC de votre instance ApsaraMQ for Kafka. Vous n'avez pas besoin de spécifier de valeur.
vpc-bp1xpdnd3l***
vSwitch ID
L'ID du vSwitch dans lequel la tâche de synchronisation des données s'exécute. Ce paramètre s'affiche lorsque vous cliquez sur Configure Runtime Environment. Le vSwitch doit se trouver dans le même VPC que l'instance ApsaraMQ for Kafka. Par défaut, il s'agit du vSwitch que vous avez spécifié lors du déploiement de l'instance ApsaraMQ for Kafka.
vsw-bp1d2jgg81***
Failure Handling Policy
Contrôle le comportement en cas d'échec de livraison d'un message provenant d'une partition. Ce paramètre s'affiche lorsque vous cliquez sur Configure Runtime Environment. Valeurs possibles :
-
Continue Subscription : poursuivre la consommation depuis la partition et enregistrer l'erreur.
-
Stop Subscription : arrêter la consommation depuis la partition et enregistrer l'erreur.
Remarque-
Pour savoir comment consulter les journaux, reportez-vous à la rubrique Opérations sur les connecteurs.
-
Pour identifier les solutions en fonction des codes d'erreur, consultez la section Codes d'erreur.
Continue Subscription
Resource Creation Method
La méthode utilisée pour créer le topic et le consumer group requis par le connecteur. Ce paramètre s'affiche lorsque vous cliquez sur Configure Runtime Environment.
-
Auto
-
Manual
Auto
Connector Consumer Group
Le consumer group utilisé par la tâche de synchronisation des données du connecteur. Ce paramètre s'affiche lorsque vous cliquez sur Configure Runtime Environment. Le nom du consumer group doit suivre le format connect-task-name.
connect-kafka-fc-sink
Task Offset Topic
Le topic utilisé pour stocker les offsets des consommateurs. Ce paramètre s'affiche lorsque vous cliquez sur Configure Runtime Environment.
-
Le nom du topic doit commencer par
connect-offset. -
Partitions : le nombre de partitions du topic doit être supérieur à 1.
-
Moteur de stockage : le moteur de stockage du topic doit être Local Storage.
-
cleanup.policy : la politique de nettoyage des journaux du topic doit être
compact.
connect-offset-kafka-fc-sink
Task Configuration Topic
Le topic utilisé pour stocker les configurations des tâches. Ce paramètre s'affiche lorsque vous cliquez sur Configure Runtime Environment.
-
Topic : nous vous recommandons de faire commencer le nom du topic par
connect-config. -
Partitions : le nombre de partitions du topic doit être égal à 1.
-
Moteur de stockage : le moteur de stockage du topic doit être Local Storage.
-
cleanup.policy : la politique de nettoyage des journaux du topic doit être
compact.
connect-config-kafka-fc-sink
Task Status Topic
Le topic utilisé pour stocker l'état des tâches. Ce paramètre s'affiche lorsque vous cliquez sur Configure Runtime Environment.
-
Topic : nous vous recommandons de faire commencer le nom du topic par
connect-status. -
Partitions : nous vous recommandons de définir le nombre de partitions sur 6.
-
Moteur de stockage : le moteur de stockage du topic doit être Local Storage.
-
cleanup.policy: la politique de nettoyage des journaux du topic doit être compact.
connect-status-kafka-fc-sink
Dead-letter Queue Topic
Le topic utilisé pour stocker les données d'erreur du framework Connect. Ce paramètre s'affiche lorsque vous cliquez sur Configure Runtime Environment. Pour économiser les ressources de topic, vous pouvez utiliser le même topic pour ce paramètre et pour le paramètre Error data topic.
-
Topic : nous vous recommandons de faire commencer le nom du topic par
connect-error. -
Partitions : nous vous recommandons de définir le nombre de partitions sur 6.
-
Moteur de stockage : le moteur de stockage du topic peut être Local Storage ou Cloud Storage.
connect-error-kafka-fc-sink
Error Data Topic
Le topic utilisé pour stocker les données d'erreur du sink. Ce paramètre s'affiche lorsque vous cliquez sur Configure Runtime Environment. Pour économiser les ressources de topic, vous pouvez utiliser le même topic pour ce paramètre et pour le paramètre dead-letter queue topic.
-
Topic : nous vous recommandons de faire commencer le nom du topic par
connect-error. -
Partitions : nous vous recommandons de définir le nombre de partitions sur 6.
-
Moteur de stockage : le moteur de stockage du topic peut être Local Storage ou Cloud Storage.
connect-error-kafka-fc-sink
-
-
Sous l'onglet Configure Destination Service, sélectionnez Function Compute pour le paramètre Target Service, configurez les paramètres suivants, puis cliquez sur Create.
RemarqueSi l'instance du connecteur se trouve dans la région Chine (Hangzhou) ou Chine (Chengdu), une boîte de dialogue Service Authorization s'affiche pour créer les rôles liés au service AliyunServiceRoleForEventBridgeSourceKafka et AliyunServiceRoleForEventBridgeConnectVPC lorsque vous sélectionnez Function Compute comme Target Service. Dans la boîte de dialogue Service Authorization, cliquez sur Confirm, puis configurez les paramètres suivants et cliquez sur Create. Si les rôles liés au service ont déjà été créés, la boîte de dialogue Service Authorization ne s'affiche pas.
Parameter
Description
Example
Cross-account/Cross-region
Indique si le connecteur FC Sink synchronise les données vers un service Function Compute situé dans un autre compte ou une autre région. Valeur par défaut : No. Valeurs possibles :
-
No : mode même compte, même région.
-
Yes : comptes croisés, régions croisées, ou les deux.
No
Region
La région du service Function Compute. Par défaut, il s'agit de la région du connecteur FC Sink. Pour synchroniser des données entre différentes régions, vous devez activer l'accès public pour le connecteur, puis sélectionner la région de destination. Pour plus d'informations, consultez la section Activation de l'accès public pour un connecteur FC Sink.
ImportantLorsque le paramètre Cross-account/Cross-region est défini sur Yes, le paramètre Region s'affiche.
cn-hangzhou
Service Endpoint
Le endpoint du service Function Compute. Vous pouvez obtenir le endpoint depuis la section Common Info de la page Overview de la console Function Compute.
-
Endpoint interne : recommandé pour une latence réduite. Utilisez ce type de endpoint si votre instance ApsaraMQ for Kafka et votre service Function Compute se trouvent dans la même région.
-
Endpoint public : non recommandé en raison d'une latence plus élevée. Utilisez ce type de endpoint si votre instance ApsaraMQ for Kafka et votre service Function Compute se trouvent dans des régions différentes. Pour utiliser un endpoint public, vous devez activer l'accès public pour le connecteur. Pour plus d'informations, consultez la section Activation de l'accès public pour un connecteur FC Sink.
ImportantLorsque le paramètre Cross-account/Cross-region est défini sur Yes, le paramètre Service Endpoint s'affiche.
http://188***.cn-hangzhou.fc.aliyuncs.com
Alibaba Cloud Account
L'ID du compte Alibaba Cloud auquel appartient le service Function Compute. Vous pouvez obtenir l'ID depuis la section Common Info de la page Overview de la console Function Compute.
ImportantSi le paramètre Cross-account/Cross-region est défini sur Yes, le paramètre Alibaba Cloud Account s'affiche.
188***
RAM Role Name
Le nom du RAM role qu'ApsaraMQ for Kafka assume pour accéder au service Function Compute.
-
Pour un accès au sein du même compte, vous devez créer un rôle RAM dans votre compte, lui accorder des autorisations, puis saisir le nom du rôle. Pour plus d'informations, consultez les rubriques Création d'une stratégie personnalisée, Création d'un rôle RAM pour un service Alibaba Cloud de confiance et Attribution d'autorisations à un rôle RAM.
-
Pour un accès entre comptes différents, vous devez créer un rôle RAM dans le compte de destination, lui accorder des autorisations, puis saisir le nom du rôle. Pour plus d'informations, consultez les rubriques Création d'une stratégie personnalisée, Création d'un rôle RAM pour un service Alibaba Cloud de confiance et Attribution d'autorisations à un rôle RAM.
ImportantLorsque le paramètre Cross-account/Cross-region est défini sur Yes, le paramètre RAM Role Name s'affiche.
AliyunKafkaConnectorRole
Service Name
Le nom du service dans Function Compute.
guide-hello_world
Function Name
Le nom de la fonction dans le service Function Compute.
hello_world
Version or Alias
La version ou l'alias du service Function Compute.
Important-
Si le paramètre Cross-account/Cross-region est défini sur No, vous devez sélectionner soit Specified Version, soit Specified Alias.
-
Si le paramètre Cross-account/Cross-region est défini sur Yes, vous devez saisir manuellement une version ou un alias de service.
LATEST
Service Version
La version du service Function Compute.
ImportantSi le paramètre Cross-account/Cross-region est défini sur No et que le paramètre Version or Alias est défini sur Specified Version, le paramètre Service Version s'affiche.
LATEST
Service Alias
L'alias du service Function Compute.
ImportantLorsque le paramètre Cross-account/Cross-region est défini sur No et que le paramètre Version or Alias est défini sur Specified Alias, le paramètre Service Alias s'affiche.
jy
Transmission Mode
Le mode d'envoi des messages. Valeurs possibles :
-
Asynchronous : recommandé.
-
Synchronous : non recommandé. Dans ce mode, un traitement lent des messages par Function Compute ralentit également ApsaraMQ for Kafka. Si le traitement d'un lot de messages prend plus de 5 minutes, un rééquilibrage client est déclenché dans ApsaraMQ for Kafka.
Asynchronous
Data Size
Le nombre maximal de messages à inclure dans un seul lot. Le connecteur regroupe les messages en lots qui respectent à la fois ce nombre et les limites de taille de requête sous-jacentes (6 Mo pour le mode synchrone, 128 Ko pour le mode asynchrone). Par exemple, si le mode de livraison est asynchrone, la taille du lot est de 20 et que vous souhaitez envoyer 18 messages dont 17 ont une taille totale de 127 Ko et un message a une taille de 200 Ko, le connecteur agrège et envoie les 17 messages dans un lot. Le connecteur envoie le message restant, dont la taille dépasse 128 Ko, dans un lot séparé.
RemarqueSi vous définissez la key sur null lors de l'envoi d'un message, la requête n'inclut pas la key. Si vous définissez la value sur null, la requête n'inclut pas la value.
-
Si la taille totale des messages dans un lot ne dépasse pas la limite de taille de requête, la requête contient le contenu du message. L'exemple de code suivant illustre une requête :
[ { "key":"this is the message's key2", "offset":8, "overflowFlag":false, "partition":4, "timestamp":1603785325438, "topic":"Test", "value":"this is the message's value2", "valueSize":28 }, { "key":"this is the message's key9", "offset":9, "overflowFlag":false, "partition":4, "timestamp":1603785325440, "topic":"Test", "value":"this is the message's value9", "valueSize":28 }, { "key":"this is the message's key12", "offset":10, "overflowFlag":false, "partition":4, "timestamp":1603785325442, "topic":"Test", "value":"this is the message's value12", "valueSize":29 }, { "key":"this is the message's key38", "offset":11, "overflowFlag":false, "partition":4, "timestamp":1603785325464, "topic":"Test", "value":"this is the message's value38", "valueSize":29 } ] -
Si la taille d'un seul message dépasse la limite de taille de requête, la requête n'inclut pas le contenu du message. L'exemple de code suivant illustre une requête :
[ { "key":"123", "offset":4, "overflowFlag":true, "partition":0, "timestamp":1603779578478, "topic":"Test", "value":"1", "valueSize":272687 } ]RemarquePour obtenir le contenu du message, vous devez extraire le message en fonction de son offset.
50
Retries
Le nombre de tentatives après l'échec de l'envoi d'un message. La valeur par défaut est 2. La plage de valeurs est comprise entre 1 et 3. Certaines erreurs provoquant l'échec de l'envoi des messages ne prennent pas en charge les nouvelles tentatives. La correspondance entre les codes d'erreur et la prise en charge des nouvelles tentatives est la suivante :
-
4XX : les nouvelles tentatives ne sont pas prises en charge pour les erreurs 4xx, à l'exception de l'erreur 429.
-
5XX : les nouvelles tentatives sont prises en charge.
Remarque-
Le connecteur appelle l'opération InvokeFunction pour envoyer des messages à Function Compute.
-
Si l'envoi d'un message échoue après le nombre maximal de tentatives, il est envoyé vers le topic de file d'attente des messages morts (dead-letter queue). Les messages présents dans un topic de file d'attente des messages morts ne déclenchent pas à nouveau les tâches du connecteur Function Compute. Nous vous recommandons de configurer des alertes pour le topic de file d'attente des messages morts afin de surveiller son état en temps réel et de traiter les exceptions rapidement.
2
Une fois le connecteur créé, vous pouvez le consulter sur la page Connectors.
-
-
-
Sur la page Connectors, localisez votre connecteur nouvellement créé et cliquez sur Deploy dans la colonne Actions.
Pour configurer les ressources Function Compute, choisissez dans la colonne Actions. Vous êtes redirigé vers la console Function Compute pour finaliser la configuration.
Envoi d'un message de test
Après avoir déployé le connecteur sink Function Compute, envoyez un message au topic source de votre instance ApsaraMQ for Kafka pour vérifier que les données sont bien synchronisées vers Function Compute.
Sur la page Connectors, localisez le connecteur cible et cliquez sur Test dans la colonne Actions.
-
Dans le panneau Send Message, envoyez un message de test.
-
Pour le paramètre Sending Method, sélectionnez Console.
Dans le champ Message Key, saisissez une clé de message, par exemple
demo.Dans le champ Message Content, saisissez le contenu du message, par exemple
{"key": "test"}.-
Pour le paramètre Send to Specified Partition, choisissez une option :
Cliquez sur Yes et saisissez un ID de partition, par exemple
0, dans le champ Partition ID. Pour trouver l'ID de partition, consultez la rubrique Affichage de l'état des partitions.Cliquez sur No pour envoyer le message sans spécifier de partition.
Pour le paramètre Sending Method, sélectionnez Docker et exécutez la commande figurant dans la section Run the Docker container to produce a sample message.
Pour le paramètre Sending Method, sélectionnez SDK. Choisissez ensuite le SDK et la méthode d'intégration correspondant à votre langage ou framework préféré pour envoyer un message.
-
Journaux de fonction
Après avoir envoyé un message au topic source de votre instance ApsaraMQ for Kafka, vérifiez les journaux de la fonction pour confirmer qu'elle a bien reçu le message. Pour plus d'informations, consultez la rubrique Configuration des journaux.
Votre message de test apparaît dans les journaux.
Sur la page des détails de la fonction, cliquez sur l'onglet Log query et sélectionnez Advanced query. Sélectionnez function-log comme Logstore. Le champ message dans les journaux contient le message issu du topic Kafka, par exemple [INFO] hello world:[{"key":"1","offset":1,"partition":0,"timestamp":1605598174308,"topic":"fc-test-input","value":"1"}]. Cela confirme que la fonction a bien reçu le message depuis le déclencheur Kafka.