Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Distribute data in real time with Flink CDC and Kafka

Última atualização: Jun 27, 2026

Este guia mostra como construir dois pipelines do Flink CDC que, juntos, movem dados do MySQL para um data lake em tempo real.

Este guia mostra como construir dois pipelines do Flink CDC que, juntos, movem dados do MySQL para um data lake em tempo real:

  • Pipeline 1 — MySQL para Kafka: Capture logs binários do MySQL e transmita eventos de alteração para tópicos do ApsaraMQ for Kafka.

  • Pipeline 2 — Kafka para DLF: Leia esses eventos de alteração do Kafka e ingira-os no Data Lake Formation (DLF) por meio de um sink Paimon.

MySQL (binlog) → [Pipeline 1] → ApsaraMQ for Kafka → [Pipeline 2] → DLF (Paimon)

O Kafka atua como um buffer durável e reproduzível entre os dois jobs, desacoplando a captura da ingestão e permitindo dimensionamento, auditoria e reprodução independentes.

Pré-requisitos

  • O MySQL deve ter o log binário ativado no formato ROW.

  • Um cluster do ApsaraMQ for Kafka deve estar disponível e acessível pelo job do Flink.

Pipeline 1: Transmita logs binários do MySQL para o Kafka

Visão geral

Este pipeline conecta-se a uma instância do MySQL, lê eventos de alteração de log binário de todas as tabelas no banco de dados kafka_test e os grava em tópicos do ApsaraMQ for Kafka. Cada tabela de origem é mapeada para seu próprio tópico do Kafka.

Configuração mínima

A configuração abaixo é o mínimo necessário para executar o pipeline:

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: kafka_test\.\*
  server-id: 8601-8604

sink:
  type: kafka
  name: Kafka Sink
  properties.bootstrap.servers: ${kafka.bootstraps.server}
  properties.enable.idempotence: false
Aviso

Defina properties.enable.idempotence: false na configuração do sink. O ApsaraMQ for Kafka não suporta gravações idempotentes ou transacionais.

Configuração completa com parâmetros opcionais

A configuração a seguir estende a configuração mínima acima com parâmetros opcionais. Adicione apenas os parâmetros relevantes para o seu caso de uso.

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: kafka_test\.\*
  server-id: 8601-8604
  # (Optional) Sync existing and incremental data from newly added tables.
  scan.newly-added-table.enabled: true
  # (Optional) Sync table and column comments.
  include-comments.enabled: true
  # (Optional) Prioritize distributing unbounded chunks during snapshot.
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (Optional) Deserialize changelog only for captured tables.
  scan.only.deserialize.captured.tables.changelog.enabled: true
  # (Optional) Include the data change timestamp as a metadata column.
  metadata-column.include-list: op_ts

sink:
  type: kafka
  name: Kafka Sink
  properties.bootstrap.servers: ${kafka.bootstraps.server}
  # Required: ApsaraMQ for Kafka does not support idempotent or transactional writes.
  properties.enable.idempotence: false
  # (Optional) Set the mapping between upstream tables and Kafka topics.
  sink.tableId-to-topic.mapping: kafka_test.customers:customers;kafka_test.products:products

Nomenclatura e mapeamento de tópicos

Por padrão, cada tabela de origem é mapeada para um tópico do Kafka nomeado no formato database.table. Por exemplo, dados de kafka_test.customers são gravados em um tópico chamado kafka_test.customers.

Para usar nomes de tópicos personalizados, defina sink.tableId-to-topic.mapping. Essa opção mapeia identificadores de tabelas de origem para nomes de tópicos do Kafka, preservando o nome original da tabela no payload do evento para fins de rastreabilidade. O exemplo acima mapeia kafka_test.customers para um tópico chamado customers e kafka_test.products para um tópico chamado products.

O ApsaraMQ for Kafka suporta os formatos de saída json, canal-json e debezium-json (padrão).

Para todas as opções de configuração disponíveis, consulte Conector do Kafka.

Estratégia de particionamento

