Todos os produtos
Search
Central de documentação

Data Lake Formation:Connect Flink CDC to DLF

Última atualização: Jul 11, 2026

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.
  1. Faça login no Console de Gerenciamento do Realtime Compute for Apache Flink.

  2. Na coluna Actions do seu workspace, clique em Console.

  3. No painel de navegação à esquerda, clique em Development > Scripts.

  4. 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 exemplo cn-hangzhou ou ap-southeast-1. Para todas as regiões suportadas e seus valores de endpoint, consulte Regiões e endpoints.

  5. 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

type

Tipo do catálogo. Defina como iceberg.

Sim

iceberg

catalog-type

Tipo de conexão do catálogo. Defina como rest.

Sim

rest

token.provider

Provedor de token para autenticação no DLF. Defina como dlf.

Sim

dlf

uri

Endpoint do Iceberg REST para o seu catálogo do DLF. Use o formato http://{region-id}-vpc.dlf.aliyuncs.com/iceberg. Para valores específicos por região, consulte Regiões e endpoints.

Sim

http://ap-southeast-1-vpc.dlf.aliyuncs.com/iceberg

warehouse

Nome do seu catálogo do DLF.

Sim

iceberg_test

rest.signing-region

ID da região do seu catálogo do DLF. Para IDs de região, consulte Regiões e endpoints.

Sim

ap-southeast-1

io-impl

Implementação FileIO para o DLF. Defina como org.apache.iceberg.rest.DlfFileIO.

Sim

org.apache.iceberg.rest.DlfFileIO

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
Nota

Parâmetros opcionais recomendados para sources MySQL (MySQL):

  1. Parâmetro: scan.binlog.newly-added-table.enabled

    Função: Sincroniza tabelas criadas durante a fase incremental.

  2. Parâmetro: include-comments.enabled

    Função: Sincroniza comentários de tabelas e campos.

  3. Parâmetro: scan.incremental.snapshot.unbounded-chunk-first.enabled

    Função: Previne erros de OOM no TaskManager.

  4. 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
Nota

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
Nota
  • A source do Kafka suporta os formatos canal-json, debezium-json (padrão) e json.

  • 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-tables ou canal-json.distributed-tables como true.

  • 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.