Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:SelectDB

Dernière mise à jour :Aug 09, 2026

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 :

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 :

  1. Téléchargez le package JAR depuis Maven Central (versions Flink 1.15–1.17).

  2. Importez le fichier JAR dans votre console de développement Realtime Compute for Apache Flink. Consultez Gérer les connecteurs personnalisés.

  3. 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

Étapes suivantes