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 |
|
|
Sim |
Tipo de Metastore. Defina como |
|
|
Sim |
Provedor de token. Defina como |
|
|
Sim |
URI do DLF Rest Catalog Server. Formato: |
|
|
Sim |
Nome do DLF Catalog. |
|
|
Não |
Defina como |
Observações de uso
Não adicione parâmetros de mesclagem de arquivos ou relacionados a buckets (como
bucketenum-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) ejson.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-jsonnã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-tablesoucanal-json.distributed-tablescomotrue.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.