Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:PostgreSQL CDC

Dernière mise à jour :Aug 20, 2026

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

Métriques de surveillance

  • currentFetchEventTimeLag : l'intervalle entre la génération des données et leur extraction par l'opérateur source.

  • currentEmitEventTimeLag : l'intervalle entre la génération des données et leur sortie de l'opérateur source.

  • sourceIdleTime : la durée pendant laquelle la source n'a pas produit de nouvelles données.

Remarque
  • Les métriques currentFetchEventTimeLag et currentEmitEventTimeLag ne sont valides que pendant la phase incrémentielle. Pendant la phase d'instantané, leur valeur est toujours de 0.

  • Pour plus d'informations sur les métriques, consultez la rubrique Description des 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.

Important

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 logical pour activer le décodage logique.

  • L'option REPLICA IDENTITY de chaque table abonnée est définie sur FULL afin que les événements INSERT et UPDATE incluent les valeurs précédentes des colonnes pour garantir la cohérence des données.

    Remarque

    REPLICA IDENTITY est un paramètre au niveau de la table PostgreSQL qui contrôle si les événements INSERT et UPDATE incluent les valeurs précédentes des colonnes. Pour plus de détails, consultez la documentation REPLICA IDENTITY.

  • Les valeurs des paramètres max_wal_senders et max_replication_slots doivent être supérieures à la somme des emplacements utilisés et des emplacements requis par la tâche Flink.

  • Le compte dispose des privilèges SUPERUSER ou des autorisations LOGIN et REPLICATION, ainsi que de l'autorisation SELECT sur les tables abonnées.

  • Si votre table Postgres contient des colonnes générées, définissez le paramètre publish_generated_columns sur stored lors 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.

      Remarque

      Gestion 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.

      Remarque

      Convention de nommage : Lorsque vous personnalisez slot.name, évitez les noms avec des suffixes numériques tels que my_slot_1 pour 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 = true pour 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é.

    Configuration pour éviter les délais d'expiration pendant la phase d'instantané

    Lorsque les instantanés incrémentiels sont désactivés, les points de contrôle pendant la phase d'instantané peuvent provoquer des basculements dus à des délais d'expiration. Configurez ces paramètres dans Other Configuration (Configuration des paramètres d'exécution personnalisés) :

    execution.checkpointing.interval: 10min
    execution.checkpointing.tolerable-failed-checkpoints: 100
    restart-strategy: fixed-delay
    restart-strategy.fixed-delay.attempts: 2147483647

    Paramètres :

    Paramètre

    Description

    Remarques

    execution.checkpointing.interval

    L'intervalle entre les points de contrôle.

    L'unité est une valeur de durée, telle que 10min ou 30s.

    execution.checkpointing.tolerable-failed-checkpoints

    Le nombre d'échecs de point de contrôle tolérés avant l'échec de la tâche.

    Le produit de ce paramètre et de l'intervalle de planification des points de contrôle correspond au temps de lecture d'instantané autorisé.

    Remarque

    Si la table est très volumineuse, attribuez une valeur plus élevée à ce paramètre.

    restart-strategy

    La stratégie de redémarrage de la tâche.

    Valeurs valides :

    • fixed-delay : stratégie de redémarrage à délai fixe.

    • failure-rate : stratégie de redémarrage basée sur le taux d'échec.

    • exponential-delay : stratégie de redémarrage à délai exponentiel.

    Stratégies de redémarrage.

    restart-strategy.fixed-delay.attempts

    Nombre maximal de tentatives de redémarrage pour la stratégie de redémarrage fixed-delay.

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

  1. 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;
    Remarque

    L'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.

  2. 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 postgres-cdc.

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 :

  • decoderbufs : pris en charge sur PostgreSQL 9.6 et versions ultérieures. Ce plugin doit être installé.

  • pgoutput (recommandé) : plugin intégré officiel pour PostgreSQL 10 et versions ultérieures.

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.

flink (versions antérieures à 8.0.1)

Définissez un slot.name unique pour chaque table afin d'éviter l'erreur PSQLException: ERROR: replication slot "debezium" is active for PID 974. Slots de réplication.

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, 'debezium.snapshot.mode' = 'never'. Propriétés de configuration.

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 :

  • initial : analyse l'intégralité des données historiques lors du premier démarrage, puis lit les dernières données WAL.

  • latest-offset : n'analyse pas l'intégralité des données historiques lors du premier démarrage. La lecture commence à la fin du journal WAL, ce qui signifie que seules les modifications les plus récentes effectuées après le démarrage du connecteur sont lues.

  • snapshot : analyse l'intégralité des données historiques, lit les nouvelles données WAL générées pendant la phase de snapshot, puis le job s'arrête.

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 :

  • ALL : prend en charge tous les types, y compris INSERT, DELETE, UPDATE_BEFORE et UPDATE_AFTER.

  • UPSERT : prend uniquement en charge le type upsert, qui inclut INSERT, DELETE et UPDATE_AFTER.

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 execution.checkpointing.checkpoints-after-tasks-finish.enabled sur true.

scan.incremental.snapshot.backfill.skip

Indique s'il faut ignorer la lecture des journaux pendant la phase de snapshot.

Boolean

Non

false

Valeurs valides :

  • true : ignore la lecture.

    Dans la phase incrémentielle, les journaux sont lus à partir du filigrane bas (low watermark).

    Si les opérateurs en aval ou le stockage prennent en charge l'idempotence, nous vous recommandons d'ignorer la lecture des journaux lors de la phase complète. Cela réduit le nombre de slots WAL, mais seule la sémantique au moins une fois (at-least-once) peut être garantie.

  • false : n'ignore pas la lecture.

    Lors de la lecture des splits dans la phase complète, les journaux situés entre le filigrane bas et le filigrane haut sont lus pour garantir la cohérence.

    Si la requête SQL effectue des agrégations, des jointures ou des opérations similaires, nous vous déconseillons d'ignorer la lecture des journaux lors de la phase complète.

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 postgres.

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 : bdb.schema_\..order_\..

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

pgoutput

Valeurs valides : decoderbufs et pgoutput.

tables.exclude

Les tables à exclure. Cette option prend effet après l'option tables. Les expressions régulières sont prises en charge.

Non

STRING

Voir l'option tables.

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 (ZoneId.systemDefault()) est utilisé.

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 :

  • initial (par défaut) : Analyse le snapshot au premier démarrage, puis bascule vers les dernières données WAL.

  • latest-offset : Ignore la lecture du snapshot ; commence la lecture à partir de la fin du WAL, ce qui signifie qu'il lit uniquement les dernières modifications effectuées après le démarrage du connecteur.

  • committed-offset : Ignore la lecture du snapshot ; il consomme les données WAL à partir d'un offset spécifié.

  • snapshot : Consomme uniquement le snapshot et non les données incrémentielles.

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 execution.checkpointing.checkpoints-after-tasks-finish.enabled sur true.

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 'jdbc.properties.useSSL' = 'false'.

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 : op_ts.

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.

Références