Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Connecteur Apache Iceberg

Dernière mise à jour :Aug 09, 2026

Cette rubrique explique comment utiliser le connecteur Apache Iceberg.

Informations générales

Apache Iceberg est un format de table de lac de données ouvert. Vous pouvez utiliser Apache Iceberg pour créer rapidement votre propre service de stockage de lac de données sur HDFS ou sur OSS cloud, et exploiter les moteurs de calcul de l'écosystème big data open source, tels que Flink, Spark, Hive et Presto, pour effectuer des analyses de lac de données.

Catégorie

Description

Types pris en charge

Table source, table de destination (sink) et ingestion de données vers une table de destination

Modes d'exécution

Mode par lots et mode en continu

Format des données

Sans objet

Métriques spécifiques

Aucune

Types d'API

SQL, job YAML d'ingestion de données

Prend en charge la mise à jour ou la suppression des données dans les tables de destination

Oui

Fonctionnalités

Le connecteur Apache Iceberg offre les fonctionnalités suivantes :

  • Crée un service de stockage de lac de données léger et peu coûteux basé sur HDFS ou sur le stockage d'objets.

  • Fournit une sémantique ACID complète.

  • Prend en charge les requêtes de voyage dans le temps pour accéder aux versions historiques des données.

  • Prend en charge le filtrage efficace des données.

  • Prend en charge l'évolution du schéma.

  • Prend en charge l'évolution du partitionnement.

Remarque

Utilisez les capacités de tolérance aux pannes et de traitement en continu de Flink pour importer de grands volumes de données de journal dans un lac de données Apache Iceberg en temps réel. Vous pouvez ensuite exploiter Flink ou d'autres moteurs d'analyse pour extraire de la valeur de ces données.

Limites

  • Le connecteur Apache Iceberg est pris en charge uniquement dans Realtime Compute for Apache Flink avec Ververica Runtime (VVR) version 4.0.8 ou ultérieure. Le connecteur Apache Iceberg doit être utilisé avec un catalogue Data Lake Formation (DLF). Pour plus d'informations, consultez la section Gérer les catalogues DLF-Legacy.

  • Le connecteur Apache Iceberg prend en charge les formats de table Apache Iceberg v1 et v2. Pour plus d'informations, consultez la section Spécification des tables Iceberg.

    Remarque

    Seul Realtime Compute for Apache Flink utilisant VVR version 8.0.7 ou ultérieure prend en charge le format de table v2.

  • En mode de lecture en continu, vous ne pouvez utiliser que des tables Iceberg en ajout seul (append-only) comme tables sources.

Syntaxe

CREATE TABLE iceberg_table (
  id    BIGINT,
  data  STRING
  PRIMARY KEY(`id`) NOT ENFORCED
)
 PARTITIONED BY (data)
 WITH (
 'connector' = 'iceberg',
  ...
);

Options WITH

Options communes

Paramètre

Description

Type

Obligatoire

Valeur par défaut

Remarques

connector

Le type du connecteur.

String

Oui

Aucune

La valeur doit êtreiceberg.

catalog-name

Le nom du catalogue.

String

Oui

Aucune

Saisissez un nom personnalisé en anglais.

catalog-database

Le nom de la base de données.

String

Oui

default

Le nom de votre base de données dans Data Lake Formation (DLF), par exemple dlf_db.

Remarque

Si vous ne disposez pas de base de données Data Lake Formation (DLF), commencez par en créer une.

io-impl

La classe d'implémentation du système de fichiers distribué.

String

Oui

Aucune

La valeur doit êtreorg.apache.iceberg.aliyun.oss.OSSFileIO.

oss.endpoint

Le point de terminaison d'Alibaba Cloud Object Storage Service (OSS).

String

Non

Aucune

Pour plus d'informations, consultez la section Régions et points de terminaison.

Remarque
  • Nous vous recommandons de définir le paramètre oss.endpoint sur le point de terminaison VPC pour OSS. Par exemple, si vous sélectionnez la région Chine (Hangzhou), définissez le paramètre oss.endpoint sur oss-cn-hangzhou-internal.aliyuncs.com.

  • Si vous devez accéder à OSS entre différents VPC, consultez la section Comment accéder à d'autres services entre différents VPC ?

  • access.key.id : pour VVR 8.0.6 et versions antérieures

  • access-key-id : pour VVR 8.0.7 et versions ultérieures

L'ID AccessKey de votre compte Alibaba Cloud.

String

Oui

Aucune

Pour plus d'informations, consultez la section Comment afficher l'ID AccessKey et la clé secrète AccessKey ?

Important

