Tous les produits
Search
Centre de documentation

DataWorks:Source de données Kafka

Dernière mise à jour :Aug 26, 2026

La source de données Kafka offre un canal bidirectionnel pour lire et écrire des données dans Kafka. Cette rubrique décrit les capacités de synchronisation des données que DataWorks fournit pour Kafka.

Versions prises en charge

DataWorks prend en charge Alibaba Cloud Kafka et les versions de Kafka gérées par l'utilisateur allant de la 0.10.2 à la 3.6.x.

Remarque

La synchronisation des données n'est pas prise en charge pour les versions de Kafka antérieures à la 0.10.2. En effet, ces versions ne permettent pas de récupérer les décalages (offsets) des partitions et leurs structures de données peuvent ne pas prendre en charge les horodatages.

Lecture en temps réel

  • Si vous utilisez un groupe de ressources Serverless par abonnement, vous devez estimer les spécifications requises à l'avance afin d'éviter les échecs de tâches dus à des ressources insuffisantes.

    Comptez 1 CU par topic. Vous devez également estimer les ressources en fonction du trafic :

    • Pour les données Kafka non compressées, comptez 1 CU pour chaque 10 Mo/s de trafic.

    • Pour les données Kafka compressées, comptez 2 CU pour chaque 10 Mo/s de trafic.

    • Pour les données Kafka compressées nécessitant une analyse JSON, comptez 3 CU pour chaque 10 Mo/s de trafic.

  • Lorsque vous utilisez un groupe de ressources Serverless par abonnement ou une ancienne version d'un groupe de ressources exclusif pour Data Integration :

    • Si votre charge de travail tolère bien le basculement, l'utilisation des slots du cluster ne doit pas dépasser 80 %.

    • Si votre charge de travail tolère mal le basculement, l'utilisation des slots du cluster ne doit pas dépasser 70 %.

Remarque

L'utilisation réelle des ressources dépend de facteurs tels que le contenu et le format des données. Après votre évaluation initiale des ressources, ajustez-les en fonction de l'utilisation réelle lors de l'exécution.

Limitations

La source de données Kafka prend en charge les groupes de ressources Serverless (recommandés) et les anciennes versions des groupes de ressources exclusifs pour Data Integration.

Lecture hors ligne depuis une seule table

Si les paramètres parameter.groupId et parameter.kafkaConfig.group.id sont tous deux configurés, parameter.groupId a priorité sur group.id dans le paramètre kafkaConfig.

Écriture en temps réel dans une seule table

Les opérations d'écriture ne prennent pas en charge la déduplication des données. Si une tâche est redémarrée après une réinitialisation du décalage (offset) ou un basculement, des données en double peuvent être écrites.

Écriture en temps réel pour une base de données entière

  • Les tâches de synchronisation de données en temps réel prennent en charge les groupes de ressources Serverless (recommandés) et les anciennes versions des groupes de ressources exclusifs pour Data Integration.

  • Si la table source possède une clé primaire, la valeur de cette clé est utilisée comme clé pour l'enregistrement Kafka. Cela garantit que les modifications apportées à la même clé primaire sont écrites dans la même partition Kafka dans l'ordre.

  • Si la table source ne possède pas de clé primaire, deux options s'offrent à vous. Si vous sélectionnez l'option de synchronisation des tables sans clés primaires, la clé de l'enregistrement Kafka reste vide. Pour garantir que les modifications de la table soient écrites dans Kafka dans l'ordre, le topic Kafka de destination ne doit comporter qu'une seule partition. Si vous sélectionnez une clé primaire personnalisée, une combinaison d'un ou plusieurs champs non clés primaires est utilisée comme clé pour l'enregistrement Kafka.

  • Pour garantir que les modifications relatives à la même clé primaire soient écrites dans la même partition Kafka dans l'ordre, même si le cluster Kafka renvoie une exception, ajoutez la configuration suivante dans le formulaire des paramètres étendus.

    {"max.in.flight.requests.per.connection":1,"buffer.memory": 100554432}

    Important

    Cette configuration dégrade considérablement les performances de réplication. Vous devez trouver un équilibre entre les performances et la nécessité d'un ordre strict et d'une fiabilité élevée.

  • Pour plus d'informations sur le format global des messages écrits dans Kafka lors de la synchronisation en temps réel, le format des messages de heartbeat et le format des messages correspondant aux modifications des données source, consultez l'Annexe : Format des messages.

Types de champs pris en charge

Kafka fournit un stockage de données non structuré. Un enregistrement Kafka comprend généralement les champs suivants : key, value, offset, timestamp, headers et partition. Lorsque DataWorks lit ou écrit des données dans Kafka, il traite les données comme suit.

Lecture des données

Lorsque DataWorks lit des données depuis Kafka, il peut analyser les données au format JSON. Le tableau suivant décrit comment chaque module de données est traité.

Module de données d'enregistrement Kafka

Type de données traité

key

Dépend de l'élément de configuration keyType dans la tâche de synchronisation des données. Pour plus d'informations sur le paramètre keyType, consultez la description complète des paramètres en annexe.

value

Dépend de l'élément de configuration valueType dans la tâche de synchronisation des données. Pour plus d'informations sur le paramètre valueType, consultez la description complète des paramètres en annexe.

offset

Long

timestamp

Long

headers

String

partition

Long

Écriture des données

DataWorks écrit les données dans Kafka au format JSON ou texte. La politique de traitement varie selon le type de tâche de synchronisation, comme décrit dans le tableau ci-dessous.

