Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:AnalyticDB for PostgreSQL connector

Dernière mise à jour :Aug 09, 2026

Le connecteur AnalyticDB for PostgreSQL vous permet d'utiliser AnalyticDB for PostgreSQL en tant que table source (bêta), de dimension ou de destination dans les jobs SQL Realtime Compute for Apache Flink. AnalyticDB for PostgreSQL est un entrepôt de données à traitement massivement parallèle (MPP) destiné à l'analytique en ligne à grande échelle.

Prise en charge : Source (bêta) · Dimension · Destination | Streaming · Batch | API SQL | La table de destination prend en charge les mises à jour et les suppressions

La lecture depuis une table source nécessite un connecteur personnalisé configuré via Flink CDC. Pour obtenir des instructions de configuration, consultez la page Utilisation de Flink CDC pour s'abonner aux données complètes et incrémentielles en temps réel .

Prérequis

Avant de commencer, assurez-vous de disposer des éléments suivants :

Limitations

  • AnalyticDB for PostgreSQL V7.0 requiert VVR 8.0.1 ou une version ultérieure.

  • Les bases de données PostgreSQL auto-gérées ne sont pas prises en charge.

Syntaxe

Utilisez connector='adbpg' pour les tables de dimension et de destination, et connector='adbpg-cdc' pour les tables sources.

CREATE TEMPORARY TABLE adbpg_table (
  id      INT,
  len     INT,
  content VARCHAR,
  PRIMARY KEY (id)
) WITH (
  'connector'  = 'adbpg',
  'url'        = 'jdbc:postgresql://<host>:<port>/<database>',
  'tableName'  = '<table>',
  'userName'   = '<username>',
  'password'   = '<password>'
);

Options du connecteur

Général

Ces options s'appliquent à tous les types de tables.

Option Obligatoire Valeur par défaut Type Description
connector Oui STRING Définissez cette option sur adbpg-cdc pour les tables sources, ou sur adbpg pour les tables de dimension et de destination.
url Oui STRING L'URL Java Database Connectivity (JDBC) au format jdbc:postgresql://<host>:<port>/<database>.
tableName Oui STRING Le nom de la table dans la base de données.
userName Oui STRING Le nom d'utilisateur pour la base de données AnalyticDB for PostgreSQL.
password Oui STRING Le mot de passe pour la base de données AnalyticDB for PostgreSQL.
maxRetryTimes Non 3 INTEGER Le nombre maximal de tentatives en cas d'échec d'écriture.
targetSchema Non public STRING Le nom du schéma de base de données.
caseSensitive Non false STRING Indique si les identifiants sont sensibles à la casse. Valeurs valides : true, false.
connectionMaxActive Non 5 INTEGER Le nombre maximal de connexions dans le pool de connexions. Le connecteur libère automatiquement les connexions inactives. Une valeur trop élevée peut entraîner un nombre anormal de connexions serveur.

Spécifique à la source (bêta)

Ces options s'appliquent uniquement aux tables sources (connector='adbpg-cdc').

Option Obligatoire Valeur par défaut Type Description
schema-name Oui STRING Le nom du schéma. Prend en charge les expressions régulières pour s'abonner à plusieurs schémas simultanément.
port Oui 5432 INTEGER Le port de l'instance AnalyticDB for PostgreSQL.
decoding.plugin.name Oui pgoutput STRING Le plug-in de décodage logique PostgreSQL. Définissez cette option sur pgoutput.
slot.name Oui STRING Le nom du slot de décodage logique. Consultez la section Conseils sur le nom du slot.
debezium.* Oui STRING Propriétés de configuration du client Debezium. Par exemple, définissez 'debezium.snapshot.mode' = 'never' pour ignorer l'instantané initial. Consultez la page Propriétés du connecteur.
scan.incremental.snapshot.enabled Non false BOOLEAN Indique si l'instantané incrémentiel est activé. Valeurs valides : true, false.
scan.startup.mode Non initial STRING Le mode de démarrage pour la consommation des données. Consultez la section Modes de démarrage.
changelog-mode Non ALL STRING Méthode d'encodage des événements de modification dans le flux de modifications. Valeurs valides : ALL, UPSERT.
heartbeat.interval.ms Non 30 secondes DURATION L'intervalle d'envoi des paquets de heartbeat, en millisecondes. Consultez la section Heartbeat et rétention WAL.
scan.incremental.snapshot.chunk.key-column Non Première colonne de clé primaire STRING La colonne utilisée comme clé de fragment lors des lectures d'instantanés.