Pour éviter toute fuite d'informations concernant votre AccessKey, nous vous recommandons d'utiliser des variables pour spécifier les valeurs AccessKey. Pour plus d'informations, consultez la section Variables de projet.

  • access.key.secret : pour VVR 8.0.6 et versions antérieures

  • access-key-secret : pour VVR 8.0.7 et versions ultérieures

La clé secrète AccessKey de votre compte Alibaba Cloud.

String

Oui

Aucune

catalog-impl

Le nom de la classe du catalogue.

String

Oui

Aucune

La valeur doit êtreorg.apache.iceberg.aliyun.dlf.DlfCatalog.

warehouse

Le chemin OSS pour stocker les données de la table.

String

Oui

Aucune

Aucune

dlf.catalog-id

L'ID de votre compte Alibaba Cloud.

String

Oui

Aucune

Vous pouvez obtenir l'ID du compte sur la page Informations utilisateur.

dlf.endpoint

Le point de terminaison de Data Lake Formation (DLF).

String

Oui

Aucune

.

Remarque
  • Nous vous recommandons de définir le paramètre dlf.endpoint sur le point de terminaison VPC de DLF. Par exemple, si vous sélectionnez la région Chine (Hangzhou), définissez le paramètre dlf.endpoint sur dlf-vpc.cn-hangzhou.aliyuncs.com.

  • Si vous devez accéder à DLF entre différents VPC, consultez la section Gestion et opérations de l'espace de travail

dlf.region-id

La région de Data Lake Formation (DLF).

String

Oui

Aucune

.

Remarque

Assurez-vous que la région est identique à celle spécifiée pour le paramètre dlf.endpoint.

Options spécifiques au sink

Paramètre

Description

Type

Obligatoire

Valeur par défaut

Remarques

write.operation

Le mode d'opération d'écriture.

String

Non

upsert

  • upsert (par défaut) : Met à jour les données.

  • insert : Ajoute les données.

  • bulk_insert : Effectue une insertion groupée sans mettre à jour les données existantes.

hive_sync.enable

Indique s'il faut synchroniser les métadonnées avec Hive.

Boolean

Non

false

Valeurs valides :

  • true : Active la synchronisation.

  • false (par défaut) : Désactive la synchronisation.

hive_sync.mode

Le mode de synchronisation des métadonnées Hive.

String

Non

hms

  • hms (par défaut) : Définissez cette valeur si vous utilisez un catalogue DLF.

  • jdbc : Définissez cette valeur si vous utilisez un catalogue JDBC.

hive_sync.db

Le nom de la base de données Hive vers laquelle les données sont synchronisées.

String

Non

Le nom de la base de données de la table actuelle dans le catalogue.

Aucune

hive_sync.table

Le nom de la table Hive vers laquelle les données sont synchronisées.

String

Non

Le nom de la table actuelle.

Aucune

dlf.catalog.region

La région de Data Lake Formation (DLF).

String

Non

Aucune

.

Remarque
  • Le paramètre dlf.catalog.region prend effet uniquement lorsque le paramètre hive_sync.mode est défini surhms.

  • Assurez-vous que la région est identique à celle spécifiée pour le paramètre dlf.catalog.endpoint.

dlf.catalog.endpoint

Le point de terminaison de Data Lake Formation (DLF).

String

Non

Aucune

.

Remarque
  • Le paramètre dlf.catalog.endpoint prend effet uniquement lorsque le paramètre hive_sync.mode est défini sur hms.

  • Nous vous recommandons de définir le paramètre dlf.catalog.endpoint sur le point de terminaison VPC de DLF. Par exemple, si vous sélectionnez la région Chine (Hangzhou), définissez le paramètre dlf.catalog.endpoint sur dlf-vpc.cn-hangzhou.aliyuncs.com.

  • Si vous devez accéder à DLF entre différents VPC, consultez la section Gestion et opérations de l'espace de travail

Mappage des types de données

Type Iceberg

Type Flink

BOOLEAN

BOOLEAN

INT

INT

LONG

BIGINT

FLOAT

FLOAT

DOUBLE

DOUBLE

DECIMAL(P,S)

DECIMAL(P,S)

DATE

DATE

TIME

TIME

Remarque

Les horodatages Iceberg ont une précision à la microseconde, tandis que les horodatages Flink ont une précision à la milliseconde. Lorsque vous utilisez Flink pour lire des données depuis Iceberg, Flink convertit la précision temporelle en millisecondes.

TIMESTAMP

TIMESTAMP

TIMESTAMPTZ

TIMESTAMP_LTZ

STRING

STRING

FIXED(L)

BYTES

BINARY

VARBINARY

STRUCT<...>

ROW

LIST<E>

LIST

MAP<K,V>

MAP

Exemples

