Tous les produits
Search
Centre de documentation

AnalyticDB:Synchroniser les données Kafka à l'aide de la fonctionnalité de synchronisation des données APS (Recommandé)

Dernière mise à jour :Aug 10, 2026

La fonctionnalité de synchronisation des données d'AnalyticDB Pipeline Service (APS) vous permet d'ingérer les messages ApsaraMQ for Kafka dans AnalyticDB for MySQL en quasi-temps réel, à partir du décalage (offset) de votre choix. Cette fonctionnalité prend en charge la sortie des données en quasi-temps réel, l'archivage complet des données historiques et l'analyse élastique. Une fois une tâche de synchronisation démarrée, les données sont validées toutes les 5 minutes par défaut ; le premier lot de données ingérées est donc interrogeable après environ 5 minutes.

Seuls les messages Kafka au format JSON sont pris en charge.

Limites

  • Seuls les messages Kafka au format JSON sont pris en charge. Les autres formats provoquent des erreurs de synchronisation.

  • Les modifications du schéma de table Kafka ne sont pas automatiquement propagées vers AnalyticDB for MySQL. Vous devez appliquer manuellement les modifications DDL.

  • Les exemples de données Kafka supérieurs à 8 Ko sont tronqués par l'API Kafka, ce qui empêche le système d'analyser l'exemple et de générer automatiquement les mappages de champs.

  • Les données écrites dans AnalyticDB for MySQL ne sont visibles qu'après l'exécution d'une opération de validation (Commit). L'intervalle de validation par défaut étant de 5 minutes, attendez au moins 5 minutes après le démarrage d'une tâche avant d'interroger le premier lot de données.

  • Les chemins OSS utilisés dans différentes tâches de synchronisation ne peuvent pas partager le même préfixe. Par exemple, oss://testBucketName/test/sls1/ et oss://testBucketName/test/ entrent en conflit et entraînent l'écrasement des données.

  • Si une tâche de synchronisation échoue et que les données du topic Kafka ont expiré, les données effacées ne peuvent pas être récupérées lors du redémarrage de la tâche. Pour réduire ce risque, augmentez la période de rétention des données du topic. En cas d'échec d'une tâche, contactez rapidement le support technique.

Prérequis

Avant de commencer, assurez-vous de disposer des éléments suivants :

  • Un cluster AnalyticDB for MySQL Enterprise Edition, Basic Edition ou Data Lakehouse Edition

  • Un groupe de ressources pour les jobs

  • Un compte de base de données pour le cluster (voir le tableau ci-dessous)

  • Une instance ApsaraMQ for Kafka située dans la même région que le cluster AnalyticDB for MySQL

  • Un topic Kafka contenant déjà des messages envoyés (voir Démarrage rapide pour ApsaraMQ for Kafka)

Exigences relatives aux comptes de base de données selon le type de compte :

Type de compte Comptes requis Étapes supplémentaires
Compte Alibaba Cloud Compte privilégié uniquement Aucune
Utilisateur Resource Access Management (RAM) Compte privilégié + compte standard Associer le compte standard à l'utilisateur RAM

Facturation

L'utilisation de la fonctionnalité de synchronisation des données entraîne les frais suivants :

Fonctionnement

  1. Ajoutez une source de données Kafka pour identifier l'instance et le topic ApsaraMQ for Kafka à lire.

  2. Créez un lien de données (job de synchronisation) qui mappe les messages Kafka à une table AnalyticDB for MySQL, configure l'analyse JSON, les clés de partition et le décalage du consommateur.

  3. Démarrez la tâche. Celle-ci commence à consommer les données Kafka à partir du décalage sélectionné et les écrit dans la table de destination.

  4. Interrogez les données ingérées à l'aide de Spark SQL.

Étape 1 : Créer une source de données

Ignorez cette étape si vous avez déjà ajouté une source de données Kafka . Accédez directement à Étape 2 : Créer un lien de données .
  1. Connectez-vous à la console AnalyticDB for MySQL. Dans le coin supérieur gauche, sélectionnez une région. Dans le volet de navigation de gauche, cliquez sur Clusters, puis cliquez sur l'ID du cluster.

  2. Dans le volet de navigation de gauche, choisissez Data Ingestion > Data Sources.

  3. Cliquez sur Create Data Source.

  4. Sur la page Create Data Source, configurez les paramètres suivants :

    Paramètre Description
    Data Source Type Sélectionnez Kafka.
    Data Source Name Généré automatiquement à partir du type de source et de l'heure actuelle. Modifiez-le si nécessaire.
    Data Source Description Saisissez une description, telle que le scénario métier ou la portée des données.
    Deployment Mode Seul le mode Alibaba Cloud Instance est pris en charge.
    Kafka Instance L'ID de l'instance Kafka. Retrouvez-le sur la page Instances de la console ApsaraMQ for Kafka.
    Kafka Topic Le nom du topic. Retrouvez-le sur la page Topics de l'instance de destination dans la console ApsaraMQ for Kafka.
    Message Data Format Seul le format JSON est pris en charge.
  5. Cliquez sur Create.