Important
  • Lorsque les données sont écrites au format texte, les noms de champs ne sont pas inclus. Les valeurs des champs sont séparées par un séparateur.

  • Lorsqu'une tâche de synchronisation en temps réel écrit des données dans Kafka, elle utilise le format JSON intégré. Les données incluent des informations telles que les messages de modification de base de données, l'heure métier et les informations DDL (Data Definition Language). Pour plus de détails sur le format des données, consultez l'Annexe : Format des messages.

Type de tâche de synchronisation

Format de la valeur écrite dans Kafka

Type de champ source

Méthode de traitement pour les opérations d'écriture

Synchronisation hors ligne

Nœud de synchronisation hors ligne dans DataStudio

json

String

Chaîne encodée en UTF-8

Boolean

Converti en chaîne encodée en UTF-8 « true » ou « false »

Heure/Date

Chaîne encodée en UTF-8 au format yyyy-MM-dd HH:mm:ss

Numérique

Chaîne numérique encodée en UTF-8

Flux d'octets

Le flux d'octets est traité comme une chaîne encodée en UTF-8 et converti en chaîne.

text

String

Chaîne encodée en UTF-8

Boolean

Converti en chaîne encodée en UTF-8 « true » ou « false »

Heure/Date

Chaîne encodée en UTF-8 au format yyyy-MM-dd HH:mm:ss

Numérique

Chaîne numérique encodée en UTF-8

Flux d'octets

Le flux d'octets est traité comme une chaîne encodée en UTF-8 et converti en chaîne.

Synchronisation en temps réel : ETL en temps réel vers Kafka

Nœud de synchronisation en temps réel dans DataStudio

json

String

Chaîne encodée en UTF-8

Boolean

Type Boolean JSON

Heure/Date

  • Pour les valeurs temporelles avec une précision inférieure à la milliseconde : Converties en un entier JSON de 13 chiffres représentant l'horodatage en millisecondes.

  • Pour les valeurs temporelles avec une précision en microsecondes ou nanosecondes : Converties en un nombre à virgule flottante JSON comprenant un entier de 13 chiffres pour l'horodatage en millisecondes et une décimale de 6 chiffres pour l'horodatage en nanosecondes.

Numérique

Type numérique JSON

Flux d'octets

Le flux d'octets est encodé en Base64, puis converti en une chaîne encodée en UTF-8.

text

String

Chaîne encodée en UTF-8

Boolean

Converti en chaîne encodée en UTF-8 « true » ou « false »

Heure/Date

Chaîne encodée en UTF-8 au format yyyy-MM-dd HH:mm:ss

Numérique

Chaîne numérique encodée en UTF-8

Flux d'octets

Le flux d'octets est encodé en Base64, puis converti en une chaîne encodée en UTF-8.

Synchronisation en temps réel : Synchronisation en temps réel d'une base de données entière vers Kafka

Synchronisation en temps réel des données incrémentielles uniquement

Format JSON intégré

String

Chaîne encodée en UTF-8

Boolean

Type Boolean JSON

Heure/Date

Horodatage en millisecondes sur 13 chiffres

Numérique

Valeur numérique JSON

Flux d'octets

Le flux d'octets est encodé en Base64, puis converti en une chaîne encodée en UTF-8.

Solution de synchronisation : Synchronisation en temps réel en un clic vers Kafka

Synchronisation hors ligne complète + synchronisation en temps réel incrémentielle

Format JSON intégré

String

Chaîne encodée en UTF-8

Boolean

Type Boolean JSON

Heure/Date

Horodatage en millisecondes sur 13 chiffres

Numérique

Valeur numérique JSON

Flux d'octets

Le flux d'octets est encodé en Base64, puis converti en une chaîne encodée en UTF-8.

Ajouter une source de données

Avant de développer une tâche de synchronisation dans DataWorks, vous devez ajouter la source de données requise à DataWorks en suivant les instructions de la rubrique Configuration de la source de données. Vous pouvez consulter les descriptions des paramètres dans la console DataWorks pour comprendre la signification des paramètres lors de l'ajout d'une source de données.

Développer une tâche de synchronisation des données

Pour obtenir des informations sur le point d'entrée et la procédure de configuration d'une tâche de synchronisation, consultez les guides de configuration suivants.

Configurer une tâche de synchronisation hors ligne pour une seule table

Configurer une tâche de synchronisation en temps réel pour une seule table ou une base de données entière

Pour obtenir des instructions, consultez les rubriques Configurer une tâche de synchronisation en temps réel pour une seule table et Configurer une tâche de synchronisation en temps réel pour une base de données entière.

Configuration de l'authentification

SSL

Lorsque vous configurez une source de données Kafka, si vous définissez la Special Authentication Method sur SSL ou SASL_SSL, vous activez l'authentification SSL pour le cluster Kafka. Vous devez télécharger le fichier de certificat du magasin de confiance (truststore) client et saisir la phrase secrète associée.

  • Si le cluster Kafka est une instance Alibaba Cloud Kafka, consultez les instructions de mise à niveau de l'algorithme de certificat SSL pour télécharger le fichier de certificat truststore approprié. La phrase secrète du truststore est KafkaOnsClient.

  • Si le cluster Kafka est une instance EMR, reportez-vous à la rubrique Utilisation du chiffrement SSL pour les connexions Kafka afin de télécharger le fichier de certificat truststore correct et d'obtenir la phrase secrète correspondante.

  • Pour un cluster auto-géré, vous devez télécharger le certificat truststore approprié et saisir la phrase secrète correcte.

Le fichier de certificat keystore, la phrase secrète du keystore et la phrase secrète SSL ne sont requis que lorsque l'authentification SSL bidirectionnelle est activée pour le cluster Kafka. Le serveur du cluster Kafka utilise ces éléments pour authentifier l'identité du client. L'authentification SSL bidirectionnelle est activée lorsque le paramètre ssl.client.auth=required est défini dans le fichier server.properties du cluster Kafka. Pour plus d'informations, consultez la rubrique Utilisation du chiffrement SSL pour les connexions Kafka.

