Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Hudi (en cours de retrait)

Dernière mise à jour :Aug 21, 2026
Important

Le connecteur Hudi intégré n'est plus pris en charge dans les futures versions de Ververica Runtime (VVR). Utilisez des connecteurs personnalisés pour connecter Realtime Compute for Apache Flink à Apache Hudi, ou migrez vers le connecteur Paimon afin de bénéficier de fonctionnalités et de performances optimisées.

Apache Hudi est un framework open source de lac de données qui gère les données de table stockées dans Object Storage Service (OSS) ou Hadoop Distributed File System (HDFS). Il offre des garanties ACID (Atomicité, Cohérence, Isolation, Durabilité), des opérations upsert et delete au niveau des lignes, une gestion automatique des petits fichiers et des requêtes de voyage dans le temps.

Fonctionnalités principales

Fonctionnalité Description
Sémantique ACID Isolation par snapshot par défaut, garantissant la cohérence des données lors des lectures et écritures concurrentes.
Sémantique UPSERT Combine INSERT et UPDATE : si l'enregistrement n'existe pas, il est inséré ; s'il existe déjà, il est mis à jour. Cela simplifie le code de développement ETL.
Voyage dans le temps Accès aux versions historiques des données à un instant donné, permettant un audit efficace des données et un contrôle qualité.

Scénarios typiques

Scénario Description
Accélération de l'ingestion de base de données Écrivez directement les données Change Data Capture (CDC) (par exemple, les journaux binaires MySQL via le connecteur MySQL CDC) dans une table Hudi pour l'ETL en temps réel en aval. Cette approche est plus rentable que le chargement groupé hors ligne.
ETL incrémentiel Extrayez les flux de données de modification depuis Hudi de manière incrémentielle pour un ETL léger en temps réel. Utilisez Apache Presto ou Apache Spark pour l'OLAP en aval.
File d'attente de messages Utilisez Hudi comme remplacement léger d'une file d'attente de messages pour les scénarios à faible volume, simplifiant ainsi l'architecture applicative.
Rétroalimentation des données Joignez les données complètes et les données incrémentielles des tables Hudi dans un metastore Hive pour générer des tables larges avec une surcharge de calcul minimale.

Avantages par rapport à Hudi open source

  • Sans maintenance : Le connecteur Hudi intégré réduit la complexité opérationnelle et fournit des garanties de SLA.

  • Connectivité des données améliorée : Découple les données des moteurs de calcul, permettant une migration transparente entre Apache Flink, Apache Spark, Apache Presto et Apache Hive.

  • Ingestion simplifiée de la base de données vers le lac : Fonctionne avec le connecteur Flink CDC pour rationaliser le développement des données.

  • Fonctionnalités de classe entreprise : Gestion unifiée des métadonnées via Data Lake Formation (DLF) et modifications légères automatiques du schéma.

  • Stockage rentable : Données stockées au format Apache Parquet ou Apache Avro dans Alibaba Cloud OSS, avec isolation du stockage et du calcul pour une mise à l'échelle flexible des ressources.

Configurations prises en charge

Élément Valeur
Type de table Table source, table de destination
Mode d'exécution Mode streaming, mode batch
Format de données N/A
Type d'API DataStream API, SQL API
Mise à jour/suppression des données dans la destination Pris en charge
Version minimale de VVR vvr-4.0.11-flink-1.13
Systèmes de fichiers pris en charge OSS, HDFS, OSS-HDFS

Métriques

Type de table Métriques
Table source numRecordsIn, numRecordsInPerSecond
Table de destination numRecordsOut, numRecordsOutPerSecond, currentSendTime

Pour les définitions des métriques, consultez la section Métriques.

Limitations

  • Version minimale du moteur : vvr-4.0.11-flink-1.13 ou ultérieure.

  • Systèmes de fichiers pris en charge : OSS, HDFS ou OSS-HDFS uniquement.

  • Les jobs brouillons ne peuvent pas s'exécuter sur des clusters de session.

  • Les modifications de champs ne sont pas prises en charge via le connecteur Hudi. Pour modifier les champs, utilisez des instructions Spark SQL dans la console Data Lake Formation (DLF).

Syntaxe