Étape 2 : Créer un lien de données

  1. Dans le volet de navigation de gauche, cliquez sur Simple Log Service/Kafka Data Synchronization.

  2. Cliquez sur Create Synchronization Job.

  3. Sur la page Create Synchronization Job, configurez les trois sections ci-dessous.

Paramètres de source et de destination

Paramètre Description
Job Name Généré automatiquement à partir du type de source et de l'heure actuelle. Modifiez-le si nécessaire.
Data Source Sélectionnez une source de données Kafka existante ou créez-en une nouvelle.
Destination Type Choisissez l'emplacement de stockage des données synchronisées. Consultez Choisir un type de destination ci-dessous.
ADB Lake Storage Le stockage lac pour les données lac d'AnalyticDB for MySQL. Sélectionnez-le dans la liste déroulante ou cliquez sur Automatically Created pour en créer un. Requis lorsque le paramètre Destination Type est défini sur Data Lake - AnalyticDB Lake Storage.
OSS Path Le chemin de stockage OSS pour les données lac d'AnalyticDB for MySQL. Requis lorsque le paramètre Destination Type est défini sur Data Lake - User OSS. Sélectionnez un dossier vide : le chemin ne peut pas être modifié après la création et ne doit pas partager de préfixe avec le chemin OSS d'une autre tâche de synchronisation.
Storage Format PAIMON (disponible uniquement lorsque le paramètre Destination Type est défini sur Data Lake - User OSS) ou ICEBERG.

Choisir un type de destination

Type de destination Cas d'utilisation Notes
Data Lake - AnalyticDB Lake Storage (recommandé) Configurations standards où vous souhaitez qu'AnalyticDB for MySQL gère le stockage Activez d'abord la fonctionnalité de stockage lac. Seul le format ICEBERG est pris en charge.
Data Lake - User OSS Lorsque vous devez utiliser votre propre bucket OSS Les formats PAIMON et ICEBERG sont tous deux pris en charge.

Paramètres de base de données et de table de destination

Paramètre Description
Database Name La base de données de destination dans AnalyticDB for MySQL. Une nouvelle base de données est créée si aucune base portant ce nom n'existe. Si elle existe déjà, les données y sont synchronisées. Consultez les Limites pour les conventions de nommage. Si le paramètre Storage Format est défini sur PAIMON, la base de données doit respecter les exigences listées ci-dessous.
Table Name La table de destination dans AnalyticDB for MySQL. Une nouvelle table est créée si aucune table portant ce nom n'existe. Si une table portant le même nom existe déjà, la tâche de synchronisation échoue. Consultez les Limites pour les conventions de nommage.
Sample Data Les dernières données récupérées automatiquement depuis le topic Kafka. Les données doivent être au format JSON.
Parsed JSON Layers Le nombre de niveaux JSON imbriqués à analyser. Valeurs valides : 0 (aucune analyse), 1 (par défaut), 2, 3, 4. Consultez Niveaux d'analyse JSON et inférence de schéma.
Schema Field Mapping Le schéma inféré à partir des données d'exemple après analyse. Ajustez les noms et types de champs de destination, ou ajoutez et supprimez des champs selon vos besoins.
Partition Key Settings (Facultatif) Une clé de partition pour la table de destination. Partitionnez par heure de journalisation ou par logique métier pour améliorer les performances d'ingestion et d'interrogation. Si ce champ est laissé vide, la table ne comporte aucune partition.

Exigences relatives aux bases de données PAIMON

Si le paramètre Storage Format est défini sur PAIMON, la base de données de destination doit remplir toutes les conditions suivantes. À défaut, la tâche de synchronisation échouera.

  • Il doit s'agir d'une base de données externe créée avec CREATE EXTERNAL DATABASE <database_name>.

  • Le paramètre DBPROPERTIES doit inclure catalog = paimon.

  • Le paramètre DBPROPERTIES doit inclure adb.paimon.warehouse. Exemple : adb.paimon.warehouse=oss://testBucketName/aps/data.

  • Le paramètre DBPROPERTIES doit inclure LOCATION avec un suffixe .db ajouté au nom de la base de données. Exemple : LOCATION=oss://testBucketName/aps/data/kafka_paimon_external_db.db/. Le répertoire du bucket dans ce chemin doit déjà exister et le chemin doit inclure .db après le nom de la base de données ; sinon, les requêtes XIHE échoueront.

Paramètres de synchronisation