GSSAPI

Si vous définissez le Sasl Mechanism sur GSSAPI lors de la configuration d'une source de données Kafka, vous devez télécharger trois fichiers d'authentification : un fichier de configuration JAAS, un fichier de configuration Kerberos et un fichier Keytab. Vous devez également configurer les paramètres DNS/HOST pour le groupe de ressources exclusif. Les sections suivantes décrivent ces fichiers ainsi que les paramètres DNS et HOST requis.

Remarque

Pour un groupe de ressources Serverless, vous devez configurer les informations d'adresse hôte en utilisant la résolution DNS interne. Pour plus d'informations, consultez la rubrique Résolution DNS interne (PrivateZone).

  • Fichier de configuration JAAS

    Le fichier JAAS doit commencer par KafkaClient, suivi de tous les éléments de configuration enclosed entre accolades {} :

    • La première ligne à l'intérieur des accolades définit la classe du composant de connexion à utiliser. Pour les différents mécanismes d'authentification SASL, cette classe est fixe. Chaque élément de configuration suivant est écrit au format key=value.

    • Tous les éléments de configuration, à l'exception du dernier, ne doivent pas se terminer par un point-virgule.

    • Le dernier élément de configuration doit se terminer par un point-virgule, et un autre point-virgule doit suivre l'accolade fermante }.

    Si les exigences de format ne sont pas respectées, le fichier de configuration JAAS ne peut pas être analysé. Le code suivant illustre un format typique de fichier de configuration JAAS. Remplacez les espaces réservés xxx par vos informations réelles.

    KafkaClient {
       com.sun.security.auth.module.Krb5LoginModule required
       useKeyTab=true
       keyTab="xxx"
       storeKey=true
       serviceName="kafka-server"
       principal="kafka-client@EXAMPLE.COM";
    };

    Configuration item

    Description

    Logon module

    Doit être défini sur com.sun.security.auth.module.Krb5LoginModule.

    useKeyTab

    Doit être défini sur true.

    keyTab

    Vous pouvez spécifier n'importe quel chemin. Lorsque la tâche de synchronisation s'exécute, le système télécharge automatiquement le fichier keytab que vous avez téléchargé lors de la configuration de la source de données vers un chemin local. Le système utilise ensuite ce chemin local pour l'élément de configuration keytab.

    storeKey

    Indique si le client enregistre la clé. Vous pouvez définir cette option sur true ou false. Cela n'affecte pas la synchronisation des données.

    serviceName

    Correspond à l'élément de configuration sasl.kerberos.service.name dans le fichier de configuration server.properties du serveur Kafka. Configurez cet élément selon vos besoins.

    principal

    Principal Kerberos utilisé par le client Kafka. Configurez-le selon vos besoins et assurez-vous que le fichier keytab téléchargé contient la clé pour ce principal.

  • Fichier de configuration Kerberos

    Le fichier de configuration Kerberos doit contenir deux modules : [libdefaults] et [realms].

    • Le module [libdefaults] spécifie les paramètres d'authentification Kerberos. Chaque élément de configuration du module est écrit au format key=value.

    • Le module [realms] spécifie l'adresse du centre de distribution de clés (KDC). Il peut contenir plusieurs sous-modules realm. Chaque sous-module realm commence par le nom du realm suivi d'un signe égal (=).

    Ceci est suivi d'un ensemble d'éléments de configuration entre accolades. Chaque élément est également écrit au format key=value. Le code suivant illustre un format typique de fichier de configuration Kerberos. Remplacez les espaces réservés xxx par vos informations réelles.

    [libdefaults]
      default_realm = xxx
    
    [realms]
      xxx = {
        kdc = xxx
      }

    Configuration item

    Description

    [libdefaults].default_realm

    Realm par défaut utilisé lors de l'accès aux nœuds du cluster Kafka. Il s'agit généralement du même realm que celui du principal client spécifié dans le fichier de configuration JAAS.

    Other [libdefaults] parameters

    Le module [libdefaults] peut spécifier d'autres paramètres d'authentification Kerberos, tels que ticket_lifetime. Configurez-les selon vos besoins.

    [realms].realm name

    Doit être identique au realm du principal client spécifié dans le fichier de configuration JAAS et au paramètre [libdefaults].default_realm. Si le realm du principal client dans le fichier de configuration JAAS diffère de [libdefaults].default_realm, vous devez inclure deux sous-modules realms. Ces sous-modules doivent correspondre respectivement au realm du principal client dans le fichier de configuration JAAS et au [libdefaults].default_realm.

    [realms].realm name.kdc

    Spécifie l'adresse et le port du KDC au format ip:port, par exemple, kdc=10.0.0.1:88. Si le port est omis, le système utilise le port par défaut 88, par exemple, kdc=10.0.0.1.

  • Fichier Keytab

    Le fichier keytab doit contenir la clé pour le principal spécifié dans le fichier de configuration JAAS et doit être vérifiable par le KDC. Par exemple, s'il existe un fichier nommé client.keytab dans le répertoire de travail actuel, vous pouvez exécuter la commande suivante pour vérifier si le fichier keytab contient la clé pour le principal spécifié.

    klist -ket ./client.keytab
    
    Keytab name: FILE:client.keytab
    KVNO Timestamp           Principal
    ---- ------------------- ------------------------------------------------------
       7 2018-07-30T10:19:16 te**@**.com (des-cbc-md5)
  • Configuration DNS et HOST pour un groupe de ressources exclusif

    Lorsqu'un cluster Kafka utilise l'authentification Kerberos, le KDC enregistre le principal de chaque nœud en utilisant le nom d'hôte du nœud. Lorsqu'un client se connecte à un nœud du cluster Kafka, il utilise les paramètres DNS et HOST locaux pour déduire le principal du nœud, puis demande au KDC des informations d'identification d'accès pour ce nœud. Lorsque vous utilisez un groupe de ressources exclusif pour accéder à un cluster Kafka avec l'authentification Kerberos activée, vous devez configurer correctement les paramètres DNS et HOST afin de garantir que les informations d'identification d'accès pour les nœuds du cluster puissent être obtenues auprès du KDC :

    • Paramètres DNS

      Si vous utilisez une instance PrivateZone pour la résolution de noms de domaine des nœuds du cluster Kafka dans le VPC auquel le groupe de ressources exclusif est attaché, vous pouvez ajouter une route personnalisée pour les adresses IP 100.100.2.136 et 100.100.2.138 à l'attachement VPC. Cela garantit que les paramètres de résolution de noms de domaine PrivateZone pour les nœuds du cluster Kafka s'appliquent au groupe de ressources exclusif. Dans le volet de navigation de gauche de la console DataWorks, cliquez sur resource group list. Dans la colonne Actions de votre groupe de ressources exclusif, cliquez sur network settings. Sous l'onglet VPC attachment, cliquez sur custom route dans la colonne Actions. Dans la boîte de dialogue qui s'affiche, cliquez sur Add Route, définissez Destination Type sur IDC et Connection Method sur Direct IP, saisissez l'adresse IP directe, puis cliquez sur Generate Route.

    • Paramètres HOST

      Si vous n'utilisez pas d'instance PrivateZone pour la résolution de noms de domaine des nœuds du cluster Kafka dans le VPC auquel le groupe de ressources exclusif est attaché, vous devez ajouter les mappages adresse IP-nom de domaine pour chaque nœud du cluster Kafka à la configuration hôte. Dans le volet de navigation de gauche de la console DataWorks, cliquez sur resource group list. Recherchez votre groupe de ressources, cliquez sur network settings dans la colonne Actions, puis sélectionnez l'onglet host configuration. Cliquez sur Add pour ajouter un domaine hôte. La configuration hôte est prioritaire sur les paramètres DNS.

