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.
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 %.
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}ImportantCette 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.
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 |
|
||
|
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
Pour plus d'informations sur la procédure, consultez les rubriques Configurer une tâche de synchronisation hors ligne dans l'interface sans code et Configurer une tâche de synchronisation hors ligne dans l'éditeur de code.
Pour la liste complète des paramètres et un exemple de script pour l'éditeur de code, consultez l'Annexe : Exemples de scripts et descriptions des paramètres.
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.
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.
|
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 :
|
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
|
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.
|
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.
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 : 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,
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 :
|
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.
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 |
|
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 |
|
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 |
|
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 |
|
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 |
|
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 |
|
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 |