Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Implementar ingestão em tempo real

Última atualização: Jun 27, 2026

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

include-comments.enabled: true

Sincroniza comentários de tabela e de campo

Conector do MySQL

scan.incremental.snapshot.unbounded-chunk-first.enabled: true

Prioriza shards ilimitados para evitar erros de OutOfMemory no TaskManager

Conector do MySQL

scan.only.deserialize.captured.tables.changelog.enabled: true

Ativa filtros de análise para acelerar a leitura de dados

Conector do MySQL

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

STANDARD

Mapeamento de tipos padrão

BROADEN

Mapeia para tipos compatíveis mais amplos

ONLY_BIGINT_OR_TEXT

Converte todos os tipos para int8 ou text

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

json

Formato JSON padrão

canal-json

Formato JSON compatível com Canal

debezium-json

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.