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.
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.
NotaApenas 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 ser |
|
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 ser |
|
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
|
|
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. |
|
AccessKey secret da sua conta Alibaba Cloud. |
String |
Sim |
Nenhum |
|
|
catalog-impl |
Nome da classe do catálogo. |
String |
Sim |
Nenhum |
O valor deve ser |
|
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
|
|
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 |
|
|
hive_sync.enable |
Define se os metadados devem ser sincronizados com o Hive. |
Boolean |
Não |
false |
Valores válidos:
|
|
hive_sync.mode |
Modo de sincronização de metadados do Hive. |
String |
Não |
hms |
|
|
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
|
|
dlf.catalog.endpoint |
Endpoint do Data Lake Formation (DLF). |
String |
Não |
Nenhum |
. Nota
|
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.
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 é |
|
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 ( Para partições que exigem transformações implícitas, adicione a função de transformação diretamente ao campo de partição. Exemplo: |
|
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
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.