Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Conector Apache Iceberg

Última atualização: Jun 27, 2026

Este tópico descreve como usar o conector Apache Iceberg.

Informações básicas

O Apache Iceberg é um formato de tabela de data lake aberto. Use o Apache Iceberg para criar rapidamente seu próprio serviço de armazenamento de data lake em HDFS ou OSS na nuvem e utilize mecanismos de computação do ecossistema open source de big data, como Flink, Spark, Hive e Presto, para analisar dados no data lake.

Categoria

Descrição

Tipos suportados

Tabela de origem, tabela de destino e destino de ingestão de dados

Modos de execução

Modo em lote e modo streaming

Formato de dados

Não aplicável

Métricas específicas

Nenhuma

Tipos de API

SQL, job YAML de ingestão de dados

Suporte a atualização ou exclusão de dados nas tabelas de destino

Sim

Recursos

O conector Apache Iceberg oferece os seguintes recursos:

  • Cria um serviço de armazenamento de data lake leve e de baixo custo, baseado em HDFS ou armazenamento de objetos.

  • Fornece semântica ACID completa.

  • Suporta consultas time travel para acessar versões históricas dos dados.

  • Filtra dados com eficiência.

  • Suporta evolução de schema.

  • Suporta evolução de partição.

Nota

Use as capacidades de tolerância a falhas e processamento de stream do Flink para importar grandes volumes de dados de log para um data lake Apache Iceberg em tempo real. Em seguida, use o Flink ou outros mecanismos de análise para extrair valor desses dados.

Limites

  • O conector Apache Iceberg tem suporte apenas no Realtime Compute for Apache Flink com Ververica Runtime (VVR) 4.0.8 ou posterior. Use esse conector com um catálogo Data Lake Formation (DLF). Para mais informações, consulte Gerenciar catálogos DLF-Legacy.

  • O conector Apache Iceberg suporta os formatos de tabela v1 e v2 do Apache Iceberg. Para mais detalhes, consulte a Especificação de Tabela Iceberg.

    Nota

    Apenas o Realtime Compute for Apache Flink com VVR 8.0.7 ou posterior suporta o formato de tabela v2.

  • No modo de leitura streaming, use apenas tabelas Iceberg do tipo append-only como tabelas de origem.

Sintaxe

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

Opções WITH

Opções comuns

Parâmetro

Descrição

Tipo

Obrigatório

Padrão

Observações

connector

Tipo do conector.

String

Sim

Nenhum

O valor deve sericeberg.

catalog-name

Nome do catálogo.

String

Sim

Nenhum

Insira um nome personalizado em inglês.

catalog-database

Nome do banco de dados.

String

Sim

default

Nome do banco de dados no Data Lake Formation (DLF), como dlf_db.

Nota

Se você não tiver um banco de dados do Data Lake Formation (DLF), crie um primeiro.

io-impl

Classe de implementação do sistema de arquivos distribuído.

String

Sim

Nenhum

O valor deve serorg.apache.iceberg.aliyun.oss.OSSFileIO.

oss.endpoint

Endpoint do Object Storage Service (OSS) da Alibaba Cloud.

String

Não

Nenhum

Para mais informações, consulte Regiões e endpoints.

Nota
  • Recomendamos definir o parâmetro oss.endpoint como o endpoint de VPC do OSS. Por exemplo, se selecionar a região China (Hangzhou), defina o parâmetro oss.endpoint como oss-cn-hangzhou-internal.aliyuncs.com.

  • Para acessar o OSS entre diferentes VPCs, consulte Como acesso outros serviços entre VPCs?

  • access.key.id: para VVR 8.0.6 e anteriores

  • access-key-id: para VVR 8.0.7 e posteriores

AccessKey ID da sua conta Alibaba Cloud.

String

Sim

Nenhum

Para mais informações, consulte Como visualizo o AccessKey ID e o AccessKey secret?

Importante

Para evitar vazamento das suas informações de AccessKey, recomendamos usar variáveis para especificar os valores de AccessKey. Para mais detalhes, consulte Variáveis de projeto.

  • access.key.secret: para VVR 8.0.6 e anteriores

  • access-key-secret: para VVR 8.0.7 e posteriores

AccessKey secret da sua conta Alibaba Cloud.

