Configurez l'évolution du schéma pour les jobs d'ingestion de données Flink CDC afin de contrôler l'application des modifications de schéma de la source au sink.
Événements de modification de schéma pris en charge
Les jobs d'ingestion de données Flink CDC peuvent synchroniser les modifications de schéma suivantes de la source vers le sink :
Création de table —
CREATE TABLE ...Ajout de colonne —
ALTER TABLE ... ADD COLUMN ...Modification du type de colonne —
ALTER TABLE ... MODIFY COLUMN ...Suppression de colonne —
ALTER TABLE ... DROP COLUMN ...Renommage de colonne —
ALTER TABLE ... RENAME COLUMN ...Troncation de table —
TRUNCATE TABLE ...Suppression de table —
DROP TABLE ...
Le framework ne prend en charge que les types de modification de schéma répertoriés ci-dessus. Les modifications non prises en charge provoquent des exceptions de job et nécessitent un redémarrage sans état pour récupérer.
Comportement de l'évolution du schéma
Définissez schema.change.behavior dans le module pipeline pour contrôler la gestion des modifications de schéma par Flink CDC :
pipeline:
schema.change.behavior: EVOLVE
LENIENT (par défaut)
Convertit les modifications de schéma non prises en charge en opérations compatibles avec le sink :
rename.column: Envoie un événementalter.column.typepour rendre la colonne d'origine nullable, puis un événementadd.columnpour ajouter une nouvelle colonne nullable portant le nouveau nom. La colonne d'origine est conservée.drop.column: Envoie un événementalter.column.typeet définit le type de colonne comme nullable au lieu de la supprimer.Nouvelles colonnes : Le système envoie toujours l'événement d'ajout de colonne, mais le type de champ devient nullable.
drop.tableettruncate.table: Ne sont pas envoyés au sink.
Utilisez LENIENT pour synchroniser les modifications de schéma de manière aussi automatique que possible, avec une compatibilité maximale.
IGNORE
Toute évolution du schéma est ignorée. Les modifications de schéma en amont ne sont pas appliquées à la table sink en aval et les données continuent de circuler à partir des colonnes existantes.
Utilisez IGNORE si votre sink ne prend pas en charge les modifications de schéma ou si vous souhaitez continuer à recevoir des données des colonnes existantes sans modifier le schéma du sink.
EVOLVE
Toutes les modifications de schéma sont appliquées à la table sink exactement telles qu'elles se produisent dans la source. Si l'application d'une modification échoue, le job génère une exception et déclenche un redémarrage après échec.
Si le sink ne peut pas traiter un événement de modification de schéma, le job peut échouer sans pouvoir se rétablir automatiquement.
Utilisez EVOLVE pour une synchronisation stricte et exacte du schéma.
TRY_EVOLVE
Tente d'appliquer les modifications de schéma à la table sink. Si le sink ne peut pas traiter la modification, le job n'échoue pas et ne redémarre pas ; il continue et tente de gérer les données affectées en les transformant.
Si une modification de schéma ne peut pas être appliquée, les données suivantes peuvent perdre des colonnes ou être tronquées pour correspondre au schéma du sink existant.
Utilisez TRY_EVOLVE pour une synchronisation stricte du schéma avec une certaine tolérance aux pannes.
EXCEPTION
Génère une exception lors de la réception de tout événement de modification de schéma. Aucune modification de schéma n'est autorisée.
Utilisez EXCEPTION pour garantir que seules les données — et non le schéma — sont synchronisées.
Contrôle des modifications de schéma au niveau du sink
Pour un contrôle plus fin, utilisez include.schema.changes et exclude.schema.changes dans le module sink pour filtrer les types d'événements de modification de schéma qui atteignent le sink.
| Option de configuration | Requis | Type de données | Valeur par défaut | Remarque |
|---|---|---|---|---|
include.schema.changes |
Non | List<String> |
Aucune valeur par défaut | Toutes les modifications sont prises en charge par défaut |
exclude.schema.changes |
Non | List<String> |
Aucune valeur par défaut | Prioritaire sur include.schema.changes |
Types d'événements configurables
| Type d'événement | Description |
|---|---|
add.column |
Ajouter une colonne |
alter.column.type |
Modifier le type de colonne |
create.table |
Créer une table |
drop.column |
Supprimer une colonne |
drop.table |
Supprimer une table |
rename.column |
Renommer une colonne |
truncate.table |
Tronquer une table |
La correspondance partielle est prise en charge. Par exemple, spécifier drop correspond à la fois à drop.column et à drop.table. Spécifier table correspond à create.table, truncate.table et drop.table. Spécifier column correspond à add.column, alter.column.type, rename.column et drop.column.
Exemples
Exemple 1 : Appliquer toutes les modifications de schéma
Définissez schema.change.behavior sur EVOLVE dans le module 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
Exemple 2 : Appliquer uniquement les types de modification de schéma sélectionnés
Autorisez la création de table et tous les événements de colonne, mais excluez 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