PLAIN

Lorsque vous configurez une source de données Kafka, si vous définissez le Sasl Mechanism sur PLAIN, le fichier JAAS doit commencer par KafkaClient, suivi de tous les éléments de configuration entre accolades {}.

  • La première ligne à l'intérieur des accolades définit la classe du composant de connexion à utiliser. Pour les différents mécanismes d'authentification SASL, cette classe est fixe. Chaque élément de configuration suivant est écrit au format key=value.

  • Tous les éléments de configuration, à l'exception du dernier, ne doivent pas se terminer par un point-virgule.

  • Le dernier élément de configuration doit se terminer par un point-virgule. Un point-virgule doit également être ajouté après l'accolade fermante « } ».

Si les exigences de format ne sont pas respectées, le fichier de configuration JAAS ne peut pas être analysé. Le code suivant illustre un format typique de fichier de configuration JAAS. Remplacez les espaces réservés xxx par vos informations réelles.

KafkaClient {
  org.apache.kafka.common.security.plain.PlainLoginModule required
  username="xxx"
  password="xxx";
};

Configuration item

Description

Logon module

Doit être défini sur org.apache.kafka.common.security.plain.PlainLoginModul

username

Nom d'utilisateur. Configurez cet élément selon vos besoins.

password

Mot de passe. Configurez cet élément selon vos besoins.

FAQ

Annexe : Démos de scripts et description des paramètres

Configuration d'une tâche de synchronisation par lots à l'aide de l'éditeur de code

Si vous souhaitez configurer une tâche de synchronisation par lots à l'aide de l'éditeur de code, vous devez configurer les paramètres associés dans le script en respectant les exigences de format de script unifié. Pour plus d'informations, consultez la rubrique Configuration en mode script. Les informations suivantes décrivent les paramètres que vous devez configurer pour les sources de données lors de la configuration d'une tâche de synchronisation par lots à l'aide de l'éditeur de code.

Démo de script Reader

La configuration JSON suivante lit les données depuis Kafka.

{
    "type": "job",
    "steps": [
        {
            "stepType": "kafka",
            "parameter": {
                "server": "host:9093",
                "column": [
                    "__key__",
                    "__value__",
                    "__partition__",
                    "__offset__",
                    "__timestamp__",
                    "'123'",
                    "event_id",
                    "tag.desc"
                ],
                "kafkaConfig": {
                    "group.id": "demo_test"
                },
                "topic": "topicName",
                "keyType": "ByteArray",
                "valueType": "ByteArray",
                "beginDateTime": "20190416000000",
                "endDateTime": "20190416000006",
                "skipExceedRecord": "true"
            },
            "name": "Reader",
            "category": "reader"
        },
        {
            "stepType": "stream",
            "parameter": {
                "print": false,
                "fieldDelimiter": ","
            },
            "name": "Writer",
            "category": "writer"
        }
    ],
    "version": "2.0",
    "order": {
        "hops": [
            {
                "from": "Reader",
                "to": "Writer"
            }
        ]
    },
    "setting": {
        "errorLimit": {
            "record": "0"
        },
        "speed": {
            "throttle": true,//When throttle is false, the mbps parameter does not take effect, which means that the rate is not limited. When throttle is true, the rate is limited.
            "concurrent": 1,//The number of concurrent threads.
            "mbps":"12"//The maximum transmission rate. 1 mbps is equal to 1 MB/s.
        }
    }
}

Paramètres du script Reader

Parameter

Description

Required

datasource