String

Sim

Nenhum

catalog-impl

Nome da classe do catálogo.

String

Sim

Nenhum

O valor deve serorg.apache.iceberg.aliyun.dlf.DlfCatalog.

warehouse

Caminho no OSS para armazenar os dados da tabela.

String

Sim

Nenhum

Nenhum

dlf.catalog-id

ID da sua conta Alibaba Cloud.

String

Sim

Nenhum

Obtenha o ID da conta na página User Information.

dlf.endpoint

Endpoint do Data Lake Formation (DLF).

String

Sim

Nenhum

.

Nota
  • Recomendamos definir o parâmetro dlf.endpoint como o endpoint de VPC do DLF. Por exemplo, se selecionar a região China (Hangzhou), defina o parâmetro dlf.endpoint como dlf-vpc.cn-hangzhou.aliyuncs.com.

  • Para acessar o DLF entre diferentes VPCs, consulte Gerenciamento e operações de workspace

dlf.region-id

Região do Data Lake Formation (DLF).

String

Sim

Nenhum

.

Nota

Certifique-se de que a região seja a mesma especificada no parâmetro dlf.endpoint.

Opções exclusivas de destino

Parâmetro

Descrição

Tipo

Obrigatório

Padrão

Observações

write.operation

Modo de operação de escrita.

String

Não

upsert

  • upsert (padrão): atualiza dados.

  • insert: adiciona dados.

  • bulk_insert: executa inserção em massa sem atualizar os dados existentes.

hive_sync.enable

Define se os metadados devem ser sincronizados com o Hive.

Boolean

Não

false

Valores válidos:

  • true: ativa a sincronização.

  • false (padrão): desativa a sincronização.

hive_sync.mode

Modo de sincronização de metadados do Hive.

String

Não

hms

  • hms (padrão): defina este valor se usar um catálogo DLF.

  • jdbc: defina este valor se usar um catálogo JDBC.

hive_sync.db

Nome do banco de dados Hive para sincronização de dados.

String

Não

Nome do banco de dados da tabela atual no catálogo.

Nenhum

hive_sync.table

Nome da tabela Hive para sincronização de dados.

String

Não

Nome da tabela atual.

Nenhum

dlf.catalog.region

Região do Data Lake Formation (DLF).

String

Não

Nenhum

.

Nota
  • O parâmetro dlf.catalog.region só tem efeito quando o parâmetro hive_sync.mode está definido comohms.

  • Certifique-se de que a região seja a mesma especificada no parâmetro dlf.catalog.endpoint.

dlf.catalog.endpoint

Endpoint do Data Lake Formation (DLF).

String

Não

Nenhum

.

Nota
  • O parâmetro dlf.catalog.endpoint só tem efeito quando o parâmetro hive_sync.mode está definido como hms.

  • Recomendamos definir o parâmetro dlf.catalog.endpoint como o endpoint de VPC do DLF. Por exemplo, se selecionar a região China (Hangzhou), defina o parâmetro dlf.catalog.endpoint como dlf-vpc.cn-hangzhou.aliyuncs.com.

  • Para acessar o DLF entre diferentes VPCs, consulte Gerenciamento e operações de workspace

Mapeamento de tipos de dados

Tipo Iceberg

Tipo Flink

BOOLEAN

BOOLEAN

INT

INT

LONG

BIGINT

FLOAT

FLOAT

DOUBLE

DOUBLE

DECIMAL(P,S)

DECIMAL(P,S)

DATE

DATE

TIME

TIME

Nota

Os timestamps do Iceberg têm precisão de microssegundos, enquanto os do Flink têm precisão de milissegundos. Ao ler dados do Iceberg com o Flink, a precisão temporal é convertida para milissegundos.

TIMESTAMP

TIMESTAMP

TIMESTAMPTZ

TIMESTAMP_LTZ

STRING

STRING

FIXED(L)

BYTES

BINARY

VARBINARY

STRUCT<...>

ROW

LIST<E>

LIST

MAP<K,V>

MAP

Exemplos

Certifique-se de ter um bucket no OSS e um banco de dados no Data Lake Formation (DLF). Para mais informações, consulte Criar um bucket e Bancos de dados, tabelas e funções.

