Este tópico aborda as melhores práticas para transmitir alterações em tempo real do MySQL ou Kafka para um data warehouse Hologres por meio de um job de ingestão de dados Change Data Capture (CDC) YAML.
Sincronizar do MySQL para o Hologres
O conector do Hologres oferece suporte a:
Criação automática de tabelas e sincronização de schema
Armazenamento híbrido linha-coluna e leituras incrementais baseadas em Binlog
Mapeamento de tipos tolerante para lidar com alterações de schema na origem sem falhas no job
O YAML a seguir define a configuração mínima para sincronizar um banco de dados MySQL inteiro com o Hologres:
source:
type: mysql
name: MySQL Source
hostname: localhost
port: 3306
username: username
password: password
tables: holo_test\.*
server-id: 8601-8604
sink:
type: hologres
name: Hologres Sink
endpoint: ****.hologres.aliyuncs.com:80
dbname: cdcyaml_test
username: ${secret_values.holo-username}
password: ${secret_values.holo-password}
sink.type-normalize-strategy: BROADEN
Parâmetros opcionais para defina na primeira execução do job:
|
Parâmetro |
Finalidade |
Referência |
|
|
Sincroniza comentários de tabela e de campo |
|
|
|
Prioriza shards ilimitados para evitar erros de OutOfMemory no TaskManager |
|
|
|
Ativa filtros de análise para acelerar a leitura de dados |
Lidar com alterações de tipo na origem
O conector do Hologres não suporta eventos de alteração de tipo de coluna. Para gerenciar mudanças de schema na origem sem causar falhas no job, configure sink.type-normalize-strategy para mapear tipos do MySQL para tipos mais abrangentes no Hologres. O valor padrão é STANDARD.
|
Estratégia |
Comportamento |
|
|
Mapeamento de tipos padrão |
|
|
Mapeia para tipos compatíveis mais amplos |
|
|
Converte todos os tipos para |
Por exemplo, ao usar ONLY_BIGINT_OR_TEXT, uma alteração de tipo de coluna de INT para BIGINT no MySQL mapeia ambos os tipos para int8 no Hologres — o job continua normalmente, sem reportar erro de conversão de tipo.
source:
type: mysql
name: MySQL Source
hostname: localhost
port: 3306
username: username
password: password
tables: holo_test\.*
server-id: 8601-8604
# (Optional) Synchronize table comments and field comments.
include-comments.enabled: true
# (Optional) Prioritize the distribution of unbounded shards to prevent potential TaskManager OutOfMemory errors.
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
# (Optional) Enable parsing filters to accelerate data reads.
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: hologres
name: Hologres Sink
endpoint: ****.hologres.aliyuncs.com:80
dbname: cdcyaml_test
username: ${secret_values.holo-username}
password: ${secret_values.holo-password}
sink.type-normalize-strategy: ONLY_BIGINT_OR_TEXT
Para obter a lista completa de mapeamentos de tipos suportados, consulte Mapeamento de tipos do conector Hologres para jobs YAML de ingestão de dados.
Gravar em tabelas particionadas
Ao utilizar o conector do Hologres como sink, é possível gravar dados em tabelas particionadas. Para mais informações, consulte Gravar dados em tabelas particionadas.
Sincronizar do Kafka para o Hologres
Esta seção pressupõe que as alterações do MySQL já estejam fluindo para o Kafka — por exemplo, usando a configuração descrita em Implementar distribuição em tempo real. O pipeline lê do Kafka e grava no Hologres.
Considere um tópico do Kafka chamado inventory contendo dados de duas tabelas, customers e products, no formato debezium-json. O job a seguir sincroniza ambas as tabelas com suas respectivas tabelas correspondentes no Hologres:
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: hologres
name: Hologres Sink
endpoint: ****.hologres.aliyuncs.com:80
dbname: cdcyaml_test
username: ${secret_values.holo-username}
password: ${secret_values.holo-password}
sink.type-normalize-strategy: ONLY_BIGINT_OR_TEXT
transform:
- source-table: \.*.\.*
projection: \*
primary-keys: id
Formatos suportados
A source do Kafka aceita três formatos de mensagem:
|
Formato |
Observações |
|
|
Formato JSON padrão |
|
|
Formato JSON compatível com Canal |
|
|
Formato padrão para dados CDC originados do Debezium |
Requisito de chave primária para debezium-json
Quando o formato é debezium-json, a source do Kafka não infere chaves primárias automaticamente. Adicione uma chave primária explicitamente usando uma regra de transformação:
transform:
- source-table: \.*.\.*
projection: \*
primary-keys: id
Tabelas distribuídas
Se os dados de uma única tabela abrangerem várias partições, ou se for necessário mesclar tabelas fragmentadas em partições diferentes, defina debezium-json.distributed-tables ou canal-json.distributed-tables como true.
Inferência de schema
A source do Kafka oferece suporte a diversas políticas de inferência de schema. Configure a política desejada por meio de schema.inference.strategy. Para mais detalhes, consulte Message Queue for Apache Kafka.