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
Vous avez activé Tablestore et créé une instance. Pour plus d'informations, consultez la rubrique Activer Tablestore et créer une instance.
Vous avez acheté et activé une instance ApsaraMQ for Kafka et créé un sujet. Pour connaître la procédure détaillée, consultez les rubriques Acheter une instance ApsaraMQ for Kafka et Créer des ressources.
Le rôle lié au service généré par la tâche du connecteur de réception Tablestore nécessite la stratégie
AliyunOTSFullAccess. Attachez manuellement cette stratégie pour accorder au rôle l'autorisation de gérer Tablestore. Pour plus d'informations, consultez la rubrique Accorder des autorisations à un rôle RAM.
É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
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.
-
Sur la page Create Task, définissez les champs Task Name et Description, configurez les paramètres suivants, puis cliquez sur Save.
-
Create Task
-
À 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_xxxest 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
-
À 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.
À 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.
-
À 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.
-
Revenez à la page Tasks. Recherchez votre tâche et cliquez sur Enable dans la colonne Actions.
-
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
Sur la page Tasks, dans la colonne Event Source correspondant à votre tâche, cliquez sur le sujet source.
Sur la page des détails du sujet, cliquez sur Send Test Message.
-
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 surots-sink-v1et Send to Specified Partition sur No. Revenez à la page Tasks et, dans la colonne Event Target correspondant à votre tâche, cliquez sur le nom de la table de destination.
-
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-k1et value=ots-sink-v1confirme que le message a bien été écrit dans Tablestore.