CREATE TEMPORARY TABLE hudi_tbl (
  uuid BIGINT,
  data STRING,
  ts   TIMESTAMP(3),
  PRIMARY KEY(uuid) NOT ENFORCED
) WITH (
  'connector' = 'hudi',
  'path' = 'oss://<yourOSSBucket>/<Custom storage directory>',
  ...
);

Paramètres de la clause WITH

Paramètres de base

Paramètres communs

Paramètre Obligatoire Par défaut Description
connector Oui Définissez sur hudi.
path Oui Chemin de stockage de la table. Formats pris en charge : OSS (oss://<bucket>/<user-defined-dir>), HDFS (hdfs://<user-defined-dir>), OSS-HDFS (oss://<bucket>.<oss-hdfs-endpoint>/<user-defined-dir>). Les chemins OSS-HDFS nécessitent VVR 8.0.3 ou ultérieur. Trouvez l'endpoint OSS-HDFS dans la section Port de la page Overview du bucket OSS.
hoodie.datasource.write.recordkey.field Non uuid Champ de clé primaire. Séparez plusieurs champs par des virgules. Vous pouvez également utiliser la syntaxe PRIMARY KEY dans le DDL.
precombine.field Non ts Champ de version utilisé pour déterminer l'ordre de mise à jour. S'il n'est pas défini, les mises à jour suivent la séquence de messages définie par le moteur.
oss.endpoint Non Requis lors du stockage des données dans OSS ou OSS-HDFS. Pour les endpoints OSS, consultez la section Régions et endpoints. Pour les endpoints OSS-HDFS, consultez la section Port de la page Overview du bucket OSS.
accessKeyId Non ID AccessKey. Requis pour OSS et OSS-HDFS. Stockez les identifiants sous forme de variables plutôt que de les coder en dur. Consultez la section Gérer les variables.
accessKeySecret Non Secret AccessKey. Requis pour OSS et OSS-HDFS.
Important

Pour protéger votre paire AccessKey, stockez l'AccessKey ID et l'AccessKey secret sous forme de variables. Consultez la section Gérer les variables.

Paramètres de la table source

Paramètre Obligatoire Par défaut Description
read.streaming.enabled Non false Définissez sur true pour activer la lecture en streaming. Par défaut, la lecture par snapshot est utilisée, ce qui renvoie le dernier snapshot complet.
read.start-commit Non (vide) Offset de départ pour la lecture en streaming. Format : yyyyMMddHHmmss pour une heure spécifique, ou earliest pour lire depuis le début. Laissez vide pour lire à partir du dernier commit.

Paramètres de la table de destination

Paramètre Obligatoire Par défaut Description
write.operation Non UPSERT Mode d'écriture. Valeurs valides : insert (ajout), upsert (insertion ou mise à jour), bulk_insert (ajout groupé).
hive_sync.enable Non false Définissez sur true pour synchroniser les métadonnées avec Apache Hive.
hive_sync.mode Non hms Mode de synchronisation. hms synchronise avec un metastore Hive ou DLF. jdbc synchronise via un pilote Java Database Connectivity (JDBC).
hive_sync.db Non default Nom de la base de données Hive cible.
hive_sync.table Non Nom de la table actuelle Nom de la table Hive cible. Ne doit pas contenir de traits d'union (-).
dlf.catalog.region Non Région où DLF est activé. Prend effet uniquement lorsque hive_sync.mode est défini sur hms. Consultez la section Régions et endpoints pris en charge. Doit correspondre à la région spécifiée par dlf.catalog.endpoint.
dlf.catalog.endpoint Non Endpoint DLF. Prend effet uniquement lorsque hive_sync.mode est défini sur hms. Utilisez l'endpoint VPC pour une latence plus faible, par exemple dlf-vpc.cn-hangzhou.aliyuncs.com pour la région Chine (Hangzhou). Consultez la section Régions et endpoints pris en charge. Pour l'accès inter-VPC, consultez la section Comment Realtime Compute for Apache Flink accède-t-il à un service entre différents VPC ?

Paramètres avancés

Paramètres de parallélisme

Paramètre Par défaut Description
write.tasks 4 Parallélisme des tâches d'écriture. Chaque tâche écrit séquentiellement dans un ou plusieurs buckets. L'augmentation de cette valeur n'augmente pas le nombre de petits fichiers.
write.bucket_assign.tasks Parallélisme du déploiement Parallélisme des opérateurs d'assignation de bucket. L'augmentation de cette valeur augmente le nombre de petits fichiers.
write.index_bootstrap.tasks Parallélisme du déploiement Parallélisme des opérateurs d'amorçage d'index. Prend effet uniquement lorsque index.bootstrap.enabled est défini sur true. L'augmentation de cette valeur améliore le débit d'amorçage, mais le checkpointing peut être bloqué pendant l'amorçage ; augmentez la tolérance aux échecs de checkpoint si nécessaire.
read.tasks 4 Parallélisme des opérateurs de lecture en streaming et par lot.
compaction.tasks 4 Parallélisme des opérateurs de compaction en ligne. La compaction en ligne consomme plus de ressources que la compaction hors ligne ; privilégiez la compaction hors ligne pour les charges de travail de production.

Paramètres de compaction en ligne

Paramètre Par défaut Description
compaction.schedule.enabled true Indique s'il faut générer des plans de compaction selon un calendrier. Maintenez cette valeur sur true même lorsque la compaction asynchrone est désactivée, afin que la compaction hors ligne puisse exécuter les plans planifiés.
compaction.async.enabled true Indique s'il faut exécuter la compaction de manière asynchrone. Définissez sur false pour désactiver la compaction en ligne tout en maintenant la génération de plans active.
compaction.tasks 4 Parallélisme des tâches de compaction.
compaction.trigger.strategy num_commits Stratégie utilisée pour déclencher la compaction. Valeurs valides : num_commits, time_elapsed, num_and_time, num_or_time.
compaction.delta_commits 5 Nombre de commits requis pour déclencher la compaction. Utilisé avec num_commits, num_and_time ou num_or_time.
compaction.delta_seconds 3600 Intervalle en secondes entre les déclenchements de compaction. Utilisé avec time_elapsed, num_and_time ou num_or_time.
compaction.max_memory 100 MB Mémoire maximale pour la table de hachage utilisée lors de la compaction et de la déduplication. Augmentez jusqu'à 1 Go si les ressources le permettent.
compaction.target_io 500 GB Débit d'E/S maximal par plan de compaction.

Paramètres de taille de fichier

Ces paramètres contrôlent la manière dont Hudi gère la taille des fichiers pour éviter l'accumulation de petits fichiers.

Paramètre Par défaut Description
hoodie.parquet.max.file.size 120 Mo (120 × 1024 × 1024 octets) Taille maximale d'un fichier Parquet. Les données dépassant ce seuil sont écrites dans un nouveau groupe de fichiers.
hoodie.parquet.small.file.limit 100 Mo (104 857 600 octets) Les fichiers inférieurs à ce seuil sont traités comme des petits fichiers. Lors des écritures, Hudi ajoute aux petits fichiers existants au lieu d'en créer de nouveaux.
hoodie.copyonwrite.record.size.estimate 1 Ko (1 024 octets) Taille estimée de l'enregistrement. Si elle n'est pas définie, Hudi calcule cette valeur dynamiquement à partir des métadonnées validées.

Paramètres de configuration Hadoop

Paramètre Par défaut Description
hadoop.${option key} Éléments de configuration Hadoop, spécifiés avec le préfixe hadoop.. Pris en charge dans Hudi 0.12.0 et versions ultérieures. Utilisez des instructions DDL pour spécifier les configurations Hadoop par job pour les scénarios inter-clusters. Plusieurs éléments peuvent être spécifiés simultanément.

Paramètres d'écriture de données

Écriture par lot

Utilisez l'écriture par lot pour importer des données existantes provenant d'autres sources dans une table Hudi.

bulk_insert ignore la sérialisation Avro, la compaction et la déduplication. Garantissez l'unicité des données source avant d'utiliser ce mode. bulk_insert n'est valide qu'en mode d'exécution par lot.
Paramètre Par défaut Description
write.operation upsert Type d'écriture. Définissez sur bulk_insert pour les écritures par lot.
write.tasks Parallélisme du déploiement Parallélisme pour les tâches bulk_insert. Le nombre final de fichiers de sortie est supérieur ou égal à cette valeur (les données basculent vers un nouveau fichier lorsque la limite Parquet de 120 Mo est atteinte).
write.bulk_insert.shuffle_input true Indique s'il faut mélanger les données d'entrée par champ de partition avant l'écriture. Disponible dans Hudi 0.11.0 et versions ultérieures. Réduit le nombre de petits fichiers mais peut provoquer un déséquilibre des données.
write.bulk_insert.sort_input true Indique s'il faut trier les données d'entrée par champ de partition avant l'écriture. Disponible dans Hudi 0.11.0 et versions ultérieures. Réduit le nombre de petits fichiers lorsqu'une seule tâche écrit dans plusieurs partitions.
write.sort.memory 128 Mémoire gérée disponible pour l'opérateur de tri, en Mo.

Mode changelog

En mode changelog, Hudi conserve tous les événements de modification : INSERT, UPDATE_BEFORE, UPDATE_AFTER et DELETE, permettant un entrepôt de données quasi temps réel de bout en bout avec le calcul stateful de Flink. Les tables Merge On Read (MOR) prennent en charge ce mode.

En mode non-changelog, les modifications intermédiaires au sein d'un lot sont fusionnées. La lecture par snapshot renvoie uniquement le résultat fusionné final ; les états intermédiaires ne sont pas visibles, quel que soit le chemin d'écriture.

Après avoir activé le mode changelog, la tâche de compaction asynchrone fusionne toujours les modifications intermédiaires. Définissez compaction.delta_commits=5 et compaction.delta_seconds=3600 pour laisser suffisamment de temps aux consommateurs en aval pour lire les enregistrements avant qu'ils ne soient compactés.

Paramètre Par défaut Description
changelog.enabled false Définissez sur true pour conserver tous les événements de modification. Lorsque la valeur est false, seul l'enregistrement fusionné final est garanti ; les modifications intermédiaires peuvent être fusionnées.

Mode append

Pris en charge dans Hudi 0.10.0 et versions ultérieures.

  • Tables MOR : La politique des petits fichiers s'applique. Les données sont écrites dans des fichiers journaux Apache Avro en mode append.

  • Tables Copy On Write (COW) : La politique des petits fichiers ne s'applique pas. Un nouveau fichier Apache Parquet est créé pour chaque écriture.

Paramètres de clustering

Hudi prend en charge le clustering pour résoudre l'accumulation de petits fichiers en mode INSERT.

Clustering en ligne (tables COW uniquement)

Paramètre Par défaut Description
write.insert.cluster false Définissez sur true pour fusionner les petits fichiers lors des écritures. Chaque opération INSERT fusionne les petits fichiers existants, mais aucune déduplication n'est effectuée et le débit d'écriture diminue.

Clustering asynchrone (Hudi 0.12.0 et versions ultérieures)

Paramètre Par défaut Description
clustering.schedule.enabled false Définissez sur true pour planifier périodiquement un plan de clustering.
clustering.delta_commits 4 Nombre de commits requis pour générer un plan de clustering. Prend effet uniquement lorsque clustering.schedule.enabled est défini sur true.
clustering.async.enabled false Définissez sur true pour exécuter le plan de clustering de manière asynchrone à intervalles réguliers.
clustering.tasks 4 Parallélisme des tâches de clustering.
clustering.plan.strategy.target.file.max.bytes 1 GiB (1 073 741 824 octets) Taille de fichier maximale cible pour la sortie du clustering.
clustering.plan.strategy.small.file.limit 600 Les fichiers inférieurs à ce seuil (en octets) sont éligibles au clustering.
clustering.plan.strategy.sort.columns Colonnes utilisées pour trier les données lors du clustering.

Stratégies de plan de clustering

Paramètre Par défaut Description
clustering.plan.partition.filter.mode NONE Mode de filtre de partition. Valeurs valides : NONE (toutes les partitions), RECENT_DAYS (partitions des N derniers jours), SELECTED_PARTITIONS (partitions spécifiques).
clustering.plan.strategy.daybased.lookback.partitions 2 Nombre de jours récents pour sélectionner les partitions pour le clustering. Prend effet uniquement lorsque filter.mode est défini sur RECENT_DAYS.
clustering.plan.strategy.cluster.begin.partition Partition de début pour le filtrage par plage. Prend effet uniquement lorsque filter.mode est défini sur SELECTED_PARTITIONS.
clustering.plan.strategy.cluster.end.partition Partition de fin pour le filtrage par plage. Prend effet uniquement lorsque filter.mode est défini sur SELECTED_PARTITIONS.
clustering.plan.strategy.partition.regex.pattern Expression régulière pour la sélection des partitions.
clustering.plan.strategy.partition.selected Liste séparée par des virgules des partitions sélectionnées.

Choisir un type d'index

Hudi prend en charge deux types d'index. Utilisez ce tableau pour choisir celui qui correspond à votre charge de travail.

Dimension FLINK_STATE BUCKET
Surcharge de stockage/calcul Oui (backend d'état) Aucune
Performance Dépend du backend d'état Meilleure (aucune surcharge d'état)
Flexibilité du groupe de fichiers Attribue dynamiquement les enregistrements en fonction de la taille du fichier Nombre fixe de buckets (ne peut pas être augmenté après la configuration initiale)
Modifications inter-partitions Pris en charge Non pris en charge (exception : entrée en streaming Change Data Capture (CDC))
Quand l'utiliser Tables contenant moins de 500 millions d'enregistrements, ou charges de travail nécessitant des mises à jour inter-partitions Tables contenant plus de 500 millions d'enregistrements où la surcharge d'état constitue un goulot d'étranglement
Lorsque index.type est défini sur BUCKET , le paramètre index.global.enabled=true n'a aucun effet : l'index bucket ne prend pas en charge la déduplication inter-partitions.
Paramètre Par défaut Description
index.type FLINK_STATE Type d'index. Valeurs valides : FLINK_STATE, BUCKET.
hoodie.bucket.index.hash.field Clé primaire Champ de clé de hachage pour l'index bucket. Peut être un sous-ensemble de la clé primaire.
hoodie.bucket.index.num.buckets 4 Nombre de buckets par partition. Ne peut pas être modifié après la création de la table.
Les paramètres d'index bucket sont pris en charge dans Hudi 0.11.0 et versions ultérieures.

Paramètres de lecture de données

Hudi prend en charge trois modèles de lecture utilisant le même ensemble de paramètres.

Modèle Configuration
Lecture en streaming Définissez read.streaming.enabled=true et éventuellement read.start-commit
Lecture incrémentielle par lot Définissez à la fois read.start-commit et read.end-commit ; l'intervalle est fermé (inclusif aux deux extrémités)
Voyage dans le temps Définissez uniquement read.end-commit ; lit un snapshot à ce commit spécifique

Paramètres de lecture en streaming

Par défaut, la lecture d'une table Hudi utilise la lecture par snapshot : le dernier snapshot complet est renvoyé en une seule fois. Définissez read.streaming.enabled=true pour passer à la lecture en streaming.

Paramètre Par défaut Description
read.streaming.enabled false Définissez sur true pour activer la lecture en streaming.
read.start-commit (vide) Offset de départ. Format : yyyyMMddHHmmss pour une heure spécifique, ou earliest pour lire depuis le début. Laissez vide pour commencer à partir du dernier commit.
clean.retain_commits 30 Nombre maximal de commits historiques conservés par le nettoyeur. Les commits dépassant cette limite sont supprimés. Par exemple, avec un intervalle de checkpointing de 5 minutes, la valeur par défaut de 30 conserve les journaux de modifications pendant au moins 150 minutes.
Important

La lecture en streaming des journaux de modifications nécessite Hudi 0.10.0 ou version ultérieure. Les tâches de compaction peuvent fusionner les journaux de modifications, supprimant les enregistrements intermédiaires et affectant potentiellement les calculs en aval.

Paramètres de lecture incrémentielle

Paramètre Par défaut Description
read.start-commit Dernier commit Début de la plage de lecture, au format yyyyMMddHHmmss.
read.end-commit Dernier commit Fin de la plage de lecture, au format yyyyMMddHHmmss. La plage est fermée (les deux points de terminaison sont inclusifs).

Exemples

Table source

CREATE TEMPORARY TABLE blackhole (
  id INT NOT NULL PRIMARY KEY NOT ENFORCED,
  data STRING,
  ts TIMESTAMP(3)
) WITH (
  'connector' = 'blackhole'
);

CREATE TEMPORARY TABLE hudi_tbl (
  id INT NOT NULL PRIMARY KEY NOT ENFORCED,
  data STRING,
  ts TIMESTAMP(3)
) WITH (
  'connector' = 'hudi',
  'oss.endpoint' = '<yourOSSEndpoint>',
  'accessKeyId' = '${secret_values.ak_id}',
  'accessKeySecret' = '${secret_values.ak_secret}',
  'path' = 'oss://<yourOSSBucket>/<Custom storage directory>',
  'table.type' = 'MERGE_ON_READ',
  'read.streaming.enabled' = 'true'
);

-- Read from the latest commit in streaming mode and write to Blackhole.
INSERT INTO blackhole SELECT * FROM hudi_tbl;

Table de destination

CREATE TEMPORARY TABLE datagen (
  id INT NOT NULL PRIMARY KEY NOT ENFORCED,
  data STRING,
  ts TIMESTAMP(3)
) WITH (
  'connector' = 'datagen',
  'rows-per-second' = '100'
);

CREATE TEMPORARY TABLE hudi_tbl (
  id INT NOT NULL PRIMARY KEY NOT ENFORCED,
  data STRING,
  ts TIMESTAMP(3)
) WITH (
  'connector' = 'hudi',
  'oss.endpoint' = '<yourOSSEndpoint>',
  'accessKeyId' = '${secret_values.ak_id}',
  'accessKeySecret' = '${secret_values.ak_secret}',
  'path' = 'oss://<yourOSSBucket>/<Custom storage directory>',
  'table.type' = 'MERGE_ON_READ'
);

INSERT INTO hudi_tbl SELECT * FROM datagen;

DataStream API

Important

Pour utiliser la DataStream API, configurez un connecteur DataStream pour Realtime Compute for Apache Flink. Consultez la section Paramètres des connecteurs DataStream.

Dépendances Maven

Faites correspondre les versions des dépendances à votre version de VVR.

<properties>
  <maven.compiler.source>8</maven.compiler.source>
  <maven.compiler.target>8</maven.compiler.target>
  <flink.version>1.15.4</flink.version>
  <hudi.version>0.13.1</hudi.version>
</properties>

<dependencies>
  <!-- Flink -->
  <dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-streaming-java</artifactId>
    <version>${flink.version}</version>
    <scope>provided</scope>
  </dependency>
  <dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-table-common</artifactId>
    <version>${flink.version}</version>
    <scope>provided</scope>
  </dependency>
  <dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-table-api-java-bridge</artifactId>
    <version>${flink.version}</version>
    <scope>provided</scope>
  </dependency>
  <dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-table-planner_2.12</artifactId>
    <version>${flink.version}</version>
    <scope>provided</scope>
  </dependency>
  <dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-clients</artifactId>
    <version>${flink.version}</version>
    <scope>provided</scope>
  </dependency>

  <!-- Hudi -->
  <dependency>
    <groupId>org.apache.hudi</groupId>
    <artifactId>hudi-flink1.15-bundle</artifactId>
    <version>${hudi.version}</version>
    <scope>provided</scope>
  </dependency>

  <!-- OSS -->
  <dependency>
    <groupId>org.apache.hadoop</groupId>
    <artifactId>hadoop-common</artifactId>
    <version>3.3.2</version>
    <scope>provided</scope>
  </dependency>
  <dependency>
    <groupId>org.apache.hadoop</groupId>
    <artifactId>hadoop-aliyun</artifactId>
    <version>3.3.2</version>
    <scope>provided</scope>
  </dependency>

  <!-- DLF -->
  <dependency>
    <groupId>com.aliyun.datalake</groupId>
    <artifactId>metastore-client-hive2</artifactId>
    <version>0.2.14</version>
    <scope>provided</scope>
  </dependency>
  <dependency>
    <groupId>org.apache.hadoop</groupId>
    <artifactId>hadoop-mapreduce-client-core</artifactId>
    <version>2.5.1</version>
    <scope>provided</scope>
  </dependency>
</dependencies>
Important

Les dépendances DLF entrent en conflit avec les versions open source d'Apache Hive (hive-common, hive-exec). Pour les tests DLF locaux, téléchargez les packages JAR personnalisés hive-common et hive-exec et importez-les manuellement dans IntelliJ IDEA.

Écrire des données dans Hudi

L'exemple suivant écrit des données dans une table Hudi MOR dans OSS et synchronise éventuellement les métadonnées avec DLF.

import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.data.GenericRowData;
import org.apache.flink.table.data.RowData;
import org.apache.flink.table.data.StringData;
import org.apache.hudi.common.model.HoodieTableType;
import org.apache.hudi.configuration.FlinkOptions;
import org.apache.hudi.util.HoodiePipeline;

import java.util.HashMap;
import java.util.Map;

public class FlinkHudiQuickStart {

  public static void main(String[] args) throws Exception {
    StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

    String dbName = "test_db";
    String tableName = "test_tbl";
    String basePath = "oss://xxx";

    Map<String, String> options = new HashMap<>();

    // Hudi configuration
    options.put(FlinkOptions.PATH.key(), basePath);
    options.put(FlinkOptions.TABLE_TYPE.key(), HoodieTableType.MERGE_ON_READ.name());
    options.put(FlinkOptions.PRECOMBINE_FIELD.key(), "ts");
    options.put(FlinkOptions.DATABASE_NAME.key(), dbName);
    options.put(FlinkOptions.TABLE_NAME.key(), tableName);

    // OSS configuration
    // Use the public endpoint for local debugging (e.g., oss-cn-hangzhou.aliyuncs.com)
    // Use the internal endpoint for cluster submission (e.g., oss-cn-hangzhou-internal.aliyuncs.com)
    options.put("hadoop.fs.oss.accessKeyId", "xxx");
    options.put("hadoop.fs.oss.accessKeySecret", "xxx");
    options.put("hadoop.fs.oss.endpoint", "xxx");
    options.put("hadoop.fs.AbstractFileSystem.oss.impl", "org.apache.hadoop.fs.aliyun.oss.OSS");
    options.put("hadoop.fs.oss.impl", "org.apache.hadoop.fs.aliyun.oss.AliyunOSSFileSystem");

    // DLF configuration (optional — remove if not syncing to DLF)
    // Use the public endpoint for local debugging (e.g., dlf.cn-hangzhou.aliyuncs.com)
    // Use the VPC endpoint for cluster submission (e.g., dlf-vpc.cn-hangzhou.aliyuncs.com)
    options.put(FlinkOptions.HIVE_SYNC_ENABLED.key(), "true");
    options.put(FlinkOptions.HIVE_SYNC_MODE.key(), "hms");
    options.put(FlinkOptions.HIVE_SYNC_DB.key(), dbName);
    options.put(FlinkOptions.HIVE_SYNC_TABLE.key(), tableName);
    options.put("hadoop.dlf.catalog.id", "xxx");
    options.put("hadoop.dlf.catalog.accessKeyId", "xxx");
    options.put("hadoop.dlf.catalog.accessKeySecret", "xxx");
    options.put("hadoop.dlf.catalog.region", "xxx");
    options.put("hadoop.dlf.catalog.endpoint", "xxx");
    options.put("hadoop.hive.imetastoreclient.factory.class",
        "com.aliyun.datalake.metastore.hive2.DlfMetaStoreClientFactory");

    DataStream<RowData> dataStream = env.fromElements(
        GenericRowData.of(StringData.fromString("id1"), StringData.fromString("name1"), 22,
            StringData.fromString("1001"), StringData.fromString("p1")),
        GenericRowData.of(StringData.fromString("id2"), StringData.fromString("name2"), 32,
            StringData.fromString("1002"), StringData.fromString("p2"))
    );

    HoodiePipeline.Builder builder = HoodiePipeline.builder(tableName)
        .column("uuid string")
        .column("name string")
        .column("age int")
        .column("ts string")
        .column("`partition` string")
        .pk("uuid")
        .partition("partition")
        .options(options);

    // Second parameter: whether the input stream is bounded (true = batch, false = streaming)
    builder.sink(dataStream, false);
    env.execute("Flink_Hudi_Quick_Start");
  }
}

FAQ