Nom de la source de données. L'éditeur de code permet d'ajouter des sources de données. La valeur de ce paramètre doit correspondre exactement au nom de la source ajoutée.

Yes

server

Adresse du broker Kafka, au format ip:port.

Vous ne pouvez configurer qu'un seul server, mais vous devez vous assurer que DataWorks peut se connecter aux adresses IP de tous les brokers du cluster Kafka.

Yes

topic

Topic Kafka. Un topic représente un regroupement de flux de messages traités par Kafka.

Yes

column

Données Kafka à lire. Les colonnes constantes, les colonnes de données et les colonnes d'attributs sont prises en charge.

  • Colonne constante : colonne entourée de guillemets simples, par exemple ["'abc'", "'123'"].

  • Colonne de données

    • Si vos données sont au format JSON, vous pouvez extraire les propriétés de l'objet JSON, par exemple ["event_id"].

    • Si vos données sont au format JSON, vous pouvez extraire les sous-propriétés imbriquées de l'objet JSON, par exemple ["tag.desc"].

  • Colonne d'attribut

    • __key__ : clé du message.

    • __value__ : contenu intégral du message.

    • __partition__ : partition dans laquelle se trouve le message actuel.

    • __headers__ : en-têtes du message actuel.

    • __offset__ : décalage (offset) du message actuel.

    • __timestamp__ : horodatage du message actuel.

    L'exemple suivant illustre une configuration complète.

    "column": [
        "__key__",
        "__value__",
        "__partition__",
        "__offset__",
        "__timestamp__",
        "'123'",
        "event_id",
        "tag.desc"
        ]

Yes

keyType

Type de la clé Kafka. Valeurs possibles : BYTEARRAY, DOUBLE, FLOAT, INTEGER, LONG et SHORT.

No

valueType

Type de la valeur Kafka. Valeurs possibles : BYTEARRAY, DOUBLE, FLOAT, INTEGER, LONG et SHORT.

No

beginDateTime

Heure de début de la consommation des données. Ce paramètre définit la borne inférieure (inclusive) de la plage temporelle. Il s'agit d'une chaîne au format yyyymmddhhmmss. Vous pouvez l'utiliser conjointement avec les scheduling parameters. Pour plus d'informations, consultez la rubrique Formats pris en charge pour les paramètres de planification.

Remarque

Cette fonctionnalité est disponible à partir de Kafka 0.10.2.

Vous devez spécifier ce paramètre ou beginOffset.

Remarque

Les paramètres beginDateTime et endDateTime doivent être utilisés conjointement.

endDateTime

Heure de fin de la consommation des données. Ce paramètre définit la borne supérieure (exclusive) de la plage temporelle. Il s'agit d'une chaîne au format yyyymmddhhmmss. Vous pouvez l'utiliser conjointement avec les scheduling parameters. Pour plus d'informations, consultez la rubrique Formats pris en charge pour les paramètres de planification.

Remarque

Cette fonctionnalité est disponible à partir de Kafka 0.10.2.

Vous devez spécifier ce paramètre ou endOffset.

Remarque

Les paramètres endDateTime et beginDateTime doivent être utilisés conjointement.

beginOffset

Décalage (offset) de départ pour la consommation des données. Vous pouvez le configurer selon les modalités suivantes :

  • Un nombre, par exemple 15553274, qui indique le décalage de départ pour la consommation.

  • seekToBeginning : indique que la consommation commence à partir du décalage le plus ancien.

  • seekToLast : lit les données à partir du décalage enregistré pour l'ID de groupe spécifié par le paramètre group.id dans kafkaConfig. Notez que le client valide automatiquement le décalage du groupe sur le serveur Kafka à intervalles réguliers. Par conséquent, si une tâche échoue et est relancée, des duplications ou des pertes de données peuvent survenir. Si le paramètre skipExceedRecord est défini sur true, la tâche peut ignorer les dernières lignes lues. Le décalage du groupe correspondant à ces données ignorées ayant déjà été validé sur le serveur, ces données ne pourront pas être relues lors de la prochaine exécution de la tâche.

  • seekToEnd : indique que la consommation commence à partir du décalage le plus récent. Cela entraînera la lecture de données vides.

Vous devez spécifier ce paramètre ou beginDateTime.

endOffset

Décalage (offset) de fin pour la consommation des données. Ce paramètre permet de contrôler l'arrêt de la tâche de consommation.

Vous devez spécifier ce paramètre ou endDateTime.

skipExceedRecord

Kafka utilise la méthode public ConsumerRecords<K, V> poll(final Duration timeout) pour consommer les données. Un seul appel à poll peut récupérer des données au-delà du endOffset ou de l'endDateTime spécifié. Ce paramètre détermine si ces données excédentaires sont écrites dans la destination. Étant donné que la tâche utilise la validation automatique des décalages, nous recommandons les configurations suivantes :

  • Pour les versions de Kafka antérieures à 0.10.2 : définissez skipExceedRecord sur false.

  • Pour Kafka 0.10.2 et versions ultérieures : définissez skipExceedRecord sur true.

No. La valeur par défaut est false.

partition

Un topic Kafka comporte plusieurs partitions (partition). Par défaut, une tâche de synchronisation de données lit les données sur une plage de décalages couvrant toutes les partitions du topic. Vous pouvez également spécifier une partition pour limiter la lecture à la plage de décalages d'une seule partition.

No. Aucune valeur par défaut.

kafkaConfig

Lors de la création d'un client KafkaConsumer pour la consommation de données, vous pouvez spécifier des paramètres étendus tels que bootstrap.servers, auto.commit.interval.ms et session.timeout.ms. Utilisez kafkaConfig pour contrôler le comportement de consommation du KafkaConsumer.

No

encoding

