Todos os produtos
Search
Central de documentação

Data Lake Formation:Acesse o DLF com o Flink CDC

Última atualização: Sep 18, 2026

O Flink Change Data Capture (CDC) é uma ferramenta de ingestão de dados oferecida pelo Realtime Compute for Apache Flink. Ele permite sincronizar bancos de dados inteiros de fontes de origem para seu data lakehouse. Este tópico orienta você no uso do Flink CDC para ingerir dados em um catálogo DLF em tempo real via Paimon REST.

Ao concluir este tópico, você terá:

  • Criado um catálogo DLF como destino de ingestão

  • Configurado um job Flink CDC YAML com o sink Paimon apontando para o DLF

  • Executado exemplos para cenários comuns de ingestão: sincronização de banco de dados completo, tabelas particionadas, tabelas append-only e sources Kafka

Pré-requisitos

Antes de começar, verifique se você tem:

  • Um workspace do Realtime Compute for Apache Flink. Consulte Criar um workspace.

  • O workspace do Realtime Compute for Apache Flink e os catálogos DLF na mesma região.

  • O VPC do seu workspace do Realtime Compute for Apache Flink adicionado à lista de permissões de VPC do DLF. Consulte Configure uma lista de permissões de VPC.

Requisitos de engine

O job do Realtime Compute for Apache Flink deve usar o Ververica Runtime (VVR) versão 11.1.0 ou posterior.

Criar um catálogo DLF

Consulte Introdução ao DLF.

Criar e configurar um job de ingestão de dados

  1. Crie um rascunho YAML para ingerir dados com o Flink CDC. Para mais informações, consulte Develop Flink CDC jobs for data ingestion (Beta).

  2. Configure o módulo sink:

    sink:
      type: paimon
      catalog.properties.metastore: rest
      catalog.properties.uri: dlf_uri
      catalog.properties.warehouse: your_warehouse
      catalog.properties.token.provider: dlf
      # (Optional) The commit user. Set different commit users for different jobs to avoid conflicts.
      commit.user: your_job_name
      # (Optional) Enable deletion vectors to improve read performance.
      table.properties.deletion-vectors.enabled: true

    Substitua os valores de espaço reservado pelos seus valores reais:

    Opção de configuração

    Descrição

    Obrigatório

    Padrão

    Exemplo de valor

    catalog.properties.metastore

    O tipo de metastore. Defina como rest.

    Sim

    rest

    catalog.properties.token.provider

    O provedor de token. Defina como dlf.

    Sim

    dlf

    catalog.properties.uri

    O URI para acessar o servidor de catálogo DLF REST. Formato: http://[region-id]-vpc.dlf.aliyuncs.com. Para IDs de região, consulte Regions and endpoints.

    Sim

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

    catalog.properties.warehouse

    O nome do catálogo Paimon.

    Sim

    dlf_test

    commit.user

    O usuário de commit para gravações de dados. Atribua usuários de commit exclusivos a jobs distintos para evitar conflitos.

    Não

    admin

    your_job_name

    table.properties.deletion-vectors.enabled

    Ativa deletion vectors para melhorar o desempenho de leitura com impacto mínimo nas gravações.

    Não

    true

Observações de uso

Antes de executar um job, atente-se às seguintes restrições:

  • Conflitos de commit user: O usuário de commit padrão é admin. Executar jobs de gravação de dados simultâneos na mesma tabela com o mesmo commit user causa conflitos de commit e inconsistências nos dados. Atribua um commit.user exclusivo a cada job.

  • Sem opções de compactação ou bucket: O DLF fornece compactação automática de arquivos. Não configure opções de compactação ou bucket, como bucket e num-sorted-run.compaction-trigger.

  • Deletion vectors: Defina table.properties.deletion-vectors.enabled: true para acelerar significativamente as leituras com impacto mínimo nas gravações.

Exemplos

Ingerir dados de um banco de dados MySQL completo para o DLF

O job a seguir sincroniza todas as tabelas de um banco de dados MySQL para 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 data from tables created in the incremental phase.
  scan.binlog.newly-added-table.enabled: true
  # (Optional) Sync table and field comments.
  include-comments.enabled: true
  # (Optional) Prioritize unbounded shards to prevent potential TaskManager out-of-memory (OOM) errors.
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (Optional) Deserialize data only from matched tables to speed up reads.
  scan.only.deserialize.captured.tables.changelog.enabled: true

sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Optional) The commit user. Set different commit users for different jobs to avoid conflicts.
  commit.user: your_job_name
  # (Optional) Enable deletion vectors to improve read performance.
  table.properties.deletion-vectors.enabled: true

As opções recomendadas para o MySQL source são:

Opção

Descrição

scan.binlog.newly-added-table.enabled

Sincroniza dados de tabelas criadas durante a fase incremental.

include-comments.enabled

Sincroniza comentários de tabelas e campos.

scan.incremental.snapshot.unbounded-chunk-first.enabled

Evita possíveis erros de OOM no TaskManager priorizando shards ilimitados.

scan.only.deserialize.captured.tables.changelog.enabled

Desserializa dados somente das tabelas correspondentes para acelerar as leituras.

Ingerir dados em uma tabela particionada

Para ingerir dados de uma tabela de origem não particionada em uma tabela particionada no DLF, adicione a opção partition-keys ao módulo transform. Para mais informações, consulte Data ingestion with 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 data from tables created in the incremental phase.
  scan.binlog.newly-added-table.enabled: true
  # (Optional) Sync table and field comments.
  include-comments.enabled: true
  # (Optional) Prioritize unbounded shards to prevent potential TaskManager out-of-memory (OOM) errors.
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (Optional) Deserialize data only from matched tables to speed up reads.
  scan.only.deserialize.captured.tables.changelog.enabled: true

sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Optional) The commit user. Set different commit users for different jobs to avoid conflicts.
  commit.user: your_job_name
  # (Optional) Enable deletion vectors to improve read performance.
  table.properties.deletion-vectors.enabled: true

transform:
  - source-table: mysql_test.tbl1
    partition-keys: id,pt
  - source-table: mysql_test.tbl2
    partition-keys: id,pt

Ingerir dados em uma tabela append-only

Para implementar exclusões lógicas (soft delete) durante a ingestão, use converter-after-transform: SOFT_DELETE no módulo transform. Isso converte operações de exclusão em operações de inserção, fazendo com que a tabela downstream registre todas as operações de alteração de forma completa. O campo __data_event_type__ em projection grava o tipo de alteração como uma nova coluna na tabela downstream.

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 data from tables created in the incremental phase.
  scan.binlog.newly-added-table.enabled: true
  # (Optional) Sync table and field comments.
  include-comments.enabled: true
  # (Optional) Prioritize unbounded shards to prevent potential TaskManager out-of-memory (OOM) errors.
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (Optional) Deserialize data only from matched tables to speed up reads.
  scan.only.deserialize.captured.tables.changelog.enabled: true

sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Optional) The commit user. Set different commit users for different jobs to avoid conflicts.
  commit.user: your_job_name
  # (Optional) Enable deletion vectors to improve read performance.
  table.properties.deletion-vectors.enabled: true

transform:
  - source-table: mysql_test.tbl1
    partition-keys: id,pt
    projection: \*, __data_event_type__ AS op_type
    converter-after-transform: SOFT_DELETE
  - source-table: mysql_test.tbl2
    partition-keys: id,pt
    projection: \*, __data_event_type__ AS op_type
    converter-after-transform: SOFT_DELETE

Para mais informações, consulte Data ingestion with Flink CDC.

Sincronizar dados do Kafka para o DLF em tempo real

O job a seguir lê dados CDC do tópico Kafka inventory (tabelas customers e products no formato Debezium JSON) e os sincroniza com as tabelas de destino correspondentes no DLF. Como as mensagens no formato Debezium JSON não contêm informações de chave primária, a chave primária é especificada explicitamente no módulo transform.

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: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Optional) The commit user. Set different commit users for different jobs to avoid conflicts.
  commit.user: your_job_name
  # (Optional) Enable deletion vectors to improve read performance.
  table.properties.deletion-vectors.enabled: true

# Debezium JSON does not contain primary key info, so explicitly specify primary keys.
transform:
  - source-table: \.*.\.*
    projection: \*
    primary-keys: id

Observações adicionais para sources Kafka:

  • O source Kafka suporta três formatos de dados: canal-json, debezium-json (padrão) e json.

  • Ao ingerir dados de múltiplas partições Kafka em uma única tabela no DLF, defina debezium-json.distributed-tables ou canal-json.distributed-tables como true.

  • O source Kafka suporta múltiplas políticas de inferência de esquema por meio da opção schema.inference.strategy. Para mais informações, consulte Message Queue for Apache Kafka.

Para mais informações, consulte Data ingestion with Flink CDC.