Le connecteur Postgres CDC lit un instantané complet d'une base de données PostgreSQL, puis capture les données modifiées, avec une sémantique de traitement exactly-once.
Présentation
Le connecteur Postgres CDC prend en charge les fonctionnalités suivantes :
|
Catégorie |
Détails |
|
Types pris en charge |
Source SQL, source Flink CDC Remarque
Utilisez le connecteur JDBC pour les tables de destination (sink) et les tables de recherche (dimension). |
|
Mode d'exécution |
Streaming |
|
Format des données |
Non applicable |
|
Métriques |
|
|
Types d'API |
SQL et Flink CDC |
|
Mise à jour/suppression dans la table de destination |
Non applicable |
Fonctionnalités
À partir de la version VVR 8.0.6, le connecteur Postgres CDC s'intègre au framework d'instantané incrémentiel. Il lit les données historiques complètes, puis bascule automatiquement vers la lecture des journaux de modifications depuis le WAL, avec une sémantique exactly-once.
Principales fonctionnalités :
Traitement unifié en flux et par lots. Lecture des données complètes et incrémentielles dans une seule tâche.
Lectures d'instantanés concurrentes. Mise à l'échelle horizontale pour de meilleures performances.
Basculement transparent du mode complet vers le mode incrémentiel. Réduction automatique de la consommation de ressources.
Reprise des lectures. Reprise à partir des points d'arrêt pendant la phase d'instantané pour une stabilité accrue.
Lecture sans verrou. Aucun verrou requis, évitant tout impact sur les opérations en ligne.
Prérequis
Le connecteur Postgres CDC lit les flux CDC via la réplication logique de PostgreSQL. Il prend en charge ApsaraDB RDS for PostgreSQL, Amazon RDS for PostgreSQL et les instances PostgreSQL autogérées.
La configuration varie selon le type de déploiement. Consultez la rubrique Configuration de Postgres.
Après la configuration, vérifiez les points suivants :
Le paramètre wal_level est défini sur
logicalpour activer le décodage logique.-
L'option REPLICA IDENTITY de chaque table abonnée est définie sur
FULLafin que les événementsINSERTetUPDATEincluent les valeurs précédentes des colonnes pour garantir la cohérence des données.RemarqueREPLICA IDENTITYest un paramètre au niveau de la table PostgreSQL qui contrôle si les événementsINSERTetUPDATEincluent les valeurs précédentes des colonnes. Pour plus de détails, consultez la documentation REPLICA IDENTITY. Les valeurs des paramètres
max_wal_sendersetmax_replication_slotsdoivent être supérieures à la somme des emplacements utilisés et des emplacements requis par la tâche Flink.Le compte dispose des privilèges
SUPERUSERou des autorisationsLOGINetREPLICATION, ainsi que de l'autorisationSELECTsur les tables abonnées.
Si votre table Postgres contient des colonnes générées, définissez le paramètre publish_generated_columns sur
storedlors de la création de l'emplacement. Sinon, les schémas des phases d'instantané et incrémentielle peuvent différer.
Remarques d'utilisation
La fonctionnalité d'instantané incrémentiel nécessite VVR 8.0.6 ou une version ultérieure.
Emplacements de réplication
Les tâches Flink PostgreSQL CDC utilisent des emplacements de réplication pour empêcher la purge prématurée du WAL et garantir la cohérence des données. Une mauvaise gestion des emplacements peut entraîner une utilisation excessive du disque ou des retards de lecture. Bonnes pratiques :
-
Nettoyez rapidement les emplacements inutilisés
Flink ne supprime pas automatiquement les emplacements de réplication après l'arrêt d'une tâche ou un redémarrage sans état, afin d'éviter la perte de données WAL.
-
Si une tâche ne redémarre pas, supprimez manuellement son emplacement de réplication pour libérer de l'espace disque.
RemarqueGestion du cycle de vie : Traitez les emplacements de réplication comme des ressources au niveau de la tâche et gérez-les conjointement avec le démarrage et l'arrêt des tâches.
-
Évitez de réutiliser d'anciens emplacements
Utilisez toujours un nouveau nom d'emplacement. La réutilisation d'un ancien emplacement oblige la tâche à lire les données WAL historiques accumulées au démarrage, ce qui retarde le traitement des nouvelles données.
-
PostgreSQL requiert un emplacement par connexion. Chaque tâche doit utiliser un nom d'emplacement unique.
RemarqueConvention de nommage : Lorsque vous personnalisez
slot.name, évitez les noms avec des suffixes numériques tels quemy_slot_1pour prévenir les conflits avec les emplacements temporaires.
-
Comportement des emplacements avec les instantanés incrémentiels activés
Prérequis : Les points de contrôle (checkpoints) doivent être activés et la table source doit avoir une clé primaire définie.
-
Règles de création des emplacements :
Instantané incrémentiel désactivé : Seul un parallélisme de 1 est pris en charge. Un emplacement global est utilisé.
-
Instantané incrémentiel activé :
Phase d'instantané : Chaque sous-tâche source concurrente crée un emplacement temporaire. Le format de nommage est
${slot.name}_${task_id}.Phase incrémentielle : Tous les emplacements temporaires sont automatiquement récupérés. Seul un emplacement global est conservé.
Nombre maximal d'emplacements : Parallélisme de la source + 1 (pendant la phase d'instantané)
-
Ressources et performances
Si les emplacements disponibles ou l'espace disque sont limités, réduisez le parallélisme de l'instantané pour utiliser moins d'emplacements temporaires. Cela réduit la vitesse de lecture de l'instantané.
-
Si la table de destination en aval supporte les écritures idempotentes, définissez
scan.incremental.snapshot.backfill.skip = truepour ignorer le remplissage arrière du WAL pendant la phase d'instantané et accélérer le démarrage.Cela fournit uniquement une sémantique at-least-once et ne convient pas aux calculs avec état (agrégations ou jointures de recherche), car les modifications historiques requises peuvent être perdues.
-
Lorsque les instantanés incrémentiels sont désactivés, les points de contrôle ne sont pas pris en charge pendant la phase d'instantané.
Réutilisation de l'abonnement Postgres
Le connecteur Postgres CDC s'appuie sur une publication pour déterminer quelles modifications de tables sont poussées vers un emplacement. Si plusieurs tâches partagent la même publication, leurs configurations seront écrasées.
Cause
La valeur par défaut de publication.autocreate.mode est filtered ; elle n'inclut que les tables présentes dans la configuration du connecteur. Cela modifie la publication au démarrage de la tâche, ce qui peut affecter d'autres tâches.
Solution
-
Créez une publication dans PostgreSQL qui inclut toutes les tables surveillées, ou créez une publication distincte par tâche.
-- Create a publication named my_flink_pub that includes all tables (or specified tables, creating one publication per job) CREATE PUBLICATION my_flink_pub FOR TABLE table_a, table_b; -- Or more simply, include all tables in the database CREATE PUBLICATION my_flink_pub FOR ALL TABLES;RemarqueL'abonnement à toutes les tables n'est pas recommandé pour les bases de données volumineuses en raison d'une consommation excessive de bande passante et de CPU sur le cluster Flink.
-
Ajoutez les configurations Flink suivantes :
debezium.publication.name = 'my_flink_pub'(Spécifie le nom de la publication)debezium.publication.autocreate.mode = 'disabled'(Empêche Flink de tenter de créer ou de modifier la publication au démarrage)
Cela assure une isolation complète et empêche les nouvelles tâches d'affecter les tâches existantes.
SQL
Syntaxe
CREATE TABLE postgrescdc_source (
id INT NOT NULL,
name STRING,
description STRING,
weight DECIMAL(10,3)
) WITH (
'connector' = 'postgres-cdc',
'hostname' = '<host name>',
'port' = '<port>',
'username' = '<user name>',
'password' = '<password>',
'database-name' = '<database name>',
'schema-name' = '<schema name>',
'table-name' = '<table name>',
'decoding.plugin.name'= 'pgoutput',
'scan.incremental.snapshot.enabled' = 'true',
-- Skipping backfill can speed up reads and reduce resource usage, but it may cause data duplication. Enable this if the downstream sink is idempotent.
'scan.incremental.snapshot.backfill.skip' = 'false',
-- In a production environment, set this to 'filtered' or 'disabled' and manage the publication manually instead of through Flink.
'debezium-publication.autocreate.mode' = 'disabled'
-- If you have multiple sources, configure a different publication for each source.
--'debezium.publication.name' = 'my_flink_pub'
);
Options du connecteur
|
Option |
Description |
Type de données |
Obligatoire |
Valeur par défaut |
Remarques |
|
connector |
Nom du connecteur. |
STRING |
Oui |
– |
La valeur doit être |
|
hostname |
Adresse IP ou nom d'hôte de la base de données PostgreSQL. |
STRING |
Oui |
– |
– |
|
username |
Nom d'utilisateur pour le service de base de données PostgreSQL. |
STRING |
Oui |
– |
– |
|
password |
Mot de passe pour le service de base de données PostgreSQL. |
STRING |
Oui |
– |
– |
|
database-name |
Nom de la base de données PostgreSQL. |
STRING |
Oui |
– |
Nom de la base de données. |
|
schema-name |
Nom du schéma PostgreSQL. Les expressions régulières sont prises en charge. |
STRING |
Oui |
– |
Le nom du schéma prend en charge les expressions régulières pour lire les données depuis plusieurs schémas. |
|
table-name |
Nom de la table PostgreSQL. Les expressions régulières sont prises en charge. |
STRING |
Oui |
– |
Le nom de la table prend en charge les expressions régulières pour lire les données depuis plusieurs tables. |
|
port |
Numéro de port. |
INTEGER |
Non |
5432 |
– |
|
decoding.plugin.name |
Nom du plugin de décodage logique PostgreSQL. |
STRING |
Non |
decoderbufs |
Cette valeur dépend du plugin installé sur le service PostgreSQL. Les plugins pris en charge sont les suivants :
|
|
slot.name |
Nom du slot de décodage logique. |
STRING |
Obligatoire pour VVR 8.0.1 et versions ultérieures. Facultatif pour les versions antérieures. |
|
Définissez un Aucune valeur par défaut pour VVR 8.0.1 et versions ultérieures. |
|
debezium.* |
Propriétés et paramètres Debezium |
STRING |
Non |
– |
Offre un contrôle plus précis du comportement du client Debezium. Par exemple, |
|
scan.incremental.snapshot.enabled |
Indique s'il faut activer les snapshots incrémentiels. |
BOOLEAN |
Non |
false |
Remarque
|
|
scan.startup.mode |
Mode de démarrage pour la consommation des données. |
STRING |
Non |
initial |
Valeurs valides :
|
|
changelog-mode |
Mode de journal des modifications pour l'encodage des changements de flux. |
String |
Non |
all |
Modes de journal des modifications pris en charge :
|
|
heartbeat.interval.ms |
Intervalle d'envoi des paquets de heartbeat. |
Duration |
Non |
30s |
L'unité est la milliseconde. Le connecteur Postgres CDC envoie activement des heartbeats à la base de données pour faire avancer l'offset du slot. Lorsque les modifications de table sont peu fréquentes, la définition de cette valeur garantit une récupération rapide des journaux WAL. |
|
scan.incremental.snapshot.chunk.key-column |
Spécifie une colonne à utiliser comme clé de chunk pour le fractionnement des shards pendant la phase de snapshot. |
STRING |
Non |
– |
Par défaut, la première colonne de la clé primaire est sélectionnée. |
|
scan.incremental.close-idle-reader.enabled |
Indique s'il faut fermer les lecteurs inactifs une fois le snapshot terminé. |
Boolean |
Non |
false |
Pour activer cette configuration, définissez |
|
scan.incremental.snapshot.backfill.skip |
Indique s'il faut ignorer la lecture des journaux pendant la phase de snapshot. |
Boolean |
Non |
false |
Valeurs valides :
|
Mappages de types
Mappages de types de PostgreSQL vers Flink :
|
PostgreSQL CDC |
Flink |
|
SMALLINT |
SMALLINT |
|
INT2 |
|
|
SMALLSERIAL |
|
|
SERIAL2 |
|
|
INTEGER |
INT |
|
SERIAL |
|
|
BIGINT |
BIGINT |
|
BIGSERIAL |
|
|
REAL |
FLOAT |
|
FLOAT4 |
|
|
FLOAT8 |
DOUBLE |
|
DOUBLE PRECISION |
|
|
NUMERIC(p, s) |
DECIMAL(p, s) |
|
DECIMAL(p, s) |
|
|
BOOLEAN |
BOOLEAN |
|
DATE |
DATE |
|
TIME [(p)] [WITHOUT TIMEZONE] |
TIME [(p)] [WITHOUT TIMEZONE] |
|
TIMESTAMP [(p)] [WITHOUT TIMEZONE] |
TIMESTAMP [(p)] [WITHOUT TIMEZONE] |
|
CHAR(n) |
STRING |
|
CHARACTER(n) |
|
|
VARCHAR(n) |
|
|
CHARACTER VARYING(n) |
|
|
TEXT |
|
|
BYTEA |
BYTES |
Exemple
CREATE TABLE source (
id INT NOT NULL,
name STRING,
description STRING,
weight DECIMAL(10,3)
) WITH (
'connector' = 'postgres-cdc',
'hostname' = '<host name>',
'port' = '<port>',
'username' = '<user name>',
'password' = '<password>',
'database-name' = '<database name>',
'schema-name' = '<schema name>',
'table-name' = '<table name>'
);
SELECT * FROM source;
Flink CDC
VVR V11.4+ prend en charge le connecteur PostgreSQL en tant que source Flink CDC.
Syntaxe
source:
type: postgres
name: PostgreSQL Source
hostname: localhost
port: 5432
username: pg_username
password: pg_password
tables: db.scm.tbl
slot.name: test_slot
scan.startup.mode: initial
server-time-zone: UTC
connect.timeout: 120s
decoding.plugin.name: decoderbufs
sink:
type: ...
Options du connecteur
|
Option |
Description |
Obligatoire |
Type de données |
Valeur par défaut |
Remarques |
|
type |
Le nom du connecteur. |
Oui |
STRING |
– |
Doit être |
|
name |
Le nom de la source de données. |
Non |
STRING |
– |
– |
|
hostname |
Le nom de domaine ou l'adresse IP du serveur de base de données PostgreSQL. |
Oui |
STRING |
– |
– |
|
port |
Le port de la base de données PostgreSQL. |
Non |
INTEGER |
5432 |
– |
|
username |
Le nom d'utilisateur PostgreSQL. |
Oui |
STRING |
– |
– |
|
password |
Le mot de passe PostgreSQL. |
Oui |
STRING |
– |
– |
|
tables |
Les noms des tables à capturer. Les expressions régulières sont prises en charge. |
Oui |
STRING |
– |
Important
Actuellement, seules les tables appartenant à la même base de données peuvent être capturées. Un point (.) est traité comme un séparateur pour un nom entièrement qualifié. Pour utiliser un point (.) dans une expression régulière afin de correspondre à n'importe quel caractère, échappez-le avec une barre oblique inverse. Par exemple : |
|
slot.name |
Le nom du slot de réplication PostgreSQL. |
Oui |
STRING |
– |
Le nom doit respecter les règles de dénomination des slots de réplication PostgreSQL et ne peut contenir que des lettres minuscules, des chiffres et des traits de soulignement. |
|
decoding.plugin.name |
Le nom du plugin de décodage logique PostgreSQL installé sur le serveur. |
Non |
STRING |
|
Valeurs valides : |
|
tables.exclude |
Les tables à exclure. Cette option prend effet après l'option |
Non |
STRING |
– |
Voir l'option |
|
server-time-zone |
Le fuseau horaire de la session du serveur de base de données, par exemple « Asia/Shanghai ». |
Non |
STRING |
– |
Si ce paramètre n'est pas défini, le fuseau horaire par défaut du système ( |
|
scan.incremental.snapshot.chunk.size |
La taille (nombre de lignes) de chaque fragment dans le framework de snapshot incrémentiel. |
Non |
INTEGER |
8096 |
Lorsque le snapshot incrémentiel est activé, la table est divisée en plusieurs fragments pour la lecture. Les données d'un fragment sont mises en cache en mémoire avant d'être entièrement consommées. Une taille de fragment plus petite entraîne un nombre total de fragments plus élevé pour la table. Bien que cela réduise la granularité de la récupération après incident, cela peut provoquer des erreurs de dépassement de mémoire (OOM) et diminuer le débit global. Il est donc nécessaire de trouver un équilibre et de définir une taille de fragment raisonnable. |
|
scan.snapshot.fetch.size |
Le nombre maximal d'enregistrements à récupérer à la fois lors de la lecture des données complètes d'une table. |
Non |
INTEGER |
1024 |
– |
|
scan.startup.mode |
Le mode de démarrage pour la consommation des données. |
Non |
STRING |
initial |
Valeurs valides :
|
|
scan.incremental.close-idle-reader.enabled |
Indique s'il faut fermer les lecteurs inactifs une fois le snapshot terminé. |
Non |
BOOLEAN |
false |
Pour activer cette configuration, définissez |
|
scan.lsn-commit.checkpoints-num-delay |
Le nombre de checkpoints à retarder avant de commencer à valider les offsets LSN. |
Non |
INTEGER |
3 |
Les offsets LSN des checkpoints sont validés de manière rotative pour éviter l'impossibilité de récupérer l'état. |
|
connect.timeout |
Le temps maximal pendant lequel le connecteur attend la connexion au serveur de base de données PostgreSQL avant d'expirer. |
Non |
DURATION |
30s |
Cette valeur ne peut pas être inférieure à 250 millisecondes. |
|
connect.max-retries |
Nombre maximal de tentatives de réessai pour que le connecteur établisse une connexion. |
Non |
INTEGER |
3 |
– |
|
connection.pool.size |
La taille du pool de connexions. |
Non |
INTEGER |
20 |
– |
|
jdbc.properties.* |
Permet aux utilisateurs de transmettre des propriétés d'URL JDBC personnalisées. |
Non |
STRING |
20 |
Les utilisateurs peuvent transmettre des propriétés personnalisées, telles que |
|
heartbeat.interval |
L'intervalle d'envoi des événements de heartbeat pour suivre le dernier offset de journal WAL disponible. |
Non |
DURATION |
30s |
– |
|
debezium.* |
Transmet les propriétés Debezium au moteur intégré Debezium, utilisé pour capturer les modifications de données depuis le serveur PostgreSQL. |
Non |
STRING |
– |
Propriétés du connecteur Debezium PostgreSQL : Documentation Debezium. |
|
chunk-meta.group.size |
La taille des métadonnées de fragment. |
Non |
STRING |
1000 |
Si les métadonnées dépassent cette valeur, elles sont transmises en plusieurs parties. |
|
metadata.list |
Une liste de métadonnées lisibles transmises en aval, qui peuvent être utilisées dans le module de transformation. |
Non |
STRING |
false |
Utilisez des virgules (,) comme séparateurs. Actuellement, les métadonnées disponibles sont : |
|
scan.incremental.snapshot.unbounded-chunk-first.enabled |
Envoie le fragment illimité en premier pendant la phase de lecture du snapshot. |
Non |
STRING |
false |
Il s'agit d'une fonctionnalité expérimentale. Son activation peut réduire le risque d'erreurs OOM lorsque le TaskManager synchronise le dernier fragment pendant la phase de snapshot. Nous recommandons d'ajouter cette option avant le premier démarrage du job. |