Lorsque keyType ou valueType est défini sur STRING, l'encodage spécifié par ce paramètre est utilisé pour analyser la chaîne.

No. La valeur par défaut est UTF-8.

waitTIme

Durée maximale, en secondes, pendant laquelle l'objet consumer attend pour extraire des données de Kafka lors d'une tentative unique.

No. La valeur par défaut est 60.

stopWhenPollEmpty

Les valeurs valides sont true et false. Si ce paramètre est défini sur true et que le consumer extrait des données vides de Kafka (généralement parce que toutes les données du topic ont été lues, ou en raison de problèmes réseau ou de disponibilité du cluster Kafka), la tâche s'arrête immédiatement. Sinon, elle réessaie jusqu'à ce que des données soient à nouveau lues.

No. La valeur par défaut est true.

stopWhenReachEndOffset

Ce paramètre prend effet uniquement lorsque stopWhenPollEmpty est défini sur true. Les valeurs valides sont true et false.

  • Si ce paramètre est défini sur true et que le consumer ne reçoit aucune donnée lors d'une requête poll, il vérifie si le décalage le plus récent a été atteint dans les partitions du topic. Si le décalage le plus récent a été atteint pour toutes les partitions, la tâche s'arrête immédiatement. Sinon, elle continue d'essayer d'extraire des données du topic.

  • Si ce paramètre est défini sur false et que le consumer ne reçoit aucune donnée lors d'une requête poll, il n'effectue pas cette vérification et arrête la tâche immédiatement.

No. La valeur par défaut est false.

Remarque

Ce paramètre assure la rétrocompatibilité. Les versions de Kafka antérieures à 0.10.2 ne prennent pas en charge la vérification du décalage le plus récent de toutes les partitions.

Le tableau suivant décrit les paramètres kafkaConfig.

Parameter

Description

fetch.min.bytes

Quantité minimale de données, en octets, que le consumer récupère auprès du broker lors d'une seule requête. Le broker attend que cette quantité de données soit disponible avant de répondre au consumer.

fetch.max.wait.ms

Durée maximale, en millisecondes, pendant laquelle le broker attend que des données deviennent disponibles avant de répondre à une requête de récupération. La valeur par défaut est 500. Le broker répond dès que la condition fetch.min.bytes ou fetch.max.wait.ms est satisfaite.

max.partition.fetch.bytes

Spécifie le nombre maximal d'octets que le broker peut renvoyer au consumer pour chaque partition. La valeur par défaut est 1 Mo.

session.timeout.ms

Spécifie la durée pendant laquelle le consumer peut rester déconnecté du serveur avant d'arrêter de recevoir des services. La valeur par défaut est 30 secondes.

auto.offset.reset

Action entreprise par le consumer lors d'une lecture sans décalage ou avec un décalage invalide (par exemple, si le consumer est resté inactif longtemps et que l'enregistrement correspondant au décalage a expiré et a été supprimé). La valeur par défaut est none, ce qui signifie que le décalage n'est pas réinitialisé automatiquement. Vous pouvez la modifier en earliest, ce qui signifie que le consumer lit les enregistrements de la partition à partir du décalage le plus ancien.

max.poll.records

Nombre de messages pouvant être renvoyés par un seul appel à la méthode poll.

key.deserializer

Méthode de désérialisation pour la clé du message, par exemple org.apache.kafka.common.serialization.StringDeserializer.

value.deserializer

Méthode de désérialisation pour la valeur des données, par exemple org.apache.kafka.common.serialization.StringDeserializer.

ssl.truststore.location

Chemin d'accès au certificat racine SSL.

ssl.truststore.password

Mot de passe du magasin de certificats racines. Si vous utilisez Alibaba Cloud Kafka, définissez cette valeur sur KafkaOnsClient.

security.protocol

Protocole d'accès. Actuellement, seul le protocole SASL_SSL est pris en charge.

sasl.mechanism

Méthode d'authentification SASL. Si vous utilisez Alibaba Cloud Kafka, utilisez PLAIN.

java.security.auth.login.config

Chemin d'accès au fichier d'authentification SASL.

Exemple de script Writer

Le code suivant présente la configuration JSON pour l'écriture de données dans Kafka.

{
  "type":"job",
  "version":"2.0",//The version number.
  "steps":[
    {
      "stepType":"stream",
      "parameter":{},
      "name":"Reader",
      "category":"reader"
    },
    {
      "stepType":"Kafka",//The plug-in name.
      "parameter":{
          "server": "ip:9092", //The server address of Kafka.
          "keyIndex": 0, //The column to be used as the key. Must follow camel case naming conventions, with k in lowercase.
          "valueIndex": 1, //The column to be used as the value. Currently, you can only select one column from the source data or leave this parameter empty. If empty, all source data is used.
          //For example, to use the 2nd, 3rd, and 4th columns of an ODPS table as the kafkaValue, create a new ODPS table, clean and integrate the data from the original ODPS table into the new table, and then use the new table for synchronization.
          "keyType": "Integer", //The type of the Kafka key.
          "valueType": "Short", //The type of the Kafka value.
          "topic": "t08", //The Kafka topic.
          "batchSize": 1024 //The amount of data written to Kafka at a time, in bytes.
        },
      "name":"Writer",
      "category":"writer"
    }
  ],
  "setting":{
      "errorLimit":{
      "record":"0"//The number of error records.
    },
    "speed":{
        "throttle":true,//When throttle is false, the mbps parameter does not take effect, which means that the rate is not limited. When throttle is true, the rate is limited.
        "concurrent":1, //The number of concurrent jobs.
        "mbps":"12"//The maximum transmission rate. 1 mbps is equal to 1 MB/s.
    }
   },
    "order":{
    "hops":[
        {
            "from":"Reader",
            "to":"Writer"
        }
      ]
    }
}

