Le connecteur SelectDB intègre Realtime Compute for Apache Flink à ApsaraDB for SelectDB, un entrepôt de données en temps réel entièrement géré et compatible avec Apache Doris sur Alibaba Cloud. Utilisez-le pour créer des pipelines en temps réel qui lisent, écrivent ou consultent des données dans SelectDB, et pour exécuter une synchronisation complète de base de données dans les jobs d'ingestion de données basés sur YAML.
Fonctionnalités prises en charge :
| Catégorie | Détails |
|---|---|
| Types de table | Table source, table puits, table de dimension, puits d'ingestion de données |
| Mode d'exécution | Flux et lot |
| Format de données | JSON et CSV |
| Type d'API | DataStream, SQL et jobs YAML d'ingestion de données |
| Prise en charge Update/Delete | Oui |
| Métriques de surveillance | Aucune |
Principales fonctionnalités :
Synchronisation complète des données de la base de données
Sémantique exactly-once via validation en deux phases (2PC) — aucun enregistrement dupliqué ou perdu
Compatible avec Apache Doris 1.0 et versions ultérieures
Prérequis
Avant de commencer, assurez-vous de disposer des éléments suivants :
Realtime Compute for Apache Flink avec Ververica Runtime (VVR) version 8.0.10 ou ultérieure
Une instance ApsaraDB for SelectDB. Consultez Créer une instance.
Une liste d'autorisation IP configurée sur l'instance. Consultez Configurer une liste d'autorisation.
Configurer le connecteur
Le connecteur SelectDB est intégré à VVR 11.1 et aux versions ultérieures — aucune installation manuelle n'est requise.
Pour les versions VVR 8.0.10 à 11.0, installez le connecteur manuellement :
Téléchargez le package JAR depuis Maven Central (versions Flink 1.15–1.17).
Importez le fichier JAR dans votre console de développement Realtime Compute for Apache Flink. Consultez Gérer les connecteurs personnalisés.
Référencez le connecteur dans votre job SQL en utilisant
'connector' = 'doris'.
SQL
Syntaxe
Les trois types de table — source, puits et dimension — partagent la même syntaxe DDL. Spécifiez le rôle de la table via les paramètres inclus.
Pour utiliser SelectDB comme table source, activez d'abord la connexion directe au cluster. Dans la console ApsaraDB for SelectDB, accédez à Instance Details > Network Information et cliquez sur Enable Direct Cluster Connection . Cela active le protocole Arrow Flight SQL pour des lectures parallèles à haut débit.
CREATE TABLE selectdb_source (
order_id BIGINT,
user_id BIGINT,
total_amount DECIMAL(10, 2),
order_status TINYINT,
create_time TIMESTAMP(3),
product_name STRING
) WITH (
'connector' = 'doris',
'fenodes' = 'selectdb-cn-*******.selectdbfe.rds.aliyuncs.com:8080',
'table.identifier' = 'shop_db.orders',
'username' = 'admin',
'password' = '****'
);
Paramètres
Général
| Paramètre | Obligatoire | Par défaut | Description |
|---|---|---|---|
connector |
Oui | — | Défini sur doris. |
fenodes |
Oui | — | Endpoint HTTP de l'instance SelectDB : <Adresse VPC ou Adresse publique>:<Port du protocole HTTP>. Obtenez ces deux valeurs depuis Instance Details > Network Information dans la console SelectDB. Exemple : selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080. |
jdbc-url |
Non | — | Chaîne de connexion Java Database Connectivity (JDBC) pour les recherches dans les tables de dimension et les requêtes de métadonnées : jdbc:mysql://<Adresse VPC ou Adresse publique>:<Port du protocole MySQL>. Exemple : jdbc:mysql://selectdb-cn-***.selectdbfe.rds.aliyuncs.com:9030. |
table.identifier |
Oui | — | Table cible au format <base_de_données>.<table>. Exemple : db.tbl. |
username |
Oui | — | Nom d'utilisateur de la base de données. Réinitialisez le mot de passe depuis le coin supérieur droit de la page Instance Details si nécessaire. |
password |
Oui | — | Mot de passe associé au nom d'utilisateur de la base de données. |
doris.request.retries |
Non | 3 |
Nombre de tentatives pour les requêtes ayant échoué. |
doris.request.connect.timeout |
Non | 30s |
Délai de connexion. |
doris.request.read.timeout |
Non | 30s |
Délai de lecture. |
Table source
| Paramètre | Obligatoire | Par défaut | Description |
|---|---|---|---|
doris.request.query.timeout |
Non | 21600s |
Délai d'expiration de la requête (6 heures par défaut). |
doris.request.tablet.size |
Non | 1 |
Nombre de tablets par partition. Des valeurs plus faibles augmentent le parallélisme Flink mais exercent une pression accrue sur la base de données. |
doris.batch.size |
Non | 4064 |
Nombre maximal de lignes lues depuis un nœud Backend (BE) par requête. Augmentez cette valeur pour réduire la surcharge de connexion et la latence réseau. |
doris.exec.mem.limit |
Non | 8192mb |
Limite de mémoire par requête en octets (8 Go par défaut). |
source.use-flight-sql |
Non | false |
Aucune configuration requise — l'activation de Direct Cluster Connection dans la console SelectDB active automatiquement Arrow Flight SQL. |
source.flight-sql-port |
Non | — | Port Arrow Flight SQL (arrow_flight_sql_port) du nœud Frontend (FE). |
Table puits
Le mode d'écriture affecte les garanties de livraison et le comportement de vidage. Choisissez en fonction de vos exigences de cohérence :
| Écriture en flux | Écriture par lots | |
|---|---|---|
| Condition de déclenchement | Suit les intervalles de point de contrôle Flink | Vidage périodique par volume de données ou seuil temporel |
| Garantie de livraison | Exactly-once (via 2PC) | At-least-once ; atteignez l'idempotence avec le modèle Unique |
| Latence | Délimitée par l'intervalle de point de contrôle | Flexible, indépendante des points de contrôle |
| Tolérance aux pannes | Récupération complète de l'état Flink | S'appuie sur la déduplication du modèle Unique |
| Paramètre | Obligatoire | Par défaut | Description |
|---|---|---|---|
sink.label-prefix |
Non | — | Préfixe d'étiquette pour les importations Stream Load. Doit être globalement unique pour tous les jobs — la même étiquette ne peut être validée qu'une seule fois. Requis pour garantir la sémantique exactly-once lors des redémarrages de job. |
sink.properties.* |
Non | — | Paramètres d'importation Stream Load transmis directement à l'API Stream Load de SelectDB. Voir les exemples ci-dessous. |
sink.enable-delete |
Non | true |
Propager les opérations DELETE. Nécessite que la table Doris ait la suppression par lots activée et fonctionne uniquement avec le modèle Unique. |
sink.enable-2pc |
Non | true |
Activer la validation en deux phases (2PC) pour la sémantique exactly-once. Consultez Explicit Transaction Operations. |
sink.buffer-size |
Non | 1 MB |
Taille du tampon de cache d'écriture en octets. Laissez la valeur par défaut. |
sink.buffer-count |
Non | 3 |
Nombre de tampons de cache d'écriture. Laissez la valeur par défaut. |
sink.max-retries |
Non | 3 |
Nombre maximal de tentatives après un échec de validation. |
sink.enable.batch-mode |
Non | false |
Basculer vers le mode d'écriture par lots. Le vidage est contrôlé par les trois paramètres sink.buffer-flush.* ci-dessous plutôt que par les points de contrôle. La sémantique exactly-once n'est pas garantie ; utilisez le modèle Unique pour l'idempotence. |
sink.flush.queue-size |
Non | 2 |
Taille de la file d'attente du cache en mode lot. |
sink.buffer-flush.max-rows |
Non | 500000 |
Nombre maximal de lignes par vidage en mode lot. |
sink.buffer-flush.max-bytes |
Non | 100 MB |
Nombre maximal d'octets par vidage en mode lot. |
sink.buffer-flush.interval |
Non | 10s |
Intervalle de vidage en mode lot. |
sink.ignore.update-before |
Non | true |
Ignorer les événements update-before de Flink CDC. |
**Exemples pour sink.properties.* :**
Format CSV :
'sink.properties.column_separator' = ','
-- If values may contain commas, use a non-printable separator:
-- 'sink.properties.column_separator' = '\x01'
Format JSON :
'sink.properties.format' = 'json',
'sink.properties.read_json_by_line' = 'true'
-- Alternatively: 'sink.properties.strip_outer_array' = 'true'
Table de dimension
| Paramètre | Obligatoire | Par défaut | Description |
|---|---|---|---|
lookup.cache.max-rows |
Non | -1 |
Nombre maximal de lignes dans le cache de recherche. -1 désactive la mise en cache. |
lookup.cache.ttl |
Non | 10s |
Durée de vie (TTL) des entrées du cache. |
lookup.max-retries |
Non | 1 |
Tentatives après l'échec d'une requête de recherche. |
lookup.jdbc.async |
Non | false |
Activer la recherche asynchrone. |
lookup.jdbc.read.batch.size |
Non | 128 |
Taille maximale du lot par requête en mode de recherche asynchrone. |
lookup.jdbc.read.batch.queue-size |
Non | 256 |
Taille de la file d'attente du tampon intermédiaire en mode de recherche asynchrone. |
lookup.jdbc.read.thread-size |
Non | 3 |
Threads de recherche JDBC par tâche en mode de recherche asynchrone. |
Exemples
Table source
CREATE TEMPORARY TABLE selectdb_source (
order_id BIGINT,
user_id BIGINT,
total_amount DECIMAL(10, 2),
order_status TINYINT,
create_time TIMESTAMP(3),
product_name STRING
) WITH (
'connector' = 'doris',
'fenodes' = 'selectdb-cn-*******.selectdbfe.rds.aliyuncs.com:8080',
'table.identifier' = 'shop_db.orders',
'username' = 'admin',
'password' = '****'
);
Table puits
CREATE TEMPORARY TABLE selectdb_sink (
order_id BIGINT,
user_id BIGINT,
total_amount DECIMAL(10, 2),
order_status TINYINT,
create_time TIMESTAMP(3),
product_name STRING
) WITH (
'connector' = 'doris',
'fenodes' = 'selectdb-cn-*******.selectdbfe.rds.aliyuncs.com:8080',
'table.identifier' = 'shop_db.orders',
'username' = 'admin',
'password' = '****',
'sink.label-prefix' = 'flink_orders' -- Must be globally unique across jobs
);
Table de dimension
SelectDB agit comme une table de dimension de recherche jointe à une table de faits en flux.
-- Fact table from Kafka
CREATE TEMPORARY TABLE fact_table (
`id` BIGINT,
`name` STRING,
`city` STRING,
`process_time` AS proctime()
) WITH (
'connector' = 'kafka',
...
);
-- Dimension table from SelectDB
CREATE TEMPORARY TABLE dim_city (
`city` STRING,
`level` INT,
`province` STRING,
`country` STRING
) WITH (
'connector' = 'doris',
'fenodes' = 'selectdb-cn-*******.selectdbfe.rds.aliyuncs.com:8080',
'jdbc-url' = 'jdbc:mysql://selectdb-cn-***.selectdbfe.rds.aliyuncs.com:9030',
'table.identifier' = 'dim.dim_city',
'username' = 'admin',
'password' = '****'
);
-- Temporal join
SELECT a.id, a.name, a.city, c.province, c.country, c.level
FROM fact_table a
LEFT JOIN dim_city FOR SYSTEM_TIME AS OF a.process_time AS c
ON a.city = c.city;
Ingestion de données
Utilisez le connecteur SelectDB comme puits dans les jobs d'ingestion de données basés sur YAML pour la synchronisation complète de la base de données.
Syntaxe
source:
type: <source-type>
sink:
type: doris
name: Doris Sink
fenodes: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080
username: root
password: ""
Paramètres
| Paramètre | Obligatoire | Par défaut | Description |
|---|---|---|---|
type |
Oui | — | Défini sur doris. |
name |
Non | — | Nom descriptif pour le puits. |
fenodes |
Oui | — | Endpoint HTTP : <Adresse VPC ou Adresse publique>:<Port du protocole HTTP>. Obtenez ces deux valeurs depuis Instance Details > Network Information dans la console SelectDB. Exemple : selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080. |
jdbc-url |
Non | — | Chaîne de connexion JDBC. Exemple : jdbc:mysql://selectdb-cn-***.selectdbfe.rds.aliyuncs.com:9030. |
username |
Oui | — | Nom d'utilisateur de la base de données. |
password |
Oui | — | Mot de passe associé au nom d'utilisateur de la base de données. |
sink.enable.batch-mode |
Non | true |
Le mode lot est activé par défaut dans les jobs d'ingestion de données. Le vidage est contrôlé par les trois paramètres sink.buffer-flush.*. La sémantique exactly-once n'est pas garantie ; utilisez le modèle Unique pour l'idempotence. |
sink.flush.queue-size |
Non | 2 |
Taille de la file d'attente du cache. |
sink.buffer-flush.max-rows |
Non | 500000 |
Nombre maximal de lignes par vidage. |
sink.buffer-flush.max-bytes |
Non | 100 MB |
Nombre maximal d'octets par vidage. |
sink.buffer-flush.interval |
Non | 10s |
Intervalle de vidage. Minimum : 1s. |
sink.properties.* |
Non | — | Paramètres d'importation Stream Load. |
**Exemples pour sink.properties.* :**
Format CSV :
sink.properties.column_separator: ','
# If values may contain commas, use a non-printable separator:
# sink.properties.column_separator: '\x01'
Format JSON :
sink.properties.format: 'json'
sink.properties.read_json_by_line: 'true'
Mappage de types
De Flink vers SelectDB
| Type Flink CDC | Type SelectDB | Notes |
|---|---|---|
TINYINT |
TINYINT |
|
SMALLINT |
SMALLINT |
|
INT |
INT |
|
BIGINT |
BIGINT |
|
DECIMAL |
DECIMAL |
|
FLOAT |
FLOAT |
|
DOUBLE |
DOUBLE |
|
BOOLEAN |
BOOLEAN |
|
DATE |
DATE |
|
TIMESTAMP[(p)] |
DATETIME[(p)] |
|
TIMESTAMP_LTZ[(p)] |
DATETIME[(p)] |
|
CHAR(n) |
CHAR(n*3) |
SelectDB stocke les chaînes en UTF-8. Les caractères anglais occupent 1 octet ; les caractères chinois occupent 3 octets. La longueur maximale de CHAR est 255 ; les valeurs plus longues sont automatiquement converties en VARCHAR. |
VARCHAR(n) |
VARCHAR(n*3) |
Le même multiplicateur UTF-8 s'applique. La longueur maximale de VARCHAR est 65533 ; les valeurs plus longues sont automatiquement converties en STRING. |
BINARY(n) |
STRING |
|
VARBINARY(n) |
STRING |
|
STRING |
STRING |