Modes de démarrage

Valeur Comportement
initial (par défaut) Effectue une analyse complète des données historiques au premier démarrage, puis lit à partir de la dernière position du journal des transactions préalables à l'écriture (WAL).
latest-offset Ignore l'analyse des données historiques et lit uniquement à partir de la fin actuelle du WAL.
snapshot Effectue une analyse complète des données historiques et capture les données WAL générées pendant cette analyse. S'arrête une fois l'analyse terminée.

Modes de journal des modifications

Valeur Événements capturés
ALL (par défaut) INSERT, DELETE, UPDATE_BEFORE, UPDATE_AFTER
UPSERT INSERT, DELETE, UPDATE_AFTER

Conseils sur le nom du slot

Gérez les slots de réplication avec soin pour éviter les conflits :

  • Au sein d'un seul job Flink : Utilisez le même slot.name pour toutes les tables sources du job.

  • Entre différents jobs Flink : Attribuez un slot.name unique par job. La réutilisation d'un nom de slot entre plusieurs jobs provoque l'erreur PSQLException: ERROR: replication slot "debezium" is active for PID 974.

Heartbeat et rétention WAL

Lorsqu'une table source reçoit peu de mises à jour, le connecteur peut ne pas faire avancer l'offset du slot de réplication. Cela empêche AnalyticDB for PostgreSQL de récupérer l'espace disque WAL, car la base de données conserve les fichiers WAL jusqu'à ce que l'offset du slot les dépasse. L'envoi périodique de heartbeats permet de faire avancer l'offset du slot, libérant ainsi les anciens fichiers WAL. Définissez heartbeat.interval.ms sur une valeur adaptée à votre fréquence de mise à jour lorsque les modifications de table sont rares.

Spécifique à la destination

Ces options s'appliquent uniquement aux tables de destination (connector='adbpg').

Option Obligatoire Valeur par défaut Type Description
writeMode Non insert STRING La stratégie d'écriture pour la première tentative d'écriture. Consultez la section Modes d'écriture et résolution des conflits.
conflictMode Non strict STRING Méthode de gestion des conflits de clé primaire ou d'index. Consultez la section Modes d'écriture et résolution des conflits.
batchSize Non 500 INTEGER Le nombre d'enregistrements écrits par lot.
flushIntervalMs Non INTEGER Le temps maximal d'attente du connecteur avant de vider les enregistrements mis en cache, en millisecondes. Lorsque le tampon n'atteint pas batchSize dans ce délai, tous les enregistrements tamponnés sont écrits immédiatement.
retryWaitTime Non 100 INTEGER L'intervalle entre les tentatives, en millisecondes.

Modes d'écriture et résolution des conflits

Les options writeMode et conflictMode fonctionnent conjointement pour contrôler l'écriture des enregistrements et la gestion des conflits.

**writeMode** Description Notes
insert (par défaut) Insère les enregistrements directement. La gestion des conflits est déterminée par conflictMode. Convient à la plupart des cas d'utilisation.
upsert Met automatiquement à jour les enregistrements existants en cas de conflit. Requiert une clé primaire.
copy Insère les enregistrements à l'aide de la commande PostgreSQL COPY. Requiert VVR 11.1 ou une version ultérieure.
**conflictMode** Comportement en cas de conflit Notes
strict (par défaut) Signale une erreur.
ignore Ignore silencieusement l'enregistrement conflictuel.
update Met à jour l'enregistrement conflictuel. Convient aux tables sans clé primaire. Réduit le débit d'écriture.
upsert Met à jour l'enregistrement conflictuel. Requiert une clé primaire. Plus efficace que update pour les tables indexées.

Spécifique à la table de dimension

Ces options s'appliquent uniquement aux tables de dimension (connector='adbpg' utilisé dans une jointure avec FOR SYSTEM_TIME AS OF).

Option Obligatoire Valeur par défaut Type Description
cache Non ALL STRING La politique de cache pour les recherches dans la table de dimension. Consultez la section Politiques de cache.
cacheSize Non 100000 LONG Le nombre maximal de lignes conservées dans le cache. Prend effet uniquement lorsque cache=LRU.
cacheTTLMs Non Long.MAX_VALUE LONG La durée de vie des entrées de cache en millisecondes. Le comportement dépend du paramètre cache. Consultez la section Politiques de cache.
maxJoinRows Non 1024 INTEGER Le nombre maximal de lignes à joindre par enregistrement d'entrée.

