Le connecteur JDBC lit et écrit dans des bases de données relationnelles (MySQL, PostgreSQL et Oracle) en utilisant le langage DDL SQL standard au sein des jobs Flink SQL.
Types de table pris en charge : table source · table de dimension · table de sortie (sink)
Modes d'exécution pris en charge : mode streaming · mode batch · API SQL
Prérequis
Avant de commencer, assurez-vous que :
La base de données et la table cibles existent déjà.
Le fichier JAR du pilote JDBC correspondant à votre base de données est disponible pour téléchargement.
Limites
Lectures bornées : Une table source JDBC constitue une source bornée. La tâche se termine une fois toutes les lignes lues. Pour capturer les données de modification en temps réel, utilisez plutôt un connecteur Change Data Capture (CDC). Consultez les rubriques Créer une table source MySQL CDC et Créer une table source PostgreSQL CDC (aperçu public).
Version de PostgreSQL : L'écriture dans PostgreSQL nécessite la version 9.5 ou ultérieure, car le mécanisme de sortie (sink) s'appuie sur la clause
ON CONFLICT.-
Téléchargement du pilote JDBC : Téléchargez le fichier JAR du pilote JDBC en tant que dépendance avant d'exécuter votre job. Pilotes courants :
Base de données Group ID Artifact ID MySQL mysqlmysql-connector-java Oracle com.oracle.database.jdbcojdbc8 PostgreSQL org.postgresqlpostgresql Pour les pilotes non répertoriés ici, vérifiez la compatibilité avant utilisation.
-
Comportement upsert MySQL : Lors de l'écriture dans une table de sortie MySQL dotée d'une clé primaire, le connecteur émet des instructions
INSERT INTO ... ON DUPLICATE KEY UPDATE ....AvertissementL'insertion de lignes comportant des valeurs d'index unique dupliquées, même si les clés primaires diffèrent, entraîne l'écrasement des lignes existantes dans toute table physique soumise à une contrainte d'index unique, ce qui provoque une perte de données.
Créer une table JDBC
CREATE TABLE jdbc_table (
`id` BIGINT,
`name` VARCHAR,
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:<db-type>://<host>:<port>/<database>',
'table-name' = '<yourTable>',
'username' = '<yourUsername>',
'password' = '<yourPassword>'
);
Options du connecteur
Général
| Option | Type | Obligatoire | Par défaut | Description |
|---|---|---|---|---|
connector |
STRING | Oui | — | Définissez cette option sur jdbc. |
url |
STRING | Oui | — | URL JDBC de la base de données. |
table-name |
STRING | Oui | — | Nom de la table à lire ou à écrire. |
username |
STRING | Non | — | Nom d'utilisateur de la base de données. À définir conjointement avec password. |
password |
STRING | Non | — | Mot de passe de la base de données. |
Options de source
| Option | Type | Obligatoire | Par défaut | Description |
|---|---|---|---|---|
scan.partition.column |
STRING | Non | — | Colonne utilisée pour diviser les données en partitions. Doit être de type NUMERIC ou TIMESTAMP. Consultez la section Analyse partitionnée. |
scan.partition.num |
INTEGER | Non | — | Nombre de partitions. |
scan.partition.lower-bound |
LONG | Non | — | Valeur minimale de la première partition. |
scan.partition.upper-bound |
LONG | Non | — | Valeur maximale de la dernière partition. |
scan.fetch-size |
INTEGER | Non | 0 |
Nombre de lignes récupérées par aller-retour vers la base de données. Si cette option est définie sur 0, elle est ignorée. |
scan.auto-commit |
BOOLEAN | Non | true |
Active commit automatique pour les transactions de lecture. |
Options de sortie (sink)
| Option | Type | Obligatoire | Par défaut | Description |
|---|---|---|---|---|
sink.buffer-flush.max-rows |
INTEGER | Non | 100 |
Nombre maximal de lignes mises en mémoire tampon avant vidage. Définissez cette option sur 0 pour vider chaque ligne immédiatement. |
sink.buffer-flush.interval |
DURATION | Non | 1000 ms |
Durée maximale de conservation des lignes en mémoire tampon avant vidage. Définissez cette option sur 0 pour vider chaque ligne immédiatement. |
sink.max-retries |
INTEGER | Non | 3 |
Nombre maximal de tentatives d'écriture en cas d'échec. |
sink.ignore-delete |
BOOLEAN | Non | false |
Ignore les messages de suppression au lieu de les transmettre. Nécessite VVR 11.4 ou version ultérieure. |
sink.ignore-delete-mode |
STRING | Non | ALL |
Contrôle les messages de suppression ignorés lorsque sink.ignore-delete est défini sur true. Nécessite VVR 11.4 ou version ultérieure. Valeurs valides : ALL (ignore -D et -U), REAL_DELETE (ignore uniquement -D), UPDATE_BEFORE (ignore uniquement -U). |
Pour vider les lignes mises en mémoire tampon de manière asynchrone selon un minuteur plutôt que sur la base du nombre de lignes, définissez sink.buffer-flush.max-rows sur 0 et configurez sink.buffer-flush.interval selon l'intervalle de vidage souhaité.
Options de table de dimension
| Option | Type | Obligatoire | Par défaut | Description |
|---|---|---|---|---|
lookup.cache.max-rows |
INTEGER | Non | — | Nombre maximal de lignes dans le cache de recherche. Lorsque le cache est plein, la ligne la moins récemment utilisée expire. La mise en cache est désactivée sauf si les options lookup.cache.max-rows et lookup.cache.ttl sont toutes deux définies. |
lookup.cache.ttl |
DURATION | Non | — | Durée maximale de validité d'une ligne mise en cache avant son expiration. |
lookup.cache.caching-missing-key |
BOOLEAN | Non | true |
Met en cache les résultats de recherche vides afin d'éviter des interrogations répétées de la base de données pour les clés manquantes. |
lookup.max-retries |
INTEGER | Non | 3 |
Nombre maximal de tentatives en cas d'échec d'une requête vers la base de données. |
Options PostgreSQL
| Option | Type | Obligatoire | Par défaut | Description |
|---|---|---|---|---|
source.extend-type.enabled |
BOOLEAN | Non | false |
Lorsque cette option est définie sur true, elle mappe les colonnes PostgreSQL JSONB et UUID vers le type Flink STRING. Pour les jointures de recherche sur des colonnes UUID, ajoutez également stringtype=unspecified à l'URL JDBC afin que PostgreSQL effectue les requêtes selon le type réel plutôt que par conversion. |
Comportements clés
Analyse partitionnée
Pour activer les lectures parallèles depuis une table source, configurez conjointement les options scan.partition.column, scan.partition.lower-bound, scan.partition.upper-bound et scan.partition.num. La colonne de partition doit être de type numérique ou TIMESTAMP. Pour plus de détails sur le calcul des splits, consultez la documentation Apache Flink relative à l'analyse partitionnée.
Cache de recherche
Par défaut, chaque requête de jointure de recherche interroge directement la base de données. Activez le cache LRU en définissant simultanément lookup.cache.max-rows et lookup.cache.ttl afin de réduire la charge sur la base de données, au prix d'une légère obsolescence des données.
Écritures idempotentes
Lorsque la table de sortie possède une clé primaire, le connecteur émet des instructions upsert spécifiques à la base de données. Pour MySQL, le connecteur utilise l'instruction INSERT ... ON DUPLICATE KEY UPDATE ....
Exemples
Les trois exemples suivants utilisent un connecteur blackhole ou datagen comme contrepartie légère, permettant ainsi aux instructions d'être autonomes.
Lire depuis une base de données (table source)
CREATE TEMPORARY TABLE jdbc_source (
`id` INT,
`name` VARCHAR
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://localhost:3306/mydb',
'table-name' = '<yourTable>',
'username' = '<yourUsername>',
'password' = '<yourPassword>'
);
CREATE TEMPORARY TABLE blackhole_sink (
`id` INT,
`name` VARCHAR
) WITH (
'connector' = 'blackhole'
);
INSERT INTO blackhole_sink SELECT * FROM jdbc_source;
Écrire dans une base de données (table de sortie)
CREATE TEMPORARY TABLE datagen_source (
`name` VARCHAR,
`age` INT
) WITH (
'connector' = 'datagen'
);
CREATE TEMPORARY TABLE jdbc_sink (
`name` VARCHAR,
`age` INT
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://localhost:3306/mydb',
'table-name' = '<yourTable>',
'username' = '<yourUsername>',
'password' = '<yourPassword>'
);
INSERT INTO jdbc_sink SELECT * FROM datagen_source;
Enrichir un flux via des recherches en base de données (table de dimension)
CREATE TEMPORARY TABLE datagen_source (
`id` INT,
`data` BIGINT,
`proctime` AS PROCTIME()
) WITH (
'connector' = 'datagen'
);
CREATE TEMPORARY TABLE jdbc_dim (
`id` INT,
`name` VARCHAR
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://localhost:3306/mydb',
'table-name' = '<yourTable>',
'username' = '<yourUsername>',
'password' = '<yourPassword>'
);
CREATE TEMPORARY TABLE blackhole_sink (
`id` INT,
`data` BIGINT,
`name` VARCHAR
) WITH (
'connector' = 'blackhole'
);
-- Look up the matching name for each stream record at processing time
INSERT INTO blackhole_sink
SELECT T.`id`, T.`data`, H.`name`
FROM datagen_source AS T
JOIN jdbc_dim FOR SYSTEM_TIME AS OF T.proctime AS H
ON T.id = H.id;
Mappings de types de données
| MySQL | Oracle | PostgreSQL | Flink SQL |
|---|---|---|---|
TINYINT |
— | — | TINYINT |
SMALLINT, TINYINT UNSIGNED |
— | SMALLINT, INT2, SMALLSERIAL, SERIAL2 |
SMALLINT |
INT, MEDIUMINT, SMALLINT UNSIGNED |
— | INTEGER, SERIAL |
INT |
BIGINT, INT UNSIGNED |
— | BIGINT, BIGSERIAL |
BIGINT |
BIGINT UNSIGNED |
— | — | DECIMAL(20, 0) |
FLOAT |
BINARY_FLOAT |
REAL, FLOAT4 |
FLOAT |
DOUBLE, DOUBLE PRECISION |
BINARY_DOUBLE |
FLOAT8, DOUBLE PRECISION |
DOUBLE |
NUMERIC(p, s), DECIMAL(p, s) |
SMALLINT, FLOAT(s), DOUBLE PRECISION, REAL, NUMBER(p, s) |
NUMERIC(p, s), DECIMAL(p, s) |
DECIMAL(p, s) |
BOOLEAN, TINYINT(1) |
— | BOOLEAN |
BOOLEAN |
DATE |
DATE |
DATE |
DATE |
TIME [(p)] |
DATE |
TIME [(p)] [WITHOUT TIMEZONE] |
TIME [(p)] [WITHOUT TIMEZONE] |
DATETIME [(p)] |
TIMESTAMP [(p)] [WITHOUT TIMEZONE] |
TIMESTAMP [(p)] [WITHOUT TIMEZONE] |
TIMESTAMP [(p)] [WITHOUT TIMEZONE] |
CHAR(n), VARCHAR(n), TEXT |
CHAR(n), VARCHAR(n), CLOB |
CHAR(n), CHARACTER(n), VARCHAR(n), CHARACTER VARYING(n), TEXT, JSONB, UUID |
STRING |
BINARY, VARBINARY, BLOB |
RAW(s), BLOB |
BYTEA |
BYTES |
| — | — | ARRAY |
ARRAY |
Étapes suivantes
Créer une table source MySQL CDC : capturez les données de modification en temps réel depuis MySQL.
Créer une table source PostgreSQL CDC (aperçu public) : capturez les données de modification en temps réel depuis PostgreSQL.