Paramètre Description
Starting Consumer Offset for Incremental Synchronization Le point de départ dans Kafka à partir duquel la tâche commence à consommer les données. Earliest offset (begin_cursor) : consomme à partir des données les plus anciennes disponibles. Latest offset (end_cursor) : consomme uniquement à partir des données les plus récentes. Custom offset : consomme à partir du premier message postérieur ou égal à une heure spécifique que vous sélectionnez.
Job Resource Group Le groupe de ressources pour les jobs dans lequel la tâche s'exécute.
ACUs for Incremental Synchronization Le nombre d'ACU alloués depuis le groupe de ressources pour les jobs. Minimum : 2 ACU. Maximum : les ACU disponibles restants dans le groupe de ressources. Un nombre d'ACU plus élevé améliore les performances d'ingestion et la stabilité de la tâche.
Advanced Settings Configurations personnalisées. Contactez le support technique pour activer cette option.

Exemple de déduction d'ACU : Si un groupe de ressources pour les jobs dispose d'un maximum de 48 ACU et qu'une tâche existante utilise déjà 8 ACU, la nouvelle tâche peut utiliser au maximum 40 ACU.

  1. Cliquez sur Submit.

Étape 3 : Démarrer la tâche de synchronisation des données

  1. Sur la page Simple Log Service/Kafka Data Synchronization, localisez la tâche que vous avez créée et cliquez sur Start dans la colonne Actions.

  2. Cliquez sur Search. La tâche a démarré avec succès lorsque son statut passe à Running.

Le premier lot de données est interrogeable au moins 5 minutes après le démarrage de la tâche, car les données sont validées par intervalles de 5 minutes par défaut.

Étape 4 : Analyser les données

Une fois les données synchronisées, utilisez Spark SQL pour les interroger dans AnalyticDB for MySQL. Pour plus d'informations, consultez Éditeur de développement Spark et Développement d'applications Spark hors ligne.

  1. Dans le volet de navigation de gauche, choisissez Job Development > Spark JAR Development.

  2. Saisissez vos instructions Spark SQL dans le modèle par défaut et cliquez sur Run Now. Voici un exemple :

    -- Example of Spark SQL. Modify the content and run your Spark program.
    
    conf spark.driver.resourceSpec=medium;
    conf spark.executor.instances=2;
    conf spark.executor.resourceSpec=medium;
    conf spark.app.name=Spark SQL Test;
    conf spark.adb.connectors=oss;
    
    -- SQL statements
    show tables from lakehouse20220413156_adbTest;
  3. (Facultatif) Dans l'onglet Applications, cliquez sur Logs dans la colonne Actions pour afficher les journaux d'exécution Spark SQL.

Étape 5 (Facultative) : Gérer la source de données

Accédez à Data Ingestion > Data Sources. Les opérations suivantes sont disponibles dans la colonne Actions :

Opération Description
Create Job Créez une tâche de synchronisation ou de migration des données pour cette source de données.
View Consultez la configuration de la source de données.
Edit Modifiez le nom et la description de la source de données.
Delete Supprimez la source de données. Si une tâche de synchronisation ou de migration existe pour cette source de données, supprimez d'abord cette tâche sur la page Simple Log Service/Kafka Data Synchronization.

Niveaux d'analyse JSON et inférence de schéma

Le paramètre Parsed JSON Layers contrôle le nombre de niveaux imbriqués d'un message JSON développés en champs de destination distincts.

Exemple de message :

{
  "name" : "zhangle",
  "age" : 18,
  "device" : {
    "os" : {
        "test":lag,
        "member":{
             "fa":zhangsan,
             "mo":limei
           }
         },
    "brand" : "none",
    "version" : "11.4.2"
  }
}
Important

Les points (.) présents dans les noms de champs sont automatiquement remplacés par des underscores (_) dans les noms de champs de destination.

Niveau 0 — Aucune analyse. L'intégralité du message JSON est sortie sous forme d'un seul champ.

Champ JSON Valeur Nom du champ de destination
__value__ {"name":"zhangle","age":18,"device":{...}} __value__

Niveau 1 (par défaut) — Les champs de premier niveau sont développés.

Champ JSON Valeur Nom du champ de destination
name zhangle name
age 18 age
device {"os":{...},"brand":"none","version":"11.4.2"} device

Niveau 2 — Deux niveaux développés. Les champs non imbriqués sont sortis directement ; les champs imbriqués se développent en leurs sous-champs.

Champ JSON Valeur Nom du champ de destination
name zhangle name
age 18 age
device.os {"test":"lag","member":{...}} device_os
device.brand none device_brand
device.version 11.4.2 device_version

Niveau 3

Champ JSON Valeur Nom du champ de destination
name zhangle name
age 18 age
device.os.test lag device_os_test
device.os.member {"fa":"zhangsan","mo":"limei"} device_os_member
device.brand none device_brand
device.version 11.4.2 device_version

Niveau 4

Champ JSON Valeur Nom du champ de destination
name zhangle name
age 18 age
device.os.test lag device_os_test
device.os.member.fa zhangsan device_os_member_fa
device.os.member.mo lime device_os_member_mo
device.brand none device_brand
device.version 11.4.2 device_version

Étapes suivantes