Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:ApsaraDB RDS for MySQL

Dernière mise à jour :Aug 20, 2026
Important

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 :

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'
);
Remarque

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