Paramètres du script Writer

Parameter

Description

Required

datasource

Nom de la source de données. L'éditeur de code prend en charge l'ajout de sources de données. La valeur de ce paramètre doit correspondre au nom de la source de données ajoutée.

Oui

server

Adresse du serveur Kafka au format ip:port.

Oui

topic

Rubrique Kafka. Il s'agit d'une catégorie pour les différents flux de messages traités par Kafka.

Chaque message publié sur un cluster Kafka appartient à une catégorie, appelée rubrique (topic). Une rubrique est un ensemble de messages.

Oui

valueIndex

Colonne utilisée comme valeur dans le writer Kafka. Si ce paramètre n'est pas spécifié, toutes les colonnes sont concaténées par défaut pour former la valeur. Le séparateur est défini par le paramètre fieldDelimiter.

Non

writeMode

Lorsque le paramètre valueIndex n'est pas configuré, ce paramètre détermine le format de concaténation de toutes les colonnes de l'enregistrement source pour former la valeur de l'enregistrement Kafka. Les valeurs valides sont text et JSON. La valeur par défaut est text.

  • Si la valeur est text, toutes les colonnes sont concaténées à l'aide du séparateur spécifié par fieldDelimiter.

  • Si la valeur est JSON, toutes les colonnes sont concaténées en une chaîne JSON basée sur les noms de champs spécifiés par le paramètre column.

Par exemple, si un enregistrement source comporte trois colonnes avec les valeurs a, b et c, et que writeMode est défini sur text et fieldDelimiter sur #, la valeur de l'enregistrement Kafka écrit est la chaîne a#b#c. Si writeMode est défini sur JSON et column sur [{"name":"col1"},{"name":"col2"},{"name":"col3"}], la valeur de l'enregistrement Kafka écrit est la chaîne {"col1":"a","col2":"b","col3":"c"}.

Si le paramètre valueIndex est configuré, ce paramètre est ignoré.

Non

column

Champs de la table de destination vers lesquels les données doivent être écrites, séparés par des virgules. Par exemple : "column": ["id", "name", "age"].