Assurez-vous de disposer d'un compartiment OSS et d'une base de données Data Lake Formation (DLF). Pour plus d'informations, consultez les sections Créer un compartiment et Bases de données, tables et fonctions.

Remarque

Lorsque vous spécifiez un path pour votre base de données Data Lake Formation (DLF), nous vous recommandons de suivre le format ${warehouse}/${database_name}.db. Par exemple, si l'adresse de l'entrepôt est oss://iceberg-test/warehouse et que le nom de la base de données est dlf_db, définissez le chemin OSS de dlf_db sur oss://iceberg-test/warehouse/dlf_db.db.

Exemple de table de destination (sink)

Cet exemple utilise le connecteur Datagen pour générer des données en continu aléatoires et les écrire dans une table Iceberg.

CREATE TEMPORARY TABLE datagen(
  id    BIGINT,
  data  STRING
) WITH (
  'connector' = 'datagen'
);

CREATE TEMPORARY TABLE dlf_iceberg (
  id    BIGINT,
  data  STRING
) WITH (
  'connector' = 'iceberg',
  'catalog-name' = '<yourCatalogName>',
  'catalog-database' = '<yourDatabaseName>',
  'io-impl' = 'org.apache.iceberg.aliyun.oss.OSSFileIO',
  'oss.endpoint' = '<yourOSSEndpoint>',  
  'access.key.id' = '${secret_values.ak_id}',
  'access.key.secret' = '${secret_values.ak_secret}',
  'catalog-impl' = 'org.apache.iceberg.aliyun.dlf.DlfCatalog',
  'warehouse' = '<yourOSSWarehousePath>',
  'dlf.catalog-id' = '<yourCatalogId>',
  'dlf.endpoint' = '<yourDLFEndpoint>',  
  'dlf.region-id' = '<yourDLFRegionId>'
);

INSERT INTO dlf_iceberg SELECT * FROM datagen;

Exemples de tables sources

  • Utilisez un catalogue DLF pour écrire les données d'une table source Iceberg dans une table de destination Iceberg.

    CREATE TEMPORARY TABLE src_iceberg (
      id    BIGINT,
      data  STRING
    ) WITH (
      'connector' = 'iceberg',
      'catalog-name' = '<yourCatalogName>',
      'catalog-database' = '<yourDatabaseName>',
      'io-impl' = 'org.apache.iceberg.aliyun.oss.OSSFileIO',
      'oss.endpoint' = '<yourOSSEndpoint>',  
      'access.key.id' = '${secret_values.ak_id}',
      'access.key.secret' = '${secret_values.ak_secret}',
      'catalog-impl' = 'org.apache.iceberg.aliyun.dlf.DlfCatalog',
      'warehouse' = '<yourOSSWarehousePath>',
      'dlf.catalog-id' = '<yourCatalogId>',
      'dlf.endpoint' = '<yourDLFEndpoint>',  
      'dlf.region-id' = '<yourDLFRegionId>'
    );
    
    CREATE TEMPORARY TABLE dst_iceberg (
      id    BIGINT,
      data  STRING
    ) WITH (
      'connector' = 'iceberg',
      'catalog-name' = '<yourCatalogName>',
      'catalog-database' = '<yourDatabaseName>',
      'io-impl' = 'org.apache.iceberg.aliyun.oss.OSSFileIO',
      'oss.endpoint' = '<yourOSSEndpoint>',  
      'access.key.id' = '${secret_values.ak_id}',
      'access.key.secret' = '${secret_values.ak_secret}',
      'catalog-impl' = 'org.apache.iceberg.aliyun.dlf.DlfCatalog',
      'warehouse' = '<yourOSSWarehousePath>',
      'dlf.catalog-id' = '<yourCatalogId>',
      'dlf.endpoint' = '<yourDLFEndpoint>',  
      'dlf.region-id' = '<yourDLFRegionId>'
    );
    
    BEGIN STATEMENT SET;
    
    INSERT INTO src_iceberg VALUES (1, 'AAA'), (2, 'BBB'), (3, 'CCC'), (4, 'DDD'), (5, 'EEE');
    INSERT INTO dst_iceberg SELECT * FROM src_iceberg;
    
    END;

Ingestion de données

Vous pouvez utiliser le connecteur Apache Iceberg comme destination (sink) dans un job YAML pour l'ingestion de données.

Syntaxe

sink:
  type: iceberg
  name: Iceberg Sink
  catalog.properties.rest.signing-region: cn-beijing
  catalog.properties.uri: http://cn-beijing-vpc.dlf.aliyuncs.com/iceberg
  catalog.properties.warehouse: flink_iceberg
  catalog.properties.type: rest
  catalog.properties.io-impl: org.apache.iceberg.rest.DlfFileIO

