Sincronize dados em um catálogo do DLF por meio do Iceberg REST usando o Flink CDC no Realtime Compute for Apache Flink.
Pré-requisitos
Você possui um workspace totalmente gerenciado para o Realtime Compute for Apache Flink. Caso ainda não tenha criado um, consulte Ativar o Realtime Compute for Apache Flink.
Verifique se o workspace do Realtime Compute for Apache Flink e o DLF estão na mesma região. Além disso, adicione a VPC do workspace à lista de permissões do DLF. Para mais informações, consulte Configurar uma lista de permissões de VPC.
Limitações
A conectividade do Iceberg REST com o DLF exige o mecanismo VVR 11.6.0 ou posterior do Realtime Compute for Apache Flink.
Registrar o catálogo do DLF no Flink
O registro de um catálogo no Flink cria um mapeamento para o seu catálogo do DLF. Criar ou excluir o catálogo no Flink não afeta os dados reais no DLF. Todas as tabelas criadas no catálogo do DLF por meio do Iceberg REST são tabelas Iceberg.
Faça login no Console de Gerenciamento do Realtime Compute for Apache Flink.
Na coluna Actions do seu workspace, clique em Console.
No painel de navegação à esquerda, clique em Development > Scripts.
-
Crie um novo script e cole a seguinte instrução SQL no editor SQL.
CREATE CATALOG `catalog_name` WITH ( 'type' = 'iceberg', 'catalog-type' = 'rest', 'uri' = 'http://{region-id}-vpc.dlf.aliyuncs.com/iceberg', 'warehouse' = 'iceberg_test', 'rest.signing-region' = '{region-id}', 'io-impl' = 'org.apache.iceberg.rest.DlfFileIO' );Substitua
{region-id}pelo ID da região do seu catálogo do DLF, por exemplocn-hangzhououap-southeast-1. Para todas as regiões suportadas e seus valores de endpoint, consulte Regiões e endpoints. No canto inferior direito, clique em Environment, selecione um cluster de sessão executando VVR 11.2.0 ou posterior e execute a instrução SQL.
A tabela a seguir descreve as opções de configuração.
|
Opção |
Descrição |
Obrigatório |
Exemplo |
|
|
Tipo do catálogo. Defina como |
Sim |
|
|
|
Tipo de conexão do catálogo. Defina como |
Sim |
|
|
|
Provedor de token para autenticação no DLF. Defina como |
Sim |
|
|
|
Endpoint do Iceberg REST para o seu catálogo do DLF. Use o formato |
Sim |
|
|
|
Nome do seu catálogo do DLF. |
Sim |
|
|
|
ID da região do seu catálogo do DLF. Para IDs de região, consulte Regiões e endpoints. |
Sim |
|
|
|
Implementação FileIO para o DLF. Defina como |
Sim |
|
Configure o Flink CDC para usar o catálogo
Crie um job de ingestão de dados: Desenvolver um job de ingestão de dados do Flink CDC.
Se você já possui um mapeamento de Catálogo do Flink, Reutilize um Catálogo existente para obter informações de conexão e adicione esta configuração de Sink:
sink:
type: iceberg
using.built-in-catalog: catalog_name
Exemplos de configuração
Padrões comuns de sincronização de dados usando jobs YAML do Flink CDC:
Sincronizar um banco de dados MySQL completo para o DLF
Este job sincroniza um banco de dados MySQL inteiro com o DLF:
source:
type: mysql
name: MySQL Source
hostname: ${mysql.hostname}
port: ${mysql.port}
username: ${mysql.username}
password: ${mysql.password}
tables: mysql_test.\.*
server-id: 8601-8604
# (Optional) Sync tables created during incremental phase.
scan.binlog.newly-added-table.enabled: true
# (Optional) Sync table and field comments.
include-comments.enabled: true
# (Optional) Prioritize unbounded chunks to prevent TaskManager OOM errors.
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
# (Optional) Speed up reads by parsing only matched tables.
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: iceberg
using.built-in-catalog: catalog_name
Parâmetros opcionais recomendados para sources MySQL (MySQL):
-
Parâmetro: scan.binlog.newly-added-table.enabled
Função: Sincroniza tabelas criadas durante a fase incremental.
-
Parâmetro: include-comments.enabled
Função: Sincroniza comentários de tabelas e campos.
-
Parâmetro: scan.incremental.snapshot.unbounded-chunk-first.enabled
Função: Previne erros de OOM no TaskManager.
-
Parâmetro: scan.only.deserialize.captured.tables.changelog.enabled: true
Função: Acelera a leitura processando apenas as tabelas correspondentes.
Gravar em uma tabela particionada do DLF
Especifique chaves de partição com o parâmetro partition-keys (Referência de desenvolvimento de job de ingestão de dados do Flink CDC):
source:
type: mysql
name: MySQL Source
hostname: ${mysql.hostname}
port: ${mysql.port}
username: ${mysql.username}
password: ${mysql.password}
tables: mysql_test.\.*
server-id: 8601-8604
# (Optional) Sync tables created during incremental phase.
scan.binlog.newly-added-table.enabled: true
# (Optional) Sync table and field comments.
include-comments.enabled: true
# (Optional) Prioritize unbounded chunks to prevent TaskManager OOM errors.
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
# (Optional) Speed up reads by parsing only matched tables.
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: iceberg
using.built-in-catalog: catalog_name
transform:
- source-table: mysql_test.tbl1
# (Optional) Set partition keys.
partition-keys: id,pt
- source-table: mysql_test.tbl2
partition-keys: id,pt
Gravar em uma tabela append-only do DLF
Para implementar exclusões lógicas (converter operações DELETE em INSERTs no destino), configure o job da seguinte forma (Referência de desenvolvimento de job de ingestão de dados do Flink CDC):
source:
type: mysql
name: MySQL Source
hostname: ${mysql.hostname}
port: ${mysql.port}
username: ${mysql.username}
password: ${mysql.password}
tables: mysql_test.\.*
server-id: 8601-8604
# (Optional) Sync tables created during incremental phase.
scan.binlog.newly-added-table.enabled: true
# (Optional) Sync table and field comments.
include-comments.enabled: true
# (Optional) Prioritize unbounded chunks to prevent TaskManager OOM errors.
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
# (Optional) Speed up reads by parsing only matched tables.
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: iceberg
using.built-in-catalog: catalog_name
transform:
- source-table: mysql_test.tbl1
# (Optional) Set partition keys.
partition-keys: id,pt
# (Optional) Implement soft delete.
projection: \*, __data_event_type__ AS op_type
converter-after-transform: SOFT_DELETE
- source-table: mysql_test.tbl2
# (Optional) Set partition keys.
partition-keys: id,pt
# (Optional) Implement soft delete.
projection: \*, __data_event_type__ AS op_type
converter-after-transform: SOFT_DELETE
Adicionar
__data_event_type__à projeção grava o tipo de evento de alteração como um novo campo no destino. Definirconverter-after-transformcomoSOFT_DELETEconverte DELETEs em INSERTs, registrando todos os eventos de alteração. Consulte a Referência de desenvolvimento de job de ingestão de dados do Flink CDC.
Sincronizar dados CDC do Kafka em tempo real para o DLF
Este exemplo sincroniza dados de alteração de duas tabelas (customers e products) de um tópico de inventário do Kafka no formato JSON Debezium para o DLF:
source:
type: kafka
name: Kafka Source
properties.bootstrap.servers: ${kafka.bootstrap.servers}
topic: inventory
scan.startup.mode: earliest-offset
value.format: debezium-json
debezium-json.distributed-tables: true
sink:
type: iceberg
using.built-in-catalog: catalog_name
# Debezium JSON lacks primary key info. Add it manually.
transform:
- source-table: \.*.\.*
projection: \*
primary-keys: id
A source do Kafka suporta os formatos
canal-json,debezium-json(padrão) ejson.-
Ao usar
debezium-json, adicione manualmente uma chave primária usando uma regra de transformação, pois as mensagens JSON do Debezium não contêm informações de chave primária:transform: - source-table: \.*.\.* projection: \* primary-keys: id Se os dados de uma única tabela abrangerem várias partições, ou se for necessário mesclar tabelas entre partições, defina
debezium-json.distributed-tablesoucanal-json.distributed-tablescomotrue.A source do Kafka oferece suporte a múltiplas estratégias de inferência de schema por meio do parâmetro
schema.inference.strategy. Consulte Kafka.
Sincronizar logs do Kafka em tempo real para o DLF
Para dados JSON personalizados no Kafka, o Flink CDC gerencia automaticamente a inferência de tipos de dados, a inferência de schema e a evolução de schema.
Este exemplo sincroniza uma única tabela de logs JSON do tópico de inventário para o DLF:
source:
type: kafka
name: Kafka Source
properties.bootstrap.servers: ${kafka.bootstrap.servers}
topic: inventory
scan.startup.mode: earliest-offset
value.format: json
# (Optional) Recursively flatten nested columns in JSON data.
json.infer-schema.flatten-nested-columns.enable: true
# (Optional) Skip first 100 parsing errors. Job fails if errors exceed 100.
ingestion.ignore-errors: true
ingestion.error-tolerance.max-count: 100
sink:
type: iceberg
using.built-in-catalog: catalog_name
# Add primary key to table.
transform:
- source-table: \.*.\.*
projection: \*
primary-keys: id
# Write all inventory topic data to test_database.inventory.
route:
- source-table: inventory
sink-table: test_database.inventory
pipeline:
# (Optional) Log dirty data that causes processing exceptions.
dirty-data.collector:
name: Logger Dirty Data Collector
type: logger
Este exemplo sincroniza várias tabelas de logs JSON do tópico de inventário, usando os campos databaseName e tableName para identificar as tabelas:
source:
type: kafka
name: Kafka Source
properties.bootstrap.servers: ${kafka.bootstrap.servers}
topic: inventory
scan.startup.mode: earliest-offset
value.format: json
# (Optional) Recursively flatten nested columns in JSON data.
json.infer-schema.flatten-nested-columns.enable: true
# Use databaseName and tableName field values as database and table names.
json.decode.parser-table-id.fields: databaseName,tableName
# (Optional) Skip first 100 parsing errors. Job fails if errors exceed 100.
ingestion.ignore-errors: true
ingestion.error-tolerance.max-count: 100
sink:
type: iceberg
using.built-in-catalog: catalog_name
# Add primary key to tables.
transform:
- source-table: \.*.\.*
projection: \*
primary-keys: id
# Write ods.inventory, ods.customer, and ods.user to test_database.inventory, test_database.customer, and test_database.user respectively.
route:
- source-table: ods.inventory
sink-table: test_database.inventory
- source-table: ods.customer
sink-table: test_database.customer
- source-table: ods.user
sink-table: test_database.user
pipeline:
# (Optional) Log dirty data that causes processing exceptions.
dirty-data.collector:
name: Logger Dirty Data Collector
type: logger
Estratégias de inferência e evolução de schema JSON: Estratégias de análise de schema e sincronização de alterações.
Referência completa de configuração: Referência de desenvolvimento de job de ingestão de dados do Flink CDC.