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. |
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_insertignore la sérialisation Avro, la compaction et la déduplication. Garantissez l'unicité des données source avant d'utiliser ce mode.bulk_insertn'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 |
Lorsqueindex.typeest défini surBUCKET, le paramètreindex.global.enabled=truen'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. |
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
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>
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");
}
}