Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:StarRocks

Dernière mise à jour :Aug 19, 2026

Découvrez comment utiliser le connecteur StarRocks.

Contexte

StarRocks est un entrepôt de données MPP (Massively Parallel Processing) de nouvelle génération qui offre des performances extrêmement rapides dans tous les scénarios, ainsi qu'une expérience analytique unifiée. StarRocks présente les avantages suivants :

  • StarRocks est compatible avec le protocole MySQL, ce qui vous permet d'utiliser des clients MySQL et des outils courants de Business Intelligence (BI) pour vous connecter et analyser les données.

  • StarRocks repose sur une architecture distribuée :

    • Il partitionne horizontalement les tables de données et les stocke avec plusieurs réplicas.

    • Le cluster peut être mis à l'échelle de manière flexible et peut analyser jusqu'à 10 pétaoctets (Po) de données.

    • Il utilise un framework MPP pour accélérer les calculs parallèles.

    • Il prend en charge plusieurs réplicas afin d'assurer la tolérance aux pannes.

Le connecteur Flink met les données en cache et utilise Stream Load pour les écrire par lots dans les tables de destination (sink). Il lit les tables source en récupérant les données par lots. Le tableau suivant répertorie les fonctionnalités du connecteur StarRocks.

Catégorie

Description

Types pris en charge

tables source, tables de dimension, tables de destination (sink) et cibles d'ingestion de données

Mode d'exécution

mode streaming et mode batch

Format de données

CSV

Métriques spécifiques au connecteur

Aucune

Types d'API

DataStream, SQL et YAML pour l'ingestion de données

Prise en charge des mises à jour/suppressions dans les tables de destination

Oui

Prérequis

Vous devez disposer d'un cluster StarRocks déployé sur EMR ou d'un cluster autogéré sur ECS.

Limites

  • Seules les versions Ververica Runtime (VVR) 11,1 ou ultérieures prennent en charge les jointures avec les tables de dimension.

  • Pour éviter les restrictions d'accès réseau, ajoutez les ports suivants du cluster StarRocks à un groupe de sécurité ou à une liste d'autorisation de pare-feu : 9030, 8030, 8040, 9060, 8060, 9020.

SQL

Fonctionnalités

StarRocks sur E-MapReduce prend en charge les instructions CREATE TABLE AS SELECT (CTAS) et CREATE DATABASE AS SELECT (CDAS). L'instruction CTAS synchronise le schéma et les données d'une seule table, tandis que CDAS synchronise une base de données entière ou plusieurs tables au sein de la même base de données. Pour plus d'informations, consultez la rubrique Utilisation des instructions CTAS et CDAS dans Realtime Compute for Apache Flink pour synchroniser les données d'une base de données MySQL vers StarRocks.

Syntaxe

CREATE TABLE USER_RESULT(
 name VARCHAR,
 score BIGINT
 ) WITH (
 'connector' = 'starrocks',
 'jdbc-url'='jdbc:mysql://fe1_ip:query_port,fe2_ip:query_port,fe3_ip:query_port?xxxxx',
 'load-url'='fe1_ip:http_port;fe2_ip:http_port;fe3_ip:http_port',
 'database-name' = 'xxx',
 'table-name' = 'xxx',
 'username' = 'xxx',
 'password' = 'xxx'
 );

Paramètres

Type

Paramètre

Description

Type

Obligatoire

Valeur par défaut

Remarques

Général

connector

Spécifie le connecteur à utiliser.

String

Oui

La valeur doit être starrocks.

jdbc-url

L'URL JDBC (Java Database Connectivity).

String

Oui

Spécifiez l'adresse IP et le port JDBC du nœud FE au format jdbc:mysql://ip:port.

database-name

Le nom de la base de données StarRocks.

String

Oui

table-name

Le nom de la table StarRocks.

String

Oui

username

Le nom d'utilisateur pour se connecter à StarRocks.

String

Oui

password

Le mot de passe pour se connecter à StarRocks.

String

Oui

starrocks.create.table.properties

Spécifie les propriétés pour la création automatique de table.

String

