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 :
Une instance et une table AnalyticDB for PostgreSQL. Consultez les pages Création d'une instance et CREATE TABLE
Une liste d'autorisation d'adresses IP configurée pour l'instance. Consultez la page Configuration d'une liste d'autorisation d'adresses IP
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.namepour toutes les tables sources du job.Entre différents jobs Flink : Attribuez un
slot.nameunique par job. La réutilisation d'un nom de slot entre plusieurs jobs provoque l'erreurPSQLException: 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. |