Configure a evolução de schema para jobs de ingestão de dados do Flink CDC e controle como as alterações de schema no source são aplicadas ao sink.
Eventos de alteração de schema suportados
Os jobs de ingestão de dados do Flink CDC sincronizam as seguintes alterações de schema do source para o sink:
Criação de tabela —
CREATE TABLE ...Adição de coluna —
ALTER TABLE ... ADD COLUMN ...Alteração de tipo de coluna —
ALTER TABLE ... MODIFY COLUMN ...Remoção de coluna —
ALTER TABLE ... DROP COLUMN ...Renomeação de coluna —
ALTER TABLE ... RENAME COLUMN ...Truncamento de tabela —
TRUNCATE TABLE ...Exclusão de tabela —
DROP TABLE ...
O framework suporta apenas os tipos de alteração de schema listados acima. Alterações não suportadas causam exceções no job e exigem uma reinicialização sem estado para recuperação.
Comportamento da evolução de schema
Defina schema.change.behavior no módulo pipeline para controlar como o Flink CDC lida com as alterações de schema:
pipeline:
schema.change.behavior: EVOLVE
LENIENT (padrão)
Converte alterações de schema não suportadas em operações compatíveis com o sink:
rename.column: Envia um eventoalter.column.typepara tornar a coluna original anulável e, em seguida, um eventoadd.columnpara adicionar uma nova coluna anulável com o novo nome. A coluna original é mantida.drop.column: Envia um eventoalter.column.typee define o tipo da coluna como anulável, em vez de removê-la.Novas colunas: O sistema ainda envia o evento de adição de coluna, mas o tipo do campo passa a ser anulável.
drop.tableetruncate.table: Não são enviados ao sink.
Utilize LENIENT quando desejar que o job sincronize alterações de schema da forma mais automática possível, com máxima compatibilidade.
IGNORE
Toda a evolução de schema é ignorada. As alterações de schema upstream não são aplicadas à tabela sink downstream, e os dados continuam fluindo pelas colunas existentes.
Use IGNORE quando o sink não suportar alterações de schema ou quando você quiser continuar recebendo dados das colunas existentes sem modificar o schema do sink.
EVOLVE
Todas as alterações de schema são aplicadas à tabela sink exatamente como ocorrem no source. Se a aplicação de uma alteração falhar, o job lançará uma exceção e acionará uma reinicialização por falha.
Se o sink não conseguir processar um evento de alteração de schema, o job poderá falhar e não se recuperará automaticamente.
Opte por EVOLVE quando for necessária uma sincronização de schema estrita e exata.
TRY_EVOLVE
Tenta aplicar as alterações de schema à tabela sink. Caso o sink não consiga processar a mudança, o job não falha nem reinicia — ele continua e tenta lidar com os dados afetados por meio de transformação.
Se a aplicação de uma alteração de schema falhar, os dados subsequentes poderão perder colunas ou ser truncados para corresponder ao schema existente do sink.
Escolha TRY_EVOLVE quando desejar uma sincronização de schema rigorosa, porém com alguma tolerância a falhas.
EXCEPTION
Lança uma exceção ao receber qualquer evento de alteração de schema. Nenhuma alteração de schema é permitida.
Adote EXCEPTION quando for obrigatório garantir que apenas dados — e não o schema — sejam sincronizados.
Controle de alterações de schema no sink
Para um controle refinado, utilize include.schema.changes e exclude.schema.changes no módulo sink a fim de filtrar quais tipos de eventos de alteração de schema chegam ao sink.
|
Opção de configuração |
Obrigatório |
Tipo de dado |
Valor padrão |
Observação |
|
|
Não |
|
Sem valor padrão |
Todas as alterações são suportadas por padrão |
|
|
Não |
|
Sem valor padrão |
Tem prioridade sobre |
Tipos de evento configuráveis
|
Tipo de evento |
Descrição |
|
|
Adicionar uma coluna |
|
|
Alterar o tipo da coluna |
|
|
Criar uma tabela |
|
|
Remover uma coluna |
|
|
Excluir uma tabela |
|
|
Renomear uma coluna |
|
|
Truncar uma tabela |
Há suporte para correspondência parcial. Por exemplo, especificar drop corresponde tanto a drop.column quanto a drop.table. Especificar table corresponde a create.table, truncate.table e drop.table. Já especificar column corresponde a add.column, alter.column.type, rename.column e drop.column.
Exemplos
Exemplo 1: Aplicar todas as alterações de schema
Defina schema.change.behavior como EVOLVE no módulo pipeline:
source:
type: mysql
name: MySQL Source
hostname: ${mysql.hostname}
port: ${mysql.port}
username: ${mysql.username}
password: ${mysql.password}
tables: ${mysql.source.table}
server-id: 7601-7604
sink:
type: values
name: Values Sink
print.enabled: true
sink.print.logger: true
pipeline:
name: mysql to print job
schema.change.behavior: EVOLVE
Exemplo 2: Aplicar apenas tipos selecionados de alteração de schema
Permita a criação de tabelas e todos os eventos de coluna, mas exclua drop.column:
source:
type: mysql
name: MySQL Source
hostname: ${mysql.hostname}
port: ${mysql.port}
username: ${mysql.username}
password: ${mysql.password}
tables: ${mysql.source.table}
server-id: 7601-7604
sink:
type: values
name: Values Sink
print.enabled: true
sink.print.logger: true
include.schema.changes: [create.table, column] # `column` matches add.column, alter.column.type, rename.column, and drop.column
exclude.schema.changes: [drop.column] # Excludes drop.column even though it is matched by `column`
pipeline:
name: mysql to print job
schema.change.behavior: EVOLVE