Non

Spécifie les propriétés initiales de la table, telles que le moteur et le nombre de réplicas. Par exemple : 'starrocks.create.table.properties' = 'buckets 8' ou 'starrocks.create.table.properties' = 'replication_num=1'.

Spécifique à la source

scan-url

L'URL de balayage des données.

String

Non

Spécifie l'adresse IP et le port HTTP du nœud FE. Format : fe_ip:http_port;fe_ip:http_port.

Remarque

Pour spécifier plusieurs adresses IP et ports, séparez-les par des points-virgules (;).

scan.connect.timeout-ms

Le délai d'expiration pour que flink-connector-starrocks se connecte à StarRocks.

Le connecteur signale une erreur si une connexion n'est pas établie dans ce délai.

String

Non

1000

Unité : millisecondes.

scan.params.keep-alive-min

La durée de maintien de la connexion (keep-alive) pour la tâche de requête.

String

Non

10

scan.params.query-timeout-s

Le délai d'expiration pour une tâche de requête.

Si aucun résultat n'est renvoyé dans ce délai, le système arrête la tâche de requête.

String

Non

600

Unité : secondes.

scan.params.mem-limit-byte

La limite de mémoire pour une seule requête sur un nœud BE.

String

Non

1073741824 (1 Go)

Unité : octets.

scan.max-retries

Le nombre maximal de tentatives pour une requête ayant échoué.

Le connecteur signale une erreur si cette limite est dépassée.

String

Non

1

Spécifique à la destination (sink)

load-url

L'URL d'importation des données.

String

Oui

Spécifiez les adresses IP et les ports HTTP des nœuds FE au format fe_ip:http_port;fe_ip:http_port.

Remarque

Pour spécifier plusieurs adresses IP et ports, séparez-les par des points-virgules (;).

sink.semantic

La sémantique de livraison pour les écritures.

String

Non

at-least-once

Valeurs valides :

  • at-least-once (par défaut) : Garantit que les données sont livrées au moins une fois.

  • exactly-once : Garantit que les données sont livrées exactement une fois.

sink.buffer-flush.max-bytes

La quantité maximale de données à mettre en tampon avant vidage.

String

Non

94371840 (90 Mo)

Plage valide : 64 Mo à 10 Go.

sink.buffer-flush.max-rows

Le nombre maximal de lignes à mettre en tampon avant vidage.

String

Non

500000

Plage valide : 64 000 à 5 000 000.

sink.buffer-flush.interval-ms

L'intervalle de vidage du tampon.

String

Non

300000

Plage valide : 1 000 ms à 3 600 000 ms.

sink.max-retries

Le nombre maximal de tentatives pour les écritures ayant échoué.

String

Non

3

Plage valide : 0 à 10.

sink.connect.timeout-ms

Le délai d'expiration pour la connexion à StarRocks.

String

Non

1000

Plage valide : 100 à 60 000. Unité : millisecondes.

sink.properties.*

Propriétés Stream Load supplémentaires pour la destination (sink).

String

Non

Ces paramètres contrôlent le comportement de Stream Load. Par exemple, sink.properties.format spécifie le format des données importées, tel que CSV. Pour plus de paramètres, consultez la rubrique Stream Load.

Spécifique à la dimension

lookup.cache.enabled

Indique s'il faut activer la mise en cache pour la table de dimension.

Boolean

Non

true

Valeurs valides :

  • true : Active la mise en cache. Après la première lecture des données de la table, le système les met en cache en mémoire. Les requêtes suivantes utilisent les données mises en cache pendant leur période de validité afin de réduire la charge d'E/S.

  • false : Désactive la mise en cache. Chaque requête accède directement à la source de données.

Important
  • Cette fonctionnalité nécessite le moteur Realtime Compute for Apache Flink version VVR 11.1 ou ultérieure.

  • Nous recommandons de désactiver cette fonctionnalité dans les scénarios suivants :

    • Les données de la table de dimension sont fréquemment mises à jour et des données en temps réel sont requises.

    • La table contient une grande quantité de données, ce qui pose un risque de dépassement de mémoire.

