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
Crie um rascunho YAML para ingerir dados com o Flink CDC. Para mais informações, consulte Develop Flink CDC jobs for data ingestion (Beta).
-
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: trueSubstitua 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.metastoreO tipo de metastore. Defina como
rest.Sim
—
restcatalog.properties.token.providerO provedor de token. Defina como
dlf.Sim
—
dlfcatalog.properties.uriO 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.comcatalog.properties.warehouseO nome do catálogo Paimon.
Sim
—
dlf_testcommit.userO usuário de commit para gravações de dados. Atribua usuários de commit exclusivos a jobs distintos para evitar conflitos.
Não
adminyour_job_nametable.properties.deletion-vectors.enabledAtiva 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 umcommit.userexclusivo 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
bucketenum-sorted-run.compaction-trigger.Deletion vectors: Defina
table.properties.deletion-vectors.enabled: truepara 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 |
|
|
Sincroniza dados de tabelas criadas durante a fase incremental. |
|
|
Sincroniza comentários de tabelas e campos. |
|
|
Evita possíveis erros de OOM no TaskManager priorizando shards ilimitados. |
|
|
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) ejson.Ao ingerir dados de múltiplas partições Kafka em uma única tabela no DLF, defina
debezium-json.distributed-tablesoucanal-json.distributed-tablescomotrue.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.