Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Ingest data into a DLF data lake in real time using Flink CDC

Última atualização: Jun 27, 2026

Utilize o Flink CDC para gravar dados de alteração em tempo real do MySQL ou Kafka no Data Lake Formation (DLF).

O DLF é uma plataforma totalmente gerenciada que oferece metadados unificados, armazenamento de dados e gerenciamento de dados, incluindo gestão de metadados, controle de permissões e otimização de armazenamento. Os jobs de ingestão de dados gravam no DLF por meio do Paimon Catalog, o que permite ingerir bancos de dados inteiros de grande escala no data lake em tempo real. Para obter mais informações sobre o DLF, consulte O que é o Data Lake Formation?.

Pré-requisitos

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

  • Uma instância do DLF com o Paimon Catalog ativado

  • A URI do seu DLF Rest Catalog Server (formato: http://<region-id>-vpc.dlf.aliyuncs.com, por exemplo, http://cn-hangzhou-vpc.dlf.aliyuncs.com)

  • O nome do warehouse do seu DLF Catalog

Configuração do sink do DLF

Todos os exemplos neste tópico utilizam os seguintes parâmetros de sink do Paimon para se conectar ao DLF.

Parâmetro

Obrigatório

Descrição

catalog.properties.metastore

Sim

Tipo de Metastore. Defina como rest.

catalog.properties.token.provider

Sim

Provedor de token. Defina como dlf.

catalog.properties.uri

Sim

URI do DLF Rest Catalog Server. Formato: http://<region-id>-vpc.dlf.aliyuncs.com.

catalog.properties.warehouse

Sim

Nome do DLF Catalog.

table.properties.deletion-vectors.enabled

Não

Defina como true para ativar vetores de exclusão. Melhora significativamente o desempenho de leitura com impacto mínimo na performance de gravação e atualização, permitindo atualizações quase em tempo real e consultas de alta velocidade.

Observações de uso

  • Não adicione parâmetros de mesclagem de arquivos ou relacionados a buckets (como bucket e num-sorted-run.compaction-trigger) à configuração do sink. O DLF gerencia a mesclagem de arquivos automaticamente, e a inclusão desses parâmetros causa conflitos.

  • Para parâmetros de source do MySQL, consulte MySQL.

Sincronize um banco de dados MySQL inteiro com o DLF

O job YAML de CDC a seguir sincroniza um banco de dados MySQL completo com o DLF. A seção source captura todas as tabelas correspondentes a mysql_test.*, enquanto a seção sink utiliza o Paimon Catalog do DLF como destino.

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) Synchronize data from tables newly created in the incremental phase.
  scan.binlog.newly-added-table.enabled: true
  # (Optional) Synchronize table and field comments.
  include-comments.enabled: true
  # (Optional) Prioritize the distribution of unbounded chunks to prevent potential TaskManager OutOfMemory issues.
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (Optional) Enable parsing filters to accelerate reads.
  scan.only.deserialize.captured.tables.changelog.enabled: true

sink:
  type: paimon
  # The Metastore type. Set to rest.
  catalog.properties.metastore: rest
  # The token provider. Set to dlf.
  catalog.properties.token.provider: dlf
  # The URI to access the DLF Rest Catalog Server.
  catalog.properties.uri: dlf_uri
  # The name of the DLF Catalog.
  catalog.properties.warehouse: your_warehouse
  # (Optional) Enable deletion vectors to improve read performance.
  table.properties.deletion-vectors.enabled: true

Grave em tabelas particionadas no DLF

As tabelas de origem nos jobs de ingestão de dados não incluem informações sobre campos de partição. Para gravar em uma tabela particionada no DLF, adicione um bloco transform e especifique partition-keys para cada tabela de origem. Para detalhes sobre a sintaxe de partition-keys, consulte Referência de desenvolvimento de jobs 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) Synchronize data from tables newly created in the incremental phase.
  scan.binlog.newly-added-table.enabled: true
  # (Optional) Synchronize table and field comments.
  include-comments.enabled: true
  # (Optional) Prioritize the distribution of unbounded chunks to prevent potential TaskManager OutOfMemory issues.
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (Optional) Enable parsing filters to accelerate 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) Enable deletion vectors to improve read performance.
  table.properties.deletion-vectors.enabled: true

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

Grave em tabelas append-only no DLF

As tabelas de origem contêm dados completos de alterações do CDC, incluindo operações de exclusão. Para converter operações de exclusão em inserções (exclusão lógica) e gravar todas as alterações em uma tabela append-only no DLF, adicione __data_event_type__ à projection e defina converter-after-transform como SOFT_DELETE no bloco transform. Isso garante que cada evento de alteração — inclusive exclusões — seja registrado na tabela de destino como uma linha de inserção.

Para mais informações, consulte Referência de desenvolvimento de jobs 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) Synchronize data from tables newly created in the incremental phase.
  scan.binlog.newly-added-table.enabled: true
  # (Optional) Synchronize table and field comments.
  include-comments.enabled: true
  # (Optional) Prioritize the distribution of unbounded chunks to prevent potential TaskManager OutOfMemory issues.
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (Optional) Enable parsing filters to accelerate 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) Enable deletion vectors to improve read performance.
  table.properties.deletion-vectors.enabled: true

transform:
  - source-table: mysql_test.tbl1
    # Set partition fields.
    partition-keys: id,pt
    # Add the change type as a new field, and convert deletes to inserts.
    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

Sincronize dados CDC do Kafka com o DLF em tempo real

Caso os dados do MySQL já estejam fluindo para o Kafka — por exemplo, através da distribuição em tempo real — configure um job YAML de CDC para ler do Kafka e gravar no DLF.

O exemplo a seguir lê do tópico Kafka inventory, que armazena dados de alteração das tabelas customers e products no formato Debezium JSON, e sincroniza cada tabela com sua respectiva tabela de destino no DLF. Como as mensagens debezium-json não carregam informações de chave primária, o bloco transform adiciona manualmente id como chave primária.

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) Enable deletion vectors to improve read performance.
  table.properties.deletion-vectors.enabled: true

# debezium-json does not include primary key information. Add primary keys using a transform rule.
transform:
  - source-table: \.*.\.*
    projection: \*
    primary-keys: id

Observações sobre a source do Kafka:

  • Formatos suportados: canal-json, debezium-json (padrão) e json.

  • Evolução de schema: Se o schema de uma mensagem do Kafka for alterado (por exemplo, com a adição de um novo campo), a mudança será sincronizado automaticamente com o schema da tabela Paimon.

  • Chaves primárias com debezium-json: Mensagens debezium-json não registram informações de chave primária. Adicione chaves primárias manualmente usando uma regra de transformação:

    transform:
      - source-table: \.*.\.*
        projection: \*
        primary-keys: id
  • Tabelas distribuídas: Se os dados de uma única tabela estiverem espalhados por várias partições do Kafka, ou para mesclar tabelas de diferentes partições usando sharding, defina debezium-json.distributed-tables ou canal-json.distributed-tables como true.

  • Políticas de inferência de schema: A source do Kafka suporta múltiplas políticas de inferência de schema, configuradas via schema.inference.strategy. Para detalhes, consulte Message Queue for Kafka.