Correspondance des types de données

Type de données StarRocks

Type de données Flink

NULL

NULL

BOOLEAN

BOOLEAN

TINYINT

TINYINT

SMALLINT

SMALLINT

INT

INT

BIGINT

BIGINT

BIGINT UNSIGNED

Remarque

Nécessite le moteur Realtime Compute for Apache Flink VVR 8.0.10 ou version ultérieure.

DECIMAL(20,0)

LARGEINT

DECIMAL(20,0)

FLOAT

FLOAT

DOUBLE

DOUBLE

DATE

DATE

DATETIME

TIMESTAMP

DECIMAL

DECIMAL

DECIMALV2

DECIMAL

DECIMAL32

DECIMAL

DECIMAL64

DECIMAL

DECIMAL128

DECIMAL

CHAR(m)

Remarque
  • La version VVR 8.0.10 étend automatiquement la longueur CHAR par trois (m=n*3, où n<=85) pour prendre en compte les différences d'encodage entre MySQL et StarRocks.

  • Les versions VVR 8.0.11 et ultérieures étendent automatiquement la longueur CHAR par quatre (m=n*4, où n<=63) pour prendre en compte les différences d'encodage entre MySQL et StarRocks.

  • La longueur maximale du type CHAR dans StarRocks est de 255. Par conséquent, Flink mappe un type CHAR vers un type CHAR StarRocks uniquement si sa longueur étendue automatiquement ne dépasse pas 255.

CHAR(n)

VARCHAR(m)

Remarque
  • La version VVR 8.0.10 étend automatiquement la longueur VARCHAR par trois (m=n*3, où n>85) pour prendre en compte les différences d'encodage entre MySQL et StarRocks.

  • Les versions VVR 8.0.11 et ultérieures étendent automatiquement la longueur VARCHAR par quatre (m=n*4, où n>63) pour prendre en compte les différences d'encodage entre MySQL et StarRocks.

  • La longueur maximale du type CHAR dans StarRocks est de 255. Par conséquent, si la longueur étendue automatiquement d'un type CHAR Flink dépasse 255, Flink mappe ce type vers le type VARCHAR de StarRocks.

CHAR(n)

VARCHAR

STRING

VARBINARY

Remarque

Nécessite le moteur Realtime Compute for Apache Flink VVR 8.0.10 ou version ultérieure.

VARBINARY

Exemple de code

CREATE TEMPORARY TABLE IF NOT EXISTS `runoob_tbl_source` (
  `runoob_id` BIGINT NOT NULL,
  `runoob_title` STRING NOT NULL,
  `runoob_author` STRING NOT NULL,
  `submission_date` DATE NULL
) WITH (
  'connector' = 'starrocks',
  'jdbc-url' = 'jdbc:mysql://ip:9030',
  'scan-url' = 'ip:18030',
  'database-name' = 'db_name',
  'table-name' = 'table_name',
  'password' = 'xxxxxxx',
  'username' = 'xxxxx'
);
CREATE TEMPORARY TABLE IF NOT EXISTS `runoob_tbl_sink` (
  `runoob_id` BIGINT NOT NULL,
  `runoob_title` STRING NOT NULL,
  `runoob_author` STRING NOT NULL,
  `submission_date` DATE NULL
  PRIMARY KEY(`runoob_id`)
  NOT ENFORCED
) WITH (
  'jdbc-url' = 'jdbc:mysql://ip:9030',
  'connector' = 'starrocks',
  'load-url' = 'ip:18030',
  'database-name' = 'db_name',
  'table-name' = 'table_name',
  'password' = 'xxxxxxx',
  'username' = 'xxxx',
  'sink.buffer-flush.interval-ms' = '5000'
);

INSERT INTO runoob_tbl_sink SELECT * FROM runoob_tbl_source;
Remarque

StarRocks autorise une colonne de clé primaire à être NULLABLE. Toutefois, Flink ne prend pas en charge les clés primaires contenant des colonnes acceptant les valeurs NULL. Le modèle de cohérence des données de Flink exige qu'une clé primaire soit unique et non nullable. Dans le cas contraire, Flink génère l'erreur Invalid primary key. Column 'xxx' is nullable. Pour plus d'informations, consultez la rubrique Erreur « Invalid primary key. Column 'xxx' is nullable ».

