Le connecteur ApsaraDB RDS for MySQL ne sera plus pris en charge à l'avenir. Utilisez le connecteur MySQL à la place.
Le connecteur ApsaraDB RDS for MySQL vous permet d'écrire la sortie Flink SQL dans une table de destination (sink) ApsaraDB RDS for MySQL ou de joindre un flux à une table de dimension ApsaraDB RDS for MySQL.
Types de tables pris en charge : Table de destination · Table de dimension
Supported running modes: Mode par lots · Mode streaming
Type d'API : SQL
Mises à jour et suppressions de données dans les tables de destination : Pris en charge
Prérequis
Avant de commencer, assurez-vous de disposer des éléments suivants :
Une base de données et une table ApsaraDB RDS for MySQL. Consultez la rubrique Créer des bases de données et des comptes pour une instance ApsaraDB RDS for MySQL
Une liste d'autorisation d'adresses IP configurée pour la base de données. Consultez la rubrique Se connecter à une instance ApsaraDB RDS for MySQL à l'aide d'un client de base de données ou de la CLI
Limites
Nécessite Realtime Compute for Apache Flink avec Ververica Runtime (VVR) 2.0.0 ou version ultérieure. Pour des performances et une stabilité optimales, utilisez VVR 6.X ou version ultérieure.
Seules les bases de données ApsaraDB RDS for MySQL sont prises en charge.
Le connecteur utilise une sémantique « au moins une fois » (at-least-once). Si la table de destination possède une clé primaire, l'idempotence garantit l'exactitude des données.
Fonctionnement
Comportement d'écriture dans la table de destination
Chaque ligne de sortie est convertie en instruction SQL avant d'être écrite dans la table de destination :
Sans clé primaire — exécute
INSERT INTO table_name (col1, col2, ...) VALUES (val1, val2, ...);Avec clé primaire — exécute
INSERT INTO table_name (col1, col2, ...) VALUES (val1, val2, ...) ON DUPLICATE KEY UPDATE col1 = VALUES(col1), col2 = VALUES(col2), ...;
Conflits d'index unique : Si la table physique possède une contrainte d'index unique en plus de la clé primaire, l'insertion de deux lignes ayant des clés primaires différentes mais la même valeur d'index unique entraîne l'écrasement de la première ligne, ce qui provoque une perte de données.
Clés primaires à incrémentation automatique : Ne déclarez pas les champs à incrémentation automatique dans le DDL Flink. La base de données attribue ces valeurs automatiquement. Le connecteur peut écrire et supprimer des lignes contenant des champs à incrémentation automatique, mais il ne peut pas les mettre à jour.
Politiques de mise en cache des tables de dimension
Le connecteur prend en charge trois politiques de mise en cache pour les recherches dans les tables de dimension :
| Politique | Comportement | Cas d'utilisation |
|---|---|---|
NONE |
Aucune mise en cache — chaque recherche interroge directement la base de données | Exigences de faible latence, petits ensembles de données |
LRU |
Met en cache un nombre fixe de lignes récemment utilisées par gestionnaire de tâches | Sous-ensembles fréquemment consultés de grandes tables |
ALL |
Charge la table entière en mémoire et la recharge périodiquement | Petites tables de référence statiques |
Syntaxe
Table de destination
CREATE TABLE rds_sink (
id INT,
num BIGINT,
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'rds',
'tableName' = '<your-table-name>',
'userName' = '<your-user-name>',
'password' = '<your-password>',
'url' = 'jdbc:mysql://<internal-endpoint>:<port>/<database-name>?rewriteBatchedStatements=true'
);
Ajoutez ?rewriteBatchedStatements=true à la valeur url pour les tables de destination afin d'améliorer le débit d'écriture.
Table de dimension
CREATE TABLE rds_dim (
id1 INT,
id2 VARCHAR
) WITH (
'connector' = 'rds',
'tableName' = '<your-table-name>',
'userName' = '<your-user-name>',
'password' = '<your-password>',
'url' = 'jdbc:mysql://<internal-endpoint>:<port>/<database-name>',
'cache' = 'NONE'
);
Paramètres de la clause WITH
Paramètres communs
| Paramètre | Type | Obligatoire | Par défaut | Description |
|---|---|---|---|---|
connector |
STRING | Oui | — | Définissez sur rds |
tableName |
STRING | Oui | — | Nom de la table physique dans ApsaraDB RDS for MySQL |
userName |
STRING | Oui | — | Nom d'utilisateur de la base de données |
password |
STRING | Oui | — | Mot de passe de la base de données |
url |
STRING | Oui | — | Endpoint VPC (Virtual Private Cloud) de la base de données, au format jdbc:mysql://<internal-endpoint>:<port>/<database-name>. Pour les tables de destination, ajoutez ?rewriteBatchedStatements=true. Pour plus de détails sur les endpoints, consultez la rubrique Afficher et modifier les endpoints internes et publics ainsi que les numéros de port d'une instance ApsaraDB RDS for MySQL |
maxRetryTimes |
INTEGER | Non | 10 (VVR 4.0.7+), 3 (VVR 4.0.6 et antérieures) | Nombre maximal de tentatives pour les recherches dans les tables de dimension ou les écritures dans les tables de destination ayant échoué |
Paramètres de la table de destination
| Paramètre | Type | Obligatoire | Par défaut | Description |
|---|---|---|---|---|
batchSize |
INTEGER | Non | 4096 (VVR 4.0.7+), 5000 (VVR 4.0.0–4.0.6), 100 (VVR 3.x et antérieures) | Nombre de lignes écrites par lot |
bufferSize |
INTEGER | Non | 10000 | Nombre maximal de lignes mises en cache en mémoire avant le déclenchement d'une écriture. Pris en charge dans VVR 4.0.7 et versions ultérieures. Prend effet uniquement lorsqu'une clé primaire est définie |
flushIntervalMs |
INTEGER | Non | 2000 (VVR 4.0.7+), 0 (VVR 4.0.0–4.0.6), 1000 (VVR 3.x et antérieures) | Intervalle en millisecondes auquel le buffer est vidé vers la table de destination, indépendamment du fait que les seuils batchSize ou bufferSize soient atteints. Si la valeur est définie sur 0 (valeur par défaut pour VVR 4.0.0–4.0.6), de petites quantités de données mises en buffer peuvent ne jamais être écrites ; effectuez une mise à niveau vers une version ultérieure de VVR pour éviter ce problème |
ignoreDelete |
BOOLEAN | Non | false | Définissez sur true pour ignorer les opérations de suppression. Utile lorsque plusieurs opérateurs mettent à jour différents champs de la même ligne ; sans ce paramètre, une suppression dans un opérateur suivie d'une mise à jour partielle dans un autre laisse les champs non mis à jour avec la valeur null ou leurs valeurs par défaut |
connectionMaxActive |
INTEGER | Non | 40 | Taille du pool de connexions. Pris en charge dans VVR 4.0.7 et versions ultérieures. Augmentez cette valeur si des délais d'expiration du pool de connexions se produisent ; diminuez-la si la base de données limite le nombre de connexions simultanées |
Paramètres de la table de dimension
| Paramètre | Type | Obligatoire | Par défaut | Description |
|---|---|---|---|---|
cache |
STRING | Non | NONE (VVR antérieur à 4.0.6), ALL (VVR 4.0.6+) | Politique de mise en cache. Valeurs valides : NONE, LRU, ALL. Consultez la rubrique Politiques de mise en cache |
cacheSize |
INTEGER | Non | 100000 | Nombre maximal de lignes à mettre en cache. Obligatoire lorsque cache est défini sur LRU ; ignoré pour NONE et ALL |
cacheTTLMs |
LONG | Non | Pas d'expiration pour NONE et LRU ; pas de rechargement pour ALL |
Durée de vie du cache en millisecondes. Pour LRU, les lignes expirent après cette période. Pour ALL, l'intégralité du cache est rechargée à cet intervalle |
maxJoinRows |
INTEGER | Non | 1024 | Nombre maximal de lignes de la table de dimension correspondant à chaque ligne d'entrée. Définissez cette valeur sur le nombre maximal de lignes de dimension attendues par ligne de la table principale pour éviter une analyse inutile |
Métriques
La table de destination expose les métriques suivantes. Les tables de dimension n'ont aucune métrique.
| Métrique | Description |
|---|---|
numRecordsOut |
Total des lignes écrites |
numRecordsOutPerSecond |
Lignes écrites par seconde |
numBytesOut |
Total des octets écrits |
numBytesOutPerSecond |
Octets écrits par seconde |
currentSendTime |
Latence d'écriture actuelle |
numRecordsOutErrors |
Total des erreurs d'écriture |
Pour les définitions des métriques, consultez la rubrique Métriques.
Mappages de types de données
| Type Flink | Type ApsaraDB RDS for MySQL |
|---|---|
| BOOLEAN | BOOLEAN |
| TINYINT | TINYINT |
| TINYINT(1) (tables de dimension uniquement) | BOOLEAN |
| SMALLINT | SMALLINT |
| SMALLINT | TINYINT UNSIGNED |
| INT | INT |
| INT | SMALLINT UNSIGNED |
| BIGINT | BIGINT |
| BIGINT | INT UNSIGNED |
| DECIMAL(20, 0) | BIGINT UNSIGNED |
| FLOAT | FLOAT |
| DECIMAL | DECIMAL |
| DOUBLE | DOUBLE |
| DATE | DATE |
| TIME | TIME |
| TIMESTAMP | TIMESTAMP |
| VARCHAR | VARCHAR |
| VARBINARY | VARBINARY |
Exemples
Exemple de table de destination
L'exemple suivant lit à partir d'une source DataGen et écrit dans une table de destination ApsaraDB RDS for MySQL.
CREATE TEMPORARY TABLE datagen_source (
`name` VARCHAR,
`age` INT
) WITH (
'connector' = 'datagen'
);
CREATE TEMPORARY TABLE rds_sink (
`name` VARCHAR,
`age` INT
) WITH (
'connector' = 'rds',
'tableName' = '<your-table-name>',
'userName' = '<your-user-name>',
'password' = '<your-password>',
'url' = 'jdbc:mysql://<internal-endpoint>:<port>/<database-name>?rewriteBatchedStatements=true'
);
INSERT INTO rds_sink
SELECT * FROM datagen_source;
Exemple de table de dimension
L'exemple suivant joint un flux à une table de dimension ApsaraDB RDS for MySQL à l'aide d'une jointure temporelle.
CREATE TEMPORARY TABLE datagen_source (
a INT,
b BIGINT,
c STRING,
`proctime` AS PROCTIME()
) WITH (
'connector' = 'datagen'
);
CREATE TEMPORARY TABLE rds_dim (
a INT,
b VARCHAR,
c VARCHAR
) WITH (
'connector' = 'rds',
'tableName' = '<your-table-name>',
'userName' = '<your-user-name>',
'password' = '<your-password>',
'url' = 'jdbc:mysql://<internal-endpoint>:<port>/<database-name>'
);
CREATE TEMPORARY TABLE blackhole_sink (
a INT,
b STRING
) WITH (
'connector' = 'blackhole'
);
INSERT INTO blackhole_sink
SELECT T.a, H.b
FROM datagen_source AS T
JOIN rds_dim FOR SYSTEM_TIME AS OF T.`proctime` AS H
ON T.a = H.a;
FAQ
Étapes suivantes
Connecteur MySQL — remplacement recommandé pour ce connecteur
ApsaraDB RDS for MySQL — présentation du produit et documentation des fonctionnalités
Métriques — définitions de toutes les métriques du connecteur