Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Connecteur DataHub

Dernière mise à jour :Aug 09, 2026

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.

Remarque

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

Compatibilité avec Kafka

.

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 :

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é.

Important

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.

Remarque

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

Important

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