Ingestion des données

Utilisez le connecteur Pipeline StarRocks pour écrire les enregistrements de données et les modifications de schéma provenant des sources de données en amont vers une base de données StarRocks externe. Le connecteur StarRocks prend en charge à la fois l'édition communautaire et le service entièrement géré EMR Serverless StarRocks d'Alibaba Cloud.

Fonctionnalités

  • Création automatique de bases de données et de tables.

    Si une base de données ou une table en amont n'existe pas dans l'instance StarRocks en aval, le connecteur la crée automatiquement. Vous pouvez utiliser le paramètre table.create.properties.* pour configurer les options de création automatique de table.

  • Synchronisation des modifications de schéma.

    Le connecteur StarRocks applique automatiquement les événements CreateTableEvent, AddColumnEvent et DropColumnEvent à la base de données en aval.

  • Les versions VVR 11.1 et ultérieures prennent en charge les modifications compatibles des types de colonnes. Pour plus d'informations, consultez la documentation ALTER TABLE | StarRocks.

Notes d'utilisation

  • Chaque table synchronisée doit disposer d'une clé primaire. Pour les tables sans clé primaire, vous devez en spécifier une dans le bloc transform afin d'écrire les données en aval. Exemple :

    transform:
      - source-table: ...
        primary-keys: id, ...
  • Pour les tables créées automatiquement, la clé de bucket est identique à la clé primaire et la table ne peut pas avoir de clé de partition.

  • Lors de la synchronisation des modifications de schéma, les nouvelles colonnes ne peuvent être ajoutées qu'à la fin des colonnes existantes. En mode de modification de schéma Lenient (par défaut), les insertions à d'autres positions sont automatiquement déplacées à la fin.

  • Si vous utilisez une version de StarRocks antérieure à la 2.5.7, vous devez spécifier explicitement le nombre de buckets avec le paramètre table.create.num-buckets . Les versions 2.5.7 et ultérieures de StarRocks peuvent déterminer automatiquement un nombre approprié de buckets.

  • Si vous utilisez StarRocks 3.2 ou une version ultérieure, nous vous recommandons d'activer l'option table.create.properties.fast_schema_evolution pour accélérer les modifications de schéma.

  • Des problèmes de streaming peuvent survenir lors de l'utilisation de CDC YAML pour l'ingestion de données dans EMR Serverless StarRocks. Vous pouvez appliquer l'une des solutions de contournement suivantes :

    • Utilisez le connecteur Flink SQL StarRocks et définissez le paramètre sink.version=V1 .

    • Activez le paramètre FE emr_internal_redirect .

    • Utilisez un nom de domaine Private Zone StarRocks au lieu d'un SLB.

Syntaxe

source:
  ...

sink:
  type: starrocks
  name: StarRocks Sink
  jdbc-url: jdbc:mysql://127.0.0.1:9030
  load-url: 127.0.0.1:8030
  username: root
  password: pass
  sink.buffer-flush.interval-ms: 5000   # Set the data flush interval.

Configuration

Parameter

Description

Type

Required

Default

Remarks

type

Spécifie le type de connecteur sink.

String

Yes

Définissez la valeur sur starrocks.

name

Nom d'affichage du sink.

String

No

jdbc-url

URL JDBC pour la connexion à la base de données.

String

Yes

Prend en charge plusieurs adresses séparées par des virgules (,). Exemple : jdbc:mysql://fe_host1:fe_query_port1,fe_host2:fe_query_port2,fe_host3:fe_query_port3.

load-url

URL HTTP d'un nœud FE pour Stream Load.

String

Yes

Prend en charge plusieurs adresses séparées par des points-virgules (;). Exemple : fe_host1:fe_http_port1;fe_host2:fe_http_port2.

username

Nom d'utilisateur pour la connexion StarRocks.

String

Yes

