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
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
}
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
transformde 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 |
|
|
Defina como |
|
|
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 |
|
|
Lê desde o início de cada tópico do Kafka na inicialização. |
|
|
Configure uma estratégia de inferência de schema para a origem Kafka com a opção |
Notas de uso
O conector de origem Kafka suporta leitura de dados nos formatos
canal-jsonedebezium-json(padrão).Ative
debezium-json.distributed-tables: true(oucanal-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
transformao usar esse formato com um sink Paimon.