Paramètres

Paramètre

Description

Obligatoire

Type

Valeur par défaut

Remarques

type

Le type du connecteur.

Oui

STRING

Aucune

La valeur fixe est iceberg.

name

Le nom de la destination (sink).

Non

STRING

Aucune

Le nom de la destination (sink).

catalog.properties.rest.signing-region

L'ID de région de DLF. Pour plus d'informations, consultez la section Points de terminaison de service.

Oui

STRING

Aucune

Aucune

catalog.properties.uri

L'URI utilisée pour accéder au catalogue REST DLF. Pour plus d'informations, consultez la section Iceberg REST.

Oui

STRING

Aucune

Aucune

catalog.properties.warehouse

Le répertoire racine pour le stockage des fichiers.

Oui

STRING

Aucune

Aucune

catalog.properties.warehouse

Le répertoire racine pour le stockage des fichiers.

Non

STRING

Aucune

Aucune

catalog.properties.type

Le type du catalogue. La valeur doit être rest.

Oui

STRING

rest

Aucune

catalog.properties.io-impl

La valeur doit être org.apache.iceberg.rest.DlfFileIO.

Oui

STRING

org.apache.iceberg.rest.DlfFileIO

Aucune

partition.key

La clé de partition pour chaque table partitionnée.

Non

STRING

Aucune

Vous pouvez définir des clés de partition pour plusieurs tables. Séparez les définitions de table par un point-virgule (;) et les clés de partition par une virgule (,). Par exemple, vous pouvez spécifier testdb.table1:id1,id2;testdb.table2:name pour définir les clés de partition de la table testdb.table1 sur id1 et id2, et la clé de partition de la table testdb.table2 sur name.

Pour les partitions nécessitant des transformations implicites, ajoutez directement la fonction de transformation au champ de partition. Exemple : testdb.table1:truncate[10](id);testdb.table2:hour(create_time);testdb.table3:day(create_time);testdb.table4:month(create_time);testdb.table5:year(create_time);testdb.table6:bucket[10](create_time).

table.properties.*

Les paramètres pour la création d'une table Iceberg.

Non

String

Aucune

Pour plus d'informations, consultez la section Options de table Iceberg.

Réutiliser un catalogue existant

À partir de VVR 11.5, vous pouvez référencer directement un catalogue Iceberg intégré créé sur la page Gestion des données dans un job d'ingestion de données Flink CDC. Cela simplifie votre configuration en réduisant le nombre de propriétés de connexion requises.

sink:
  type: iceberg
  using.built-in-catalog: iceberg_catalog

Les jobs d'ingestion de données peuvent réutiliser automatiquement tous les paramètres du catalogue Iceberg. Cela équivaut à configurer manuellement les paramètres avec le préfixe catalog.properties. dans votre job YAML.

Si vous souhaitez remplacer les paramètres réutilisés, vous pouvez spécifier explicitement les paramètres YAML correspondants. Ces derniers sont prioritaires.

Exemple

L'exemple suivant montre comment utiliser un catalogue DLF comme catalogue Iceberg et écrire des données dans Data Lake Formation (DLF) :

  • Pour plus d'informations sur les paramètres préfixés par catalog.properties, consultez la section Créer un catalogue DLF Iceberg.

    source:
      type: mysql
      name: MySQL Source
      hostname: ${secret_values.mysql.hostname}
      port: ${mysql.port}
      username: ${secret_values.mysql.username}
      password: ${secret_values.mysql.password}
      tables: ${mysql.source.table}
      server-id: 8601-8604
    
    sink:
      type: iceberg
      name: Iceberg Sink
      catalog.properties.rest.signing-region: cn-beijing
      catalog.properties.uri: http://cn-beijing-vpc.dlf.aliyuncs.com/iceberg
      catalog.properties.warehouse: flink_iceberg
      catalog.properties.type: rest
      catalog.properties.io-impl: org.apache.iceberg.rest.DlfFileIO

Modifications de schéma

Lorsqu'il est utilisé comme destination (sink) d'ingestion de données, le connecteur Apache Iceberg prend en charge les modifications de schéma suivantes :

  • CREATE TABLE

  • ADD COLUMN

  • ALTER COLUMN TYPE (la modification du type d'une colonne de clé primaire n'est pas prise en charge)

  • RENAME COLUMN

  • DROP COLUMN

  • TRUNCATE TABLE

  • DROP TABLE

Remarque

Si la table Iceberg en aval existe déjà, le job utilise le schéma de table existant pour les écritures et ne recrée pas la table.

Documentation connexe

Pour plus d'informations sur les connecteurs pris en charge par Flink, consultez la section Connecteurs pris en charge.