Cet utilisateur doit disposer au minimum des autorisations SELECT et INSERT sur la table cible. Vous pouvez accorder les autorisations requises à l'aide de la commande StarRocks GRANT.

password

Mot de passe pour la connexion StarRocks.

String

Yes

sink.semantic

Sémantique de livraison pour l'écriture des données.

String

No

at-least-once

Valeurs valides :

  • at-least-once (par défaut) : Garantit que les données sont livrées au moins une fois.

  • exactly-once : Garantit que les données sont livrées exactement une fois.

sink.label-prefix

Préfixe d'étiquette pour les tâches Stream Load.

String

No

sink.connect.timeout-ms

Délai d'attente pour l'établissement d'une connexion HTTP.

Integer

No

30000

Unité : millisecondes. La valeur doit être comprise entre 100 et 60 000.

sink.wait-for-continue.timeout-ms

Délai d'attente pour recevoir une réponse 100 Continue du serveur.

Integer

No

30000

Unité : millisecondes. La valeur doit être comprise entre 3 000 et 600 000.

sink.buffer-flush.max-bytes

Taille maximale du cache en mémoire, en octets, avant le déclenchement d'un vidage.

Long

No

157286400

Unité : octets. La valeur doit être comprise entre 64 Mo et 10 Go.

Remarque
  • Cette taille de cache est partagée par toutes les tables. Lorsque le tampon est plein, le connecteur sélectionne plusieurs tables à vider.

  • Une valeur plus élevée peut améliorer le débit, mais risque d'augmenter la latence d'ingestion.

sink.buffer-flush.max-rows

Nombre maximal de lignes dans le cache en mémoire avant le déclenchement d'un vidage.

Long

No

500000

La valeur doit être comprise entre 64 000 et 5 000 000.

sink.buffer-flush.interval-ms

Intervalle de temps entre les vidages du tampon pour chaque table.

Long

No

300000

Unité : millisecondes.

Remarque

Pour les tâches synchronisant de faibles volumes de données, réduisez cette valeur afin d'éviter de longs délais avant la persistance des données.

sink.max-retries

Nombre maximal de tentatives.

Long

No

3

La valeur doit être comprise entre 0 et 1 000.

sink.scan-frequency.ms

Fréquence à laquelle le connecteur vérifie s'il doit vider le tampon.

Long

No

50

Unité : millisecondes.

sink.io.thread-count

Nombre de threads utilisés pour Stream Load.

Integer

No

2

sink.at-least-once.use-transaction-stream-load

Indique s'il faut utiliser l'interface de transaction Stream Load pour l'ingestion des données.

Boolean

No

true

Cette option ne prend effet que si la base de données la prend en charge.

sink.properties.*

Propriétés supplémentaires pour le sink.

String

No

Pour connaître les propriétés prises en charge, reportez-vous à la documentation STREAM LOAD.

table.create.num-buckets

Nombre de buckets pour les tables créées automatiquement.

Integer

No

table.create.properties.*

Propriétés supplémentaires pour la création automatique de table.

String

No

Par exemple, vous pouvez transmettre 'table.create.properties.fast_schema_evolution' = 'true' pour activer l'évolution rapide du schéma. Pour plus de détails, consultez la documentation StarRocks.

table.schema-change.timeout

Délai d'attente pour les opérations de modification de schéma.

Duration

No

30 min

Doit être un nombre entier de secondes.

Remarque

Si une opération de modification de schéma dépasse cette limite, la tâche échoue.

unicode-char.max-bytes

Nombre d'octets alloués pour chaque caractère Unicode.

Integer

No

3

Dans CDC, la longueur d'un type VARCHAR est mesurée en caractères, tandis que dans StarRocks, elle est mesurée en octets.

Dans la plupart des cas, un caractère Unicode ne dépasse pas 3 octets après encodage UTF-8. Toutefois, certains caractères rares et emojis peuvent occuper 4 octets ou plus.

Réutiliser un catalogue intégré

