Le connecteur DataHub permet de lire des données en continu depuis Alibaba Cloud DataHub vers des jobs Flink et d'écrire les résultats traités dans des rubriques DataHub. Il prend en charge à la fois Flink SQL et l'API DataStream.
DataHub est compatible avec le protocole Kafka. Pour connecter Flink à DataHub via le protocole Kafka, utilisez le connecteur Kafka standard, et non le connecteur Upsert Kafka. Pour plus de détails, consultez
.
Fonctionnalités
| Élément | Description |
|---|---|
| Type pris en charge | Source et destination (sink) |
| Mode d'exécution | Streaming et batch |
| Format des données | N/A |
| Métriques | N/A |
| Type d'API | DataStream et SQL |
| Prise en charge de la mise à jour/suppression des données dans la destination | Non prise en charge. La destination écrit uniquement des lignes d'insertion dans la rubrique cible. |
Prérequis
Avant de commencer, assurez-vous de disposer des éléments suivants :
Un projet et une rubrique DataHub. Consultez Prise en main de DataHub.
Un abonnement DataHub (requis pour la source). Consultez Créer un abonnement.
Un ID AccessKey et un secret AccessKey Alibaba Cloud. Consultez Opérations dans la console.
Limites
Il est déconseillé de lire une table source DataHub avec un job batch. En mode batch, une table source DataHub ne peut pas atteindre l'état terminé et le job reste en attente au lieu de se terminer.
Syntaxe
CREATE TEMPORARY TABLE datahub_input (
`time` BIGINT,
`sequence` STRING METADATA VIRTUAL,
`shard-id` BIGINT METADATA VIRTUAL,
`system-time` TIMESTAMP METADATA VIRTUAL
) WITH (
'connector' = 'datahub',
'subId' = '<yourSubId>',
'endPoint' = '<yourEndPoint>',
'project' = '<yourProjectName>',
'topic' = '<yourTopicName>',
'accessId' = '${secret_values.ak_id}',
'accessKey' = '${secret_values.ak_secret}'
);
Options du connecteur
Général
| Option | Type | Requis | Par défaut | Description |
|---|---|---|---|---|
| connector | String | Oui | (aucun) | Le type de connecteur. Définissez cette option sur datahub. |
| endPoint | String | Oui | (aucun) | L'endpoint du projet DataHub. La valeur varie selon la région. Consultez Endpoints. |
| project | String | Oui | (aucun) | Le nom du projet DataHub. |
| topic | String | Oui | (aucun) | Le nom de la rubrique DataHub. Pour les rubriques BLOB (données non typées et non structurées), la table Flink doit contenir exactement une colonne VARBINARY. |
| accessId | String | Oui | (aucun) | L'ID AccessKey de votre compte Alibaba Cloud. Stockez-le en tant que variable plutôt que de le coder en dur. Consultez Gérer les variables. |
| accessKey | String | Oui | (aucun) | Le secret AccessKey de votre compte Alibaba Cloud. |
| retryTimeout | Integer | Non | 1800000 | Le délai d'expiration maximal pour une nouvelle tentative, en millisecondes. |
| retryInterval | Integer | Non | 1000 | L'intervalle entre les nouvelles tentatives, en millisecondes. |
| CompressType | String | Non | lz4 | L'algorithme de compression pour les lectures et les écritures. Valeurs valides : lz4, deflate, "" (désactivé). Nécessite VVR 6.0.5 ou version ultérieure. |
Spécifique à la source
| Option | Type | Requis | Par défaut | Description |
|---|---|---|---|---|
| subId | String | Oui | (aucun) | L'ID de l'abonnement DataHub. |
| maxFetchSize | Integer | Non | 50 | Le nombre d'enregistrements récupérés par requête. Augmentez cette valeur pour améliorer le débit de lecture. |
| maxBufferSize | Integer | Non | 50 | Le nombre maximal d'enregistrements mis en cache lors des lectures asynchrones. Augmentez cette valeur pour améliorer le débit de lecture. |
| fetchLatestDelay | Integer | Non | 500 | La durée d'attente en millisecondes lorsqu'aucune donnée n'est disponible. Diminuez cette valeur pour réduire la latence de lecture pour les rubriques à faible trafic. |
| lengthCheck | String | Non | NONE | La règle de gestion des lignes où le nombre de champs analysés ne correspond pas au nombre de colonnes définies. Valeurs valides : NONE, SKIP, EXCEPTION, PAD. Consultez Règles de validation du nombre de champs. |
| columnErrorDebug | Boolean | Non | false | Indique s'il faut activer la journalisation de débogage pour les erreurs d'analyse des champs. Définissez sur true pour afficher les journaux d'exceptions d'analyse. |
| startTime | String | Non | (aucun) | L'horodatage à partir duquel commencer la consommation. Format : yyyy-MM-dd hh:mm:ss. |
| endTime | String | Non | (aucun) | L'horodatage auquel arrêter la consommation. Format : yyyy-MM-dd hh:mm:ss. |
| startTimeMs | Long | Non | -1 | L'horodatage à partir duquel commencer la consommation, en millisecondes. Cette option a priorité sur startTime. Consultez Position de début de consommation. |
Position de début de consommation
L'option startTimeMs contrôle l'emplacement à partir duquel la source commence la lecture :
-1 (par défaut) : Commence à partir du dernier offset de la rubrique. Si aucun offset n'existe, revient au premier offset.
Un horodatage spécifique : Commence à partir du premier enregistrement à ou après l'horodatage spécifié.
La valeur par défaut -1 peut entraîner une perte de données. Si votre job échoue avant son premier point de contrôle, le dernier offset de la rubrique peut avoir avancé et les enregistrements écrits pendant cette fenêtre sont ignorés. Définissez explicitement startTimeMs sur un horodatage spécifique pour contrôler la position de départ.
Règles de validation du nombre de champs
L'option lengthCheck détermine ce qui se produit lorsque le nombre de champs analysés dans une ligne ne correspond pas au nombre de colonnes définies :
| Valeur | Comportement |
|---|---|
NONE (par défaut) |
Si le nombre de champs analysés > nombre de colonnes définies : lit de gauche à droite jusqu'au nombre défini. Si le nombre de champs analysés < nombre de colonnes définies : ignore la ligne. |
SKIP |
Ignore les lignes dont le nombre de champs analysés diffère du nombre de colonnes définies. |
EXCEPTION |
Lève une exception lorsque le nombre de champs analysés diffère du nombre de colonnes définies. |
PAD |
Lit de gauche à droite. Si le nombre de champs analysés > nombre de colonnes définies : lit de gauche à droite jusqu'au nombre défini. Si le nombre de champs analysés < nombre de colonnes définies : complète les champs manquants avec null. |
Spécifique à la destination (sink)
| Option | Type | Requis | Par défaut | Description |
|---|---|---|---|---|
| batchCount | Integer | Non | 500 | Le nombre maximal de lignes par lot d'écriture. |
| batchSize | Integer | Non | 512000 | La taille maximale d'un lot d'écriture, en octets. |
| flushInterval | Integer | Non | 5000 | L'intervalle de vidage, en millisecondes. |
| hashFields | String | Non | null | Une liste séparée par des virgules des noms de colonnes utilisées pour router les lignes vers les shards. Les lignes ayant les mêmes valeurs dans ces colonnes sont écrites dans le même shard. La valeur par défaut (null) utilise des écritures aléatoires. Exemple : hashFields=a,b. |
| timeZone | String | Non | (aucun) | Le fuseau horaire utilisé lors de la conversion des champs TIMESTAMP. |
| schemaVersion | Integer | Non | -1 | La version du schéma dans le registre de schémas enregistré. |
Comportement de vidage des lots d'écriture
Un lot d'écriture est vidé vers DataHub dès qu'une des conditions suivantes est remplie :
Le nombre de lignes mises en tampon atteint
batchCount.La taille totale des données mises en tampon atteint
batchSize.Le temps écoulé depuis le dernier vidage dépasse
flushInterval.
L'augmentation de batchCount, batchSize ou flushInterval améliore le débit d'écriture au prix d'une latence plus élevée.
Mappings de types de données
| Type Flink | Type DataHub |
|---|---|
| TINYINT | TINYINT |
| BOOLEAN | BOOLEAN |
| INTEGER | INTEGER |
| BIGINT | BIGINT |
| BIGINT | TIMESTAMP |
| FLOAT | FLOAT |
| DOUBLE | DOUBLE |
| DECIMAL | DECIMAL |
| VARCHAR | STRING |
| SMALLINT | SMALLINT |
| VARBINARY | BLOB |
Métadonnées
Les champs de métadonnées sont en lecture seule (R). Déclarez-les comme METADATA VIRTUAL dans la définition de la table source pour les inclure dans les requêtes sans les réécrire dans DataHub.
Les champs de métadonnées sont disponibles uniquement lors de l'utilisation de VVR 3.0.1 ou version ultérieure.
| Clé | Type de données | Description | L/E |
|---|---|---|---|
| shard-id | BIGINT METADATA VIRTUAL | L'ID du shard de l'enregistrement. | R |
| sequence | STRING METADATA VIRTUAL | Le numéro de séquence de l'enregistrement dans le shard. | R |
| system-time | TIMESTAMP METADATA VIRTUAL | L'heure à laquelle DataHub a reçu l'enregistrement. | R |
Exemples
Source
L'exemple suivant lit les données d'une rubrique DataHub et les affiche dans la console.
CREATE TEMPORARY TABLE datahub_input (
`time` BIGINT,
`sequence` STRING METADATA VIRTUAL,
`shard-id` BIGINT METADATA VIRTUAL,
`system-time` TIMESTAMP METADATA VIRTUAL
) WITH (
'connector' = 'datahub',
'subId' = '<yourSubId>',
'endPoint' = '<yourEndPoint>',
'project' = '<yourProjectName>',
'topic' = '<yourTopicName>',
'accessId' = '${secret_values.ak_id}',
'accessKey' = '${secret_values.ak_secret}'
);
CREATE TEMPORARY TABLE test_out (
`time` BIGINT,
`sequence` STRING,
`shard-id` BIGINT,
`system-time` TIMESTAMP
) WITH (
'connector' = 'print',
'logger' = 'true'
);
INSERT INTO test_out
SELECT
`time`,
`sequence`,
`shard-id`,
`system-time`
FROM datahub_input;
Destination (Sink)
L'exemple suivant lit à partir d'une rubrique DataHub, convertit le champ name en minuscules et écrit les résultats dans une autre rubrique DataHub.
CREATE TEMPORARY TABLE datahub_source (
name VARCHAR
) WITH (
'connector' = 'datahub',
'endPoint' = '<endPoint>',
'project' = '<yourProjectName>',
'topic' = '<yourTopicName>',
'subId' = '<yourSubId>',
'accessId' = '${secret_values.ak_id}',
'accessKey' = '${secret_values.ak_secret}',
'startTime' = '2018-06-01 00:00:00'
);
CREATE TEMPORARY TABLE datahub_sink (
name VARCHAR
) WITH (
'connector' = 'datahub',
'endPoint' = '<endPoint>',
'project' = '<yourProjectName>',
'topic' = '<yourTopicName>',
'accessId' = '${secret_values.ak_id}',
'accessKey' = '${secret_values.ak_secret}',
'batchSize' = '512000',
'batchCount' = '500'
);
INSERT INTO datahub_sink
SELECT
LOWER(name)
FROM datahub_source;
API DataStream
Pour utiliser l'API DataStream avec DataHub, configurez un connecteur DataStream pour Realtime Compute for Apache Flink. Consultez Paramètres des connecteurs DataStream.
Lire depuis DataHub
VVR fournit la classe DatahubSourceFunction, qui implémente l'interface SourceFunction de Flink.
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1);
// Configure the DataHub source
DatahubSourceFunction datahubSource =
new DatahubSourceFunction(
<yourEndPoint>,
<yourProjectName>,
<yourTopicName>,
<yourSubId>,
<yourAccessId>,
<yourAccessKey>,
"public",
<yourStartTime>,
<yourEndTime>
);
datahubSource.setRequestTimeout(30 * 1000);
datahubSource.enableExitAfterReadFinished();
env.addSource(datahubSource)
.map((MapFunction<RecordEntry, Tuple2<String, Long>>) this::getStringLongTuple2)
.print();
env.execute();
private Tuple2<String, Long> getStringLongTuple2(RecordEntry recordEntry) {
Tuple2<String, Long> tuple2 = new Tuple2<>();
TupleRecordData recordData = (TupleRecordData) (recordEntry.getRecordData());
tuple2.f0 = (String) recordData.getField(0);
tuple2.f1 = (Long) recordData.getField(1);
return tuple2;
}
Écrire dans DataHub
VVR fournit la classe OutputFormatSinkFunction, qui implémente l'interface DatahubSinkFunction.
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// Configure the DataHub sink
env.generateSequence(0, 100)
.map((MapFunction<Long, RecordEntry>) aLong -> getRecordEntry(aLong, "default:"))
.addSink(
new DatahubSinkFunction<>(
<yourEndPoint>,
<yourProjectName>,
<yourTopicName>,
<yourSubId>,
<yourAccessId>,
<yourAccessKey>,
"public",
<schemaVersion> // If schema registry is enabled, you must specify the schema version. Otherwise, set this to 0.
)
);
env.execute();
private RecordEntry getRecordEntry(Long message, String s) {
RecordSchema recordSchema = new RecordSchema();
recordSchema.addField(new Field("f1", FieldType.STRING));
recordSchema.addField(new Field("f2", FieldType.BIGINT));
recordSchema.addField(new Field("f3", FieldType.DOUBLE));
recordSchema.addField(new Field("f4", FieldType.BOOLEAN));
recordSchema.addField(new Field("f5", FieldType.TIMESTAMP));
recordSchema.addField(new Field("f6", FieldType.DECIMAL));
RecordEntry recordEntry = new RecordEntry();
TupleRecordData recordData = new TupleRecordData(recordSchema);
recordData.setField(0, s + message);
recordData.setField(1, message);
recordEntry.setRecordData(recordData);
return recordEntry;
}
Dépendance Maven
Ajoutez le connecteur DataStream DataHub à votre projet. Toutes les versions disponibles sont répertoriées dans le référentiel central Maven.
<dependency>
<groupId>com.alibaba.ververica</groupId>
<artifactId>ververica-connector-datahub</artifactId>
<version>${vvr-version}</version>
</dependency>
Étapes suivantes
Pour obtenir la liste complète des connecteurs pris en charge par Realtime Compute for Apache Flink, consultez Connecteurs pris en charge.
Pour vous connecter à DataHub en utilisant le connecteur Kafka à la place, consultez Message Queue for Apache Kafka.
Puis-je supprimer une rubrique DataHub en cours de consommation ?