Por padrão, todos os eventos de alteração são gravados na partição 0 de cada tópico (all-to-zero). Isso preserva a ordenação global dentro de um tópico, mas não distribui a carga entre as partições.

Para distribuir registros em várias partições com base no hash da chave primária, adicione o seguinte à configuração do seu sink:

sink:
  partition.strategy: hash-by-key

Com hash-by-key, todos os eventos de alteração para um determinado valor de chave primária sempre vão para a mesma partição. Isso preserva a ordenação por chave e permite o consumo paralelo.

Saída esperada

Ao executar uma instrução UPDATE na tabela customers, o job produz uma mensagem do Kafka no formato configurado.

Formato Debezium JSON (debezium-json, padrão):

{
  "before": {
    "id": 4,
    "name": "John",
    "address": "New York",
    "phone_number": "2222",
    "age": 12
  },
  "after": {
    "id": 4,
    "name": "John",
    "address": "New York",
    "phone_number": "1234",
    "age": 12
  },
  "op": "u",
  "source": {
    "db": "kafka_test",
    "table": "customers",
    "ts_ms": 1728528674000
  }
}

Formato Canal JSON (canal-json):

{
  "old": [
    {
      "id": 4,
      "name": "John",
      "address": "New York",
      "phone_number": "2222",
      "age": 12
    }
  ],
  "data": [
    {
      "id": 4,
      "name": "John",
      "address": "New York",
      "phone_number": "1234",
      "age": 12
    }
  ],
  "type": "UPDATE",
  "database": "kafka_test",
  "table": "customers",
  "pkNames": [
    "id"
  ],
  "ts": 1728528674000,
  "es": 0
}
Nota
  • O formato Debezium JSON omite informações de chave primária no payload da mensagem. Se um consumidor downstream precisar identificar chaves primárias, use o formato Canal JSON ou defina-as explicitamente na seção transform de um pipeline downstream (conforme mostrado no Pipeline 2 abaixo).

Pipeline 2: Ingerir eventos de alteração do Kafka no DLF

Visão geral

Este pipeline lê eventos de alteração de um tópico do Kafka e os grava no Data Lake Formation (DLF) usando um sink Paimon. Ele foi projetado para consumir a saída produzida pelo Pipeline 1.

Considere um tópico do Kafka chamado inventory que contém dados de log binário no formato debezium-json de duas tabelas: customers e products. Quando um único tópico do Kafka transporta eventos de alteração de várias tabelas de origem, defina debezium-json.distributed-tables: true para que o conector possa rotear corretamente cada registro para sua tabela de destino.

Configuração

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

transform:
  - source-table: \.*.\.*
    projection: \*
    primary-keys: id

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

Detalhes importantes da configuração

Parâmetro

Observações

debezium-json.distributed-tables

Defina como true quando um tópico do Kafka contiver registros de mais de uma tabela de origem. Sem essa configuração, o conector não consegue determinar a qual tabela de destino cada registro pertence. Para o formato Canal JSON, use canal-json.distributed-tables: true.

transform / primary-keys

O formato Debezium JSON omite informações de chave primária no payload da mensagem. É obrigatório definir as chaves primárias na seção transform para que o Paimon possa realizar upserts corretamente. O exemplo aplica id como chave primária para todas as tabelas correspondentes a \.*.\.*.

scan.startup.mode: earliest-offset

Lê desde o início de cada tópico do Kafka na inicialização.

schema.inference.strategy

Configure uma estratégia de inferência de schema para a origem Kafka com a opção schema.inference.strategy. Para mais informações, consulte Conector do Kafka.

Notas de uso

Nota
  • O conector de origem Kafka suporta leitura de dados nos formatos canal-json e debezium-json (padrão).

  • Ative debezium-json.distributed-tables: true (ou canal-json.distributed-tables: true) quando dados de várias tabelas de origem estiverem distribuídos entre partições do mesmo tópico ou quando for necessário mesclar dados distribuídos.

  • Como o Debezium JSON omite informações de chave primária, sempre defina as chaves primárias explicitamente no módulo transform ao usar esse formato com um sink Paimon.