Nota

Ao especificar um path para seu banco de dados no Data Lake Formation (DLF), recomendamos seguir o formato ${warehouse}/${database_name}.db. Por exemplo, se o endereço do warehouse for oss://iceberg-test/warehouse e o nome do banco de dados for dlf_db, defina o caminho no OSS de dlf_db como oss://iceberg-test/warehouse/dlf_db.db.

Exemplo de tabela de destino

Este exemplo usa o conector Datagen para gerar dados streaming aleatórios e gravá-los em uma tabela 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;

Exemplos de tabela de origem

  • Use um catálogo DLF para gravar dados de uma tabela de origem Iceberg em uma tabela de destino 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;

Ingestão de dados

Use o conector Apache Iceberg como destino em um job YAML para ingestão de dados.

Sintaxe

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

Parâmetros

Parâmetro

Descrição

Obrigatório

Tipo

Padrão

Observações

type

Tipo do conector.

Sim

STRING

Nenhum

O valor fixo é iceberg.

name

Nome do destino.

Não

STRING

Nenhum

Nome do destino.

catalog.properties.rest.signing-region

ID da região do DLF. Para mais informações, consulte Endpoints de serviço.

Sim

STRING

Nenhum

Nenhum

catalog.properties.uri

URI usada para acessar o catálogo REST do DLF. Para mais informações, consulte Iceberg REST.

Sim

STRING

Nenhum

Nenhum

catalog.properties.warehouse

Diretório raiz para armazenamento de arquivos.

Sim

STRING

Nenhum

Nenhum

catalog.properties.warehouse

Diretório raiz para armazenamento de arquivos.

Não

STRING

Nenhum

Nenhum

catalog.properties.type

Tipo do catálogo. O valor deve ser rest.

Sim

STRING

rest

Nenhum

catalog.properties.io-impl

O valor deve ser org.apache.iceberg.rest.DlfFileIO.

Sim

STRING

org.apache.iceberg.rest.DlfFileIO

Nenhum

partition.key

Chave de partição para cada tabela particionada.

Não

STRING

Nenhum

Defina chaves de partição para várias tabelas. Separe as definições de tabela com ponto e vírgula (;) e as chaves de partição com vírgula (,). Por exemplo, especifique testdb.table1:id1,id2;testdb.table2:name para definir as chaves de partição da tabela testdb.table1 como id1 e id2, e a chave de partição da tabela testdb.table2 como name.

Para partições que exigem transformações implícitas, adicione a função de transformação diretamente ao campo de partição. Exemplo: 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.*

Parâmetros para criar uma tabela Iceberg.

Não

String

Nenhum

Para mais informações, consulte Opções de tabela Iceberg.

Reutilizar um catálogo existente

A partir do VVR 11.5, referencie diretamente um catálogo Iceberg integrado criado na página Data Management em um job de ingestão de dados Flink CDC. Isso simplifica a configuração ao reduzir o número de propriedades de conexão necessárias.

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

Jobs de ingestão de dados reutilizam automaticamente todos os parâmetros do catálogo Iceberg. Isso equivale a configurar manualmente os parâmetros com o prefixo catalog.properties. no job YAML.

Para substituir os parâmetros reutilizados, especifique explicitamente os parâmetros YAML correspondentes. Esses parâmetros têm precedência.

Exemplo

O exemplo a seguir mostra como usar um catálogo DLF como catálogo Iceberg e gravar dados no Data Lake Formation (DLF):

  • Para informações sobre os parâmetros com prefixo catalog.properties, consulte Criar um catálogo Iceberg DLF.

    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

Alterações de schema

Quando usado como destino de ingestão de dados, o conector Apache Iceberg suporta as seguintes alterações de schema:

  • CREATE TABLE

  • ADD COLUMN

  • ALTER COLUMN TYPE (não há suporte para modificação do tipo de coluna de chave primária)

  • RENAME COLUMN

  • DROP COLUMN

  • TRUNCATE TABLE

  • DROP TABLE

Nota

Se a tabela Iceberg downstream já existir, o job usará o schema da tabela existente para gravações e não criará a tabela novamente.

Documentação relacionada

Para mais informações sobre conectores suportados pelo Flink, consulte Conectores suportados.