Lorsque le paramètre valueIndex n'est pas configuré et que writeMode est défini sur JSON, ce paramètre définit les noms de champs dans la structure JSON pour les valeurs de colonne de l'enregistrement source. Par exemple, "column": [{"name":id","type":"JSON_NUMBER"}, {"name":"name","type":"JSON_STRING"}, {"name":"age","type":"JSON_NUMBER"}].

  • Si le nombre de colonnes dans l'enregistrement source est supérieur au nombre de noms de champs configurés dans column, les données sont tronquées lors de l'écriture. Par exemple :

    Si un enregistrement source comporte trois colonnes avec les valeurs a, b et c, et que column est configuré comme suit : [{"name":"col1","type":"JSON_STRING"},{"name":"col2","type":"JSON_STRING"}], la valeur de l'enregistrement Kafka écrit est la chaîne {"col1":"a","col2":"b"}.

  • Si le nombre de colonnes dans l'enregistrement source est inférieur au nombre de champs spécifiés dans column, les champs supplémentaires sont renseignés avec null ou la chaîne spécifiée par nullValueFormat. Par exemple :

    Si un enregistrement source comporte deux colonnes avec les valeurs a et b, et que column est configuré comme suit : [{"name":"col1","type":"JSON_STRING"},{"name":"col2","type":"JSON_STRING"},{"name":"col3","type":"JSON_STRING"}], la valeur de l'enregistrement Kafka écrit est la chaîne {"col1":"a","col2":"b","col3":null}. Si le paramètre valueIndex est configuré ou si writeMode est défini sur text, ce paramètre est ignoré.

  • Si le type de champ JSON n'est pas configuré, le type de champ par défaut est JSON_STRING.

  • Pour connaître les valeurs valides du type de champ JSON, consultez la section Annexe : Types de champs JSON.

Si le paramètre valueIndex est configuré ou si writeMode est défini sur text, ce paramètre est ignoré.

Requis lorsque le paramètre valueIndex n'est pas configuré et que writeMode est défini sur JSON.

partition

Spécifie le numéro de la partition de la rubrique Kafka vers laquelle les données sont écrites. Il doit s'agir d'un entier supérieur ou égal à 0.

Non

keyIndex

Colonne utilisée comme clé dans le writer Kafka.

La valeur du paramètre keyIndex doit être un entier supérieur ou égal à 0. Dans le cas contraire, la tâche échouera.

Non

keyIndexes

Tableau des numéros d'ordre des colonnes de l'enregistrement source utilisées comme clé pour l'enregistrement Kafka.

Le numéro d'ordre de la colonne commence à 0. Par exemple, [0,1,2] concatène les valeurs de tous les numéros de colonne configurés avec des virgules pour former la clé de l'enregistrement Kafka. Si ce paramètre n'est pas spécifié, la clé de l'enregistrement Kafka est null et les données sont écrites dans les partitions de la rubrique selon un algorithme round-robin. Vous ne pouvez spécifier que ce paramètre ou keyIndex, mais pas les deux.

Non

fieldDelimiter

Lorsque writeMode est défini sur text et que valueIndex n'est pas configuré, toutes les colonnes de l'enregistrement source sont concaténées à l'aide du séparateur de colonne spécifié par ce paramètre pour former la valeur de l'enregistrement Kafka. Vous pouvez configurer un seul caractère ou plusieurs caractères comme séparateur. Vous pouvez configurer des caractères Unicode au format \u0001. Les caractères d'échappement tels que \t et \n sont pris en charge. La valeur par défaut est \t.

Si writeMode n'est pas défini sur text ou si le paramètre valueIndex est configuré, ce paramètre est ignoré.

Non

keyType

Type de la clé Kafka. Valeurs valides : BYTEARRAY, DOUBLE, FLOAT, INTEGER, LONG et SHORT.

Oui

valueType

Type de la valeur Kafka. Valeurs valides : BYTEARRAY, DOUBLE, FLOAT, INTEGER, LONG et SHORT.

Oui

nullKeyFormat

Si la valeur de la colonne source spécifiée par keyIndex ou keyIndexes est null, elle est remplacée par la chaîne spécifiée par ce paramètre. Si ce paramètre n'est pas configuré, aucun remplacement n'est effectué.

Non

nullValueFormat

Si la valeur d'une colonne source est null, elle est remplacée par la chaîne spécifiée par ce paramètre lors de l'assemblage de la valeur de l'enregistrement Kafka. Si vous ne spécifiez pas ce paramètre, aucun remplacement n'est effectué.

Non

acks

Configuration acks lors de l'initialisation du producteur Kafka. Elle détermine la méthode d'accusé de réception pour les écritures réussies. Par défaut, le paramètre acks est défini sur all. Les valeurs valides pour acks sont :

  • 0 : Aucun accusé de réception pour les écritures réussies.

  • 1 : Accusé de réception pour l'écriture réussie sur le réplica principal.

  • all : Accusé de réception pour l'écriture réussie sur tous les réplicas.

Non

Annexe : Définition du format de message pour l'écriture dans Kafka

Après avoir configuré et exécuté une tâche de synchronisation en temps réel, les données lues depuis la base de données source sont écrites dans une rubrique Kafka au format JSON. Tout d'abord, toutes les données existantes de la table source spécifiée sont écrites dans la rubrique Kafka correspondante. Ensuite, la tâche démarre la synchronisation en temps réel pour écrire continuellement les données incrémentielles dans la rubrique. Les informations de modification DDL incrémentielles de la table source sont également écrites dans la rubrique Kafka au format JSON. Vous pouvez obtenir le statut et les informations de modification des messages écrits dans Kafka. Pour plus d'informations, consultez la section Annexe : Format de message.

Remarque

Dans les données JSON issues d'une tâche de synchronisation hors ligne, les champs payload.sequenceId, payload.timestamp.eventTime et payload.timestamp.checkpointTime sont définis sur -1.

Annexe : Types de champs JSON

Lorsque le paramètre writeMode est défini sur JSON, vous pouvez définir le type de champ JSON à l'aide du champ type dans le paramètre column. Lors d'une opération d'écriture, le système tente de convertir la valeur de la colonne de l'enregistrement source vers le type spécifié. Un échec de conversion de type entraîne des données incorrectes (dirty data).

Valid value

Description

JSON_STRING

Convertit la valeur de la colonne de l'enregistrement source en chaîne et l'écrit dans le champ JSON. Par exemple, si la valeur de la colonne de l'enregistrement source est l'entier 123 et que le paramètre column est configuré comme suit : [{"name":"col1","type":"JSON_STRING"}], la valeur écrite dans l'enregistrement Kafka est la chaîne {"col1":"123"}.

JSON_NUMBER

Convertit la valeur de la colonne de l'enregistrement source en nombre et l'écrit dans le champ JSON. Par exemple, si la valeur de la colonne de l'enregistrement source est la chaîne 1.23 et que le paramètre column est configuré comme suit : [{"name":"col1","type":"JSON_NUMBER"}], la valeur écrite dans l'enregistrement Kafka est la chaîne {"col1":1.23}.

JSON_BOOL

Convertit la valeur de la colonne de l'enregistrement source en valeur booléenne et l'écrit dans le champ JSON. Par exemple, si la valeur de la colonne de l'enregistrement source est la chaîne true et que le paramètre column est configuré comme suit : [{"name":"col1","type":"JSON_BOOL"}, la valeur de l'enregistrement Kafka écrit est la chaîne {"col1":true}

JSON_ARRAY

Convertit la valeur de la colonne de l'enregistrement source en tableau JSON et l'écrit dans le champ JSON. Par exemple, si la valeur de la colonne de l'enregistrement source est la chaîne [1,2,3] et que le paramètre column est configuré comme suit : [{"name":"col1","type":"JSON_ARRAY"}], la valeur écrite dans l'enregistrement Kafka est la chaîne {"col1":[1,2,3]}.

JSON_MAP

Convertit la valeur de la colonne de l'enregistrement source en objet JSON et l'écrit dans le champ JSON. Par exemple, si la valeur de la colonne de l'enregistrement source est la chaîne {"k1":"v1"} et que le paramètre column est configuré comme suit : [{"name":"col1","type":"JSON_MAP"}], la valeur écrite dans l'enregistrement Kafka est la chaîne {"col1":{"k1":"v1"}}.

JSON_BASE64

Convertit un tableau d'octets de la colonne source en chaîne encodée en BASE64 et l'écrit dans le champ JSON. Par exemple, si la valeur de la colonne de l'enregistrement source est un tableau de 2 octets représenté en hexadécimal par 0x01 0x02, et que le paramètre column est configuré comme suit : [{"name":"col1","type":"JSON_BASE64"}], la valeur écrite dans l'enregistrement Kafka est la chaîne {"col1":"AQI="}.

JSON_HEX

Convertit un tableau d'octets de la colonne source en chaîne hexadécimale et l'écrit dans le champ JSON. Par exemple, si la valeur de la colonne de l'enregistrement source est un tableau de 2 octets représenté en hexadécimal par 0x01 0x02, et que le paramètre column est configuré comme suit : [{"name":"col1","type":"JSON_HEX"}], la valeur écrite dans l'enregistrement Kafka est la chaîne {"col1":"0102"}.