Politiques de cache

Le choix de la politique de cache appropriée permet d'équilibrer le débit des requêtes et l'actualité des données.

**Valeur cache** Comportement Cas d'utilisation
ALL (par défaut) Charge la totalité de la table de dimension en mémoire avant le démarrage du job. Toutes les recherches atteignent le cache. À l'expiration du cache (cacheTTLMs), recharge la table complète. Si une clé de jointure est introuvable dans le cache, elle n'existe pas. Tables de dimension de petite à moyenne taille où l'actualité des données peut tolérer des rechargements périodiques.
LRU Met en cache un sous-ensemble de lignes. En cas d'échec du cache, récupère les données depuis la base de données et met à jour le cache. Évince les entrées les moins récemment utilisées lorsque le cache atteint cacheSize. cacheTTLMs contrôle l'expiration par entrée. Grandes tables de dimension où la mise en cache intégrale n'est pas réalisable, ou lorsqu'une mise en cache partielle à faible latence est acceptable.
None Pas de cache. Chaque recherche accède directement à la base de données. Lorsque les données doivent toujours être lues fraîches depuis la base de données.

Mappages de types de données

Type AnalyticDB for PostgreSQL Type Flink SQL
BOOLEAN BOOLEAN
SMALLINT INT
INT INT
BIGINT BIGINT
FLOAT DOUBLE
VARCHAR VARCHAR
TEXT VARCHAR
TIMESTAMP TIMESTAMP
DATE DATE

Exemples

Table de destination

Cet exemple lit depuis une source datagen et écrit dans une table de destination AnalyticDB for PostgreSQL.

CREATE TEMPORARY TABLE datagen_source (
  `name` VARCHAR,
  `age`  INT
)
COMMENT 'datagen source table'
WITH (
  'connector' = 'datagen'
);

CREATE TEMPORARY TABLE adbpg_sink (
  name VARCHAR,
  age  INT
) WITH (
  'connector' = 'adbpg',
  'url'       = 'jdbc:postgresql://<host>:<port>/<database>',
  'tableName' = '<table>',
  'userName'  = '<username>',
  'password'  = '<password>'
);

INSERT INTO adbpg_sink
SELECT * FROM datagen_source;

Table de dimension

Cet exemple effectue une jointure entre un flux datagen et une table de dimension AnalyticDB for PostgreSQL à l'aide d'une jointure temporelle.

CREATE TEMPORARY TABLE datagen_source (
  a         INT,
  b         BIGINT,
  c         STRING,
  `proctime` AS PROCTIME()
)
COMMENT 'datagen source table'
WITH (
  'connector' = 'datagen'
);

CREATE TEMPORARY TABLE adbpg_dim (
  a INT,
  b VARCHAR,
  c VARCHAR
) WITH (
  'connector' = 'adbpg',
  'url'       = 'jdbc:postgresql://<host>:<port>/<database>',
  'tableName' = '<table>',
  'userName'  = '<username>',
  'password'  = '<password>'
);

CREATE TEMPORARY TABLE blackhole_sink (
  a INT,
  b STRING
)
COMMENT 'blackhole sink table'
WITH (
  'connector' = 'blackhole'
);

INSERT INTO blackhole_sink
SELECT T.a, H.b
FROM datagen_source AS T
JOIN adbpg_dim FOR SYSTEM_TIME AS OF T.proctime AS H
  ON T.a = H.a;

Table source (bêta)

Pour la configuration et des exemples de tables sources, consultez la page Utilisation de Flink CDC pour s'abonner aux données complètes et incrémentielles en temps réel.

Métriques

Les métriques des tables de destination sont disponibles dans la console Realtime Compute for Apache Flink. Les tables de dimension n'exposent pas de métriques. Pour les définitions des métriques, consultez la page Métriques.

Métrique Description
numRecordsOut Nombre total d'enregistrements écrits dans la destination.
numRecordsOutPerSecond Enregistrements écrits par seconde.
numBytesOut Total des octets écrits dans la destination.
numBytesOutPerSecond Octets écrits par seconde.
currentSendTime Temps nécessaire pour la dernière opération d'écriture.

Étapes suivantes