À partir de VVR 11.5, vous pouvez référencer directement un catalogue StarRocks intégré créé sur la page Data Management dans une tâche d'ingestion de données Flink CDC. Cette approche simplifie la configuration en réduisant le nombre de propriétés à définir manuellement.

sink:
  type: starrocks
  using.built-in-catalog: starrocks_catalog

Les tâches d'ingestion de données peuvent réutiliser automatiquement les options de catalogue StarRocks suivantes :

  • jdbc-url

  • http-url

  • username

  • password

  • table.num-buckets

Pour remplacer ces valeurs, définissez explicitement les options YAML correspondantes, qui ont priorité.

Mappage de types

Remarque

StarRocks ne prend pas en charge tous les types YAML CDC. L'écriture d'un type non pris en charge dans le sink entraîne l'échec de la tâche. Vous pouvez utiliser la fonction intégrée CAST dans une transformation pour convertir les données non prises en charge, ou utiliser une instruction de projection pour les exclure de la table de résultat. Pour plus d'informations, reportez-vous à la rubrique Développer une tâche d'ingestion de données Flink CDC.

CDC type

StarRocks type

Remarks

TINYINT

TINYINT

SMALLINT

SMALLINT

INT

INT

BIGINT

BIGINT

FLOAT

FLOAT

DOUBLE

DOUBLE

BOOLEAN

BOOLEAN

DATE

DATE

TIMESTAMP

DATETIME

TIMESTAMP_LTZ

DATETIME

DECIMAL(p, s)

DECIMAL(p, s)

Étant donné que StarRocks ne prend pas en charge DECIMAL pour une clé primaire, le connecteur convertit automatiquement une colonne de clé primaire DECIMAL amont en VARCHAR dans le schéma StarRocks synchronisé.

CHAR(n)

(n <= 85)

CHAR(n × 3)

CDC mesure la longueur en caractères, tandis que StarRocks utilise les octets. Le connecteur multiplie la longueur par 3 pour tenir compte des caractères UTF-8 multioctets.

Remarque

La longueur maximale du type CHAR StarRocks est de 255. Par conséquent, seuls les types CHAR CDC dont la longueur est inférieure ou égale à 85 sont mappés au type CHAR StarRocks.

Remarque

Vous pouvez définir le paramètre unicode-char.max-bytes pour allouer davantage d'octets à chaque caractère Unicode.

CHAR(n)

(n > 85)

VARCHAR(n × 3)

CDC mesure la longueur en caractères, tandis que StarRocks utilise les octets. Le connecteur multiplie la longueur par 3 pour tenir compte des caractères UTF-8 multioctets.

Remarque

CDC mesure la longueur en caractères, tandis que StarRocks utilise les octets. Le connecteur multiplie la longueur par 3. Étant donné que le résultat dépasse la limite de 255 octets pour le type CHAR StarRocks, il est mappé vers VARCHAR.

Remarque

Vous pouvez définir le paramètre unicode-char.max-bytes pour allouer davantage d'octets à chaque caractère Unicode.

VARCHAR(n)

VARCHAR(n × 3)

CDC mesure la longueur en caractères, tandis que StarRocks utilise les octets. Le connecteur multiplie la longueur par 3 pour tenir compte des caractères UTF-8 multioctets.

BINARY(n)

BINARY(n+2)

Deux octets de remplissage sont ajoutés pour garantir l'intégrité des données.

VARBINARY(n)

VARBINARY(n+1)

Un octet de remplissage est ajouté pour garantir l'intégrité des données.

Modification de schéma

En tant que sink d'ingestion de données, StarRocks prend en charge les événements de modification de schéma suivants :

  • CREATE TABLE EVENT

    Remarque

    Si la table StarRocks en aval existe déjà, le connecteur ne tente pas de la recréer. Assurez-vous que le schéma de la table en aval est compatible avec le schéma amont.

  • ADD COLUMN EVENT

    Remarque

    StarRocks exige que les colonnes de clé primaire apparaissent en premier dans une table. Toute nouvelle colonne doit être ajoutée après celles-ci.

  • DROP COLUMN EVENT

  • TRUNCATE TABLE EVENT

  • DROP TABLE EVENT