Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Connecteur MaxCompute

Dernière mise à jour :Aug 09, 2026

Le connecteur MaxCompute permet de lire des données depuis et d'écrire des données dans MaxCompute (anciennement ODPS), la plateforme d'entreposage de données entièrement gérée à l'échelle de l'exaoctet d'Alibaba Cloud, directement depuis les jobs Flink SQL et DataStream.

Fonctionnalités

Élément Description
Types de table Table source, table de dimension, table de destination (sink) et ingestion de données vers une table de destination
Modes d'exécution Mode streaming et mode batch
Types d'API API DataStream, API SQL et jobs YAML d'ingestion de données
Sémantique Au moins une fois (At-least-once)
Mise à jour ou suppression des données dans une table de destination Tunnel Batch ou Tunnel Streaming : insertion uniquement. Tunnel Upsert : insertion, mise à jour et suppression.

Métriques

Type de table Métriques
Source numRecordsIn, numRecordsInPerSecond, numBytesIn, numBytesInPerSecond
Destination (Sink) numRecordsOut, numRecordsOutPerSecond, numBytesOut, numBytesOutPerSecond
Table de dimension dim.odps.cacheSize
Pour plus d'informations, consultez la page Métriques de surveillance .

Prérequis

Avant de commencer, assurez-vous de disposer des éléments suivants :

Limites

  • Le connecteur prend en charge uniquement la sémantique « au moins une fois ». Des doublons peuvent apparaître dans MaxCompute selon le tunnel utilisé. Pour obtenir des conseils sur le choix du tunnel, consultez la section « Comment choisir un tunnel de données ? » de la rubrique FAQ sur le stockage amont et aval.

  • Par défaut, une source fonctionne en mode complet : elle lit uniquement les données de la partition spécifiée par l'option partition. Une fois toutes les données lues, le job se termine et ne surveille pas l'apparition de nouvelles partitions. Pour surveiller continuellement les nouvelles partitions, configurez une source incrémentielle en utilisant le paramètre startPartition.

  • À chaque actualisation du cache d'une table de dimension, le système vérifie la dernière partition disponible. Une fois la source démarrée, elle ne lit pas les nouvelles données ajoutées à une partition déjà en cours de lecture ; effectuez un déploiement uniquement lorsque la partition contient l'intégralité des données.

Choisir un tunnel

MaxCompute propose trois tunnels pour écrire des données depuis Flink. Choisissez celui qui correspond à votre cas d'utilisation :

Tunnel Par défaut Cas d'utilisation
Tunnel Batch MaxCompute Oui (useStreamTunnel=false, enableUpsert=false) Chargements batch ; les données sont disponibles uniquement après la création d'un point de contrôle (checkpoint). Définissez flushIntervalMs=0 pour désactiver la vidange planifiée.
Tunnel Streaming MaxCompute Non (useStreamTunnel=true) Ingestion en quasi temps réel ; les données vidées sont immédiatement disponibles dans MaxCompute.
Tunnel Upsert MaxCompute Non (enableUpsert=true) Opérations INSERT, UPDATE et DELETE sur une table Delta MaxCompute. Nécessite VVR 8.0.6 ou version ultérieure.

Pour une comparaison détaillée, consultez la section « Comment choisir un tunnel de données ? » de la rubrique FAQ sur le stockage amont et aval.

SQL

Le connecteur MaxCompute peut être utilisé comme table source, table de dimension ou table de destination (sink) dans les jobs basés sur SQL.

Syntaxe

CREATE TEMPORARY TABLE odps_source(
  id INT,
  user_name VARCHAR,
  content VARCHAR
) WITH (
  'connector' = 'odps',
  'endpoint' = '<yourEndpoint>',
  'project' = '<yourProjectName>',
  'schemaName' = '<yourSchemaName>',
  'tableName' = '<yourTableName>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}',
  'partition' = 'ds=2018****'
);

Options du connecteur

Options générales

Option Obligatoire Valeur par défaut Type Description
connector Oui STRING Définissez cette option sur odps.
endpoint Oui STRING Endpoint MaxCompute. Consultez la rubrique Endpoints.
tunnelEndpoint Non STRING Endpoint du tunnel MaxCompute. Si non spécifié, MaxCompute alloue les connexions au tunnel via Server Load Balancer (SLB).
project Oui STRING Nom du projet MaxCompute.
schemaName Non STRING Requis uniquement lorsque la fonctionnalité de schéma MaxCompute est activée. Définissez ce paramètre sur le nom du schéma de la table. Consultez la page Opérations liées aux schémas. VVR 8.0.6 ou version ultérieure.
tableName Oui STRING Nom de la table MaxCompute.
accessId Oui STRING ID AccessKey utilisé pour accéder à MaxCompute. Consultez la rubrique Comment afficher mon ID AccessKey et ma clé secrète AccessKey ?
Important

Stockez l'ID AccessKey en tant que variable. Consultez la rubrique Gérer les variables.

accessKey Oui STRING Clé secrète AccessKey utilisée pour accéder à MaxCompute.
partition Non STRING Nom de la partition dans la table MaxCompute. Non requis pour les tables non partitionnées ou les sources incrémentielles. Consultez la section « Comment configurer l'option de partition ? » de la rubrique FAQ sur le stockage amont et aval.
compressAlgorithm Non SNAPPY STRING Algorithme de compression pour le tunnel MaxCompute. Valeurs valides : RAW (aucune compression), ZLIB, SNAPPY. SNAPPY améliore le débit d'environ 50 % par rapport à ZLIB dans les scénarios de test.
quotaName Non STRING Nom du quota pour les groupes de ressources exclusifs du tunnel MaxCompute. VVR 8.0.3 ou version ultérieure. Ce paramètre prend effet uniquement lorsque l'endpoint est défini sur une adresse VPC. Si un endpoint public est utilisé ou si le tunnelEndpoint est spécifié, ce paramètre n'a aucun effet.

Options de source

Option Obligatoire Valeur par défaut Type Description
maxPartitionCount Non 100 INTEGER Nombre maximal de partitions à lire. Si cette limite est dépassée, l'erreur "The number of matched partitions exceeds the default limit" s'affiche. La lecture d'un trop grand nombre de partitions peut surcharger MaxCompute et ralentir le démarrage du job ; augmentez cette valeur uniquement si votre charge de travail l'exige.
useArrow Non false BOOLEAN Lire les données au format Arrow, qui appelle l'API de stockage MaxCompute. Uniquement pour les déploiements batch. VVR 8.0.8 ou version ultérieure.
splitSize Non 256 MB MEMORYSIZE Quantité de données extraites par fragment lors de l'utilisation du format Arrow. Uniquement pour les déploiements batch. VVR 8.0.8 ou version ultérieure.
compressCodec Non "" (aucun) STRING Algorithme de compression lors de la lecture avec le format Arrow. Valeurs valides : "" (aucun), ZSTD, LZ4_FRAME. La spécification d'un codec améliore le débit par rapport à l'absence de compression. Uniquement pour les déploiements batch. VVR 8.0.8 ou version ultérieure.
dynamicLoadBalance Non false BOOLEAN Activez l'allocation dynamique des shards pour améliorer les performances de traitement et réduire le temps de lecture global. Notez que cela peut provoquer un déséquilibre des données (data skew) car différents opérateurs lisent des quantités de données incohérentes. Uniquement pour les déploiements batch. VVR 8.0.8 ou version ultérieure.

Options de source incrémentielle

La source incrémentielle interroge MaxCompute par intermittence pour découvrir de nouvelles partitions. Avant de lire une nouvelle partition, toutes les écritures de données dans cette partition doivent être terminées. Pour plus de détails, consultez la section « Que faire si une source incrémentielle détecte une nouvelle partition alors que des données sont encore en cours d'écriture ? » de la rubrique FAQ sur le stockage amont et aval.

Ordre des partitions : la source lit les partitions dont l'ordre alphabétique est supérieur ou égal à la valeur de startPartition. Par exemple, year=2023,month=10 vient avant year=2023,month=9 dans l'ordre alphabétique ; utilisez donc un remplissage par zéro pour les valeurs de mois (utilisez year=2023,month=09 au lieu de year=2023,month=9) afin de garantir un ordre correct.

Option Obligatoire Valeur par défaut Type Description
startPartition Oui STRING Partition de départ pour les lectures incrémentielles. Lorsque ce paramètre est spécifié, partition est ignoré. Pour les tables partitionnées à plusieurs niveaux, configurez les valeurs des colonnes de partition par ordre décroissant de niveau. Consultez la section « Comment configurer startPartition ? » de la rubrique FAQ sur le stockage amont et aval.
subscribeIntervalInSec Non 30 INTEGER Intervalle d'interrogation en secondes.
modifiedTableOperation Non NONE Enum Action à entreprendre lorsqu'une partition est modifiée pendant la lecture. Les sessions de téléchargement sont enregistrées dans les points de contrôle ; si les données d'une partition changent après le démarrage d'une session, la reprise à partir du point de contrôle échoue et le déploiement redémarre en boucle. Valeurs valides : NONE — mettez à jour startPartition pour ignorer la partition indisponible et redémarrer sans état ; SKIP — ignore automatiquement la partition indisponible lors de la reprise. VVR 8.0.3 ou version ultérieure. Si l'une de ces valeurs est définie, les données déjà lues à partir de la partition modifiée sont conservées ; les données non lues sont supprimées.

Options de destination (sink)

Option Obligatoire Valeur par défaut Type Description
useStreamTunnel Non false BOOLEAN Utilisez le tunnel Streaming MaxCompute au lieu du tunnel Batch. true : tunnel Streaming ; false : tunnel Batch. Consultez la section Choisir un tunnel.
flushIntervalMs Non 30000 (30 s) LONG Intervalle de vidage du tampon de l'outil d'écriture du tunnel, en millisecondes. Pour le tunnel Streaming : les données vidées sont immédiatement disponibles. Pour le tunnel Batch : les données deviennent disponibles uniquement après la création d'un point de contrôle — définissez cette option sur 0 pour désactiver la vidange planifiée. Le déclenchement a lieu lorsque flushIntervalMs ou batchSize est atteint.
batchSize Non 67108864 (64 Mo) LONG Taille du tampon en octets. Les données sont vidées lorsque le tampon atteint cette taille. Le déclenchement a lieu lorsque batchSize ou flushIntervalMs est atteint.
numFlushThreads Non 1 INTEGER Nombre de threads utilisés pour vider le tampon de l'outil d'écriture du tunnel. Les valeurs supérieures à 1 permettent une vidange simultanée sur plusieurs partitions.
slotNum Non 0 INTEGER Nombre de slots de tunnel pour recevoir les données de Flink. Consultez la page Présentation du service de transmission de données pour connaître les limites de slots.
dynamicPartitionLimit Non 100 INTEGER Nombre maximal de partitions dynamiques écrites entre deux points de contrôle. Si cette limite est dépassée, l'erreur "Too many dynamic partitions" s'affiche. L'écriture dans un grand nombre de partitions augmente la charge sur MaxCompute et ralentit la création des points de contrôle ; augmentez cette valeur uniquement si votre charge de travail l'exige.
retryTimes Non 3 INTEGER Nombre maximal de tentatives pour les requêtes vers le serveur MaxCompute (création de session, soumission ou échecs de vidage).
sleepMillis Non 1000 INTEGER Intervalle de nouvelle tentative en millisecondes.
enableUpsert Non false BOOLEAN Utilisez le tunnel Upsert MaxCompute. true : traite les enregistrements INSERT, UPDATE_AFTER et DELETE ; false : utilise le tunnel spécifié par useStreamTunnel. VVR 8.0.6 ou version ultérieure. Si la destination rencontre des erreurs ou des pannes prolongées lors des validations de session en mode upsert, définissez le parallélisme de l'opérateur de destination sur 10 ou moins.
upsertAsyncCommit Non false BOOLEAN Utilisez le mode asynchrone lors de la validation des sessions upsert. Le mode asynchrone réduit le temps de validation, mais les données validées ne sont pas immédiatement interrogables. VVR 8.0.6 ou version ultérieure.
upsertCommitTimeoutMs Non 120000 (120 s) INTEGER Délai d'expiration pour les validations de session upsert, en millisecondes. VVR 8.0.6 ou version ultérieure.
sink.operation Non insert STRING Mode d'écriture pour une table Delta. insert : mode ajout (append) ; upsert : mode mise à jour. VVR 8.0.10 ou version ultérieure.
sink.parallelism Non INTEGER Parallélisme d'écriture pour une table Delta. La valeur par défaut correspond au parallélisme en amont. La valeur de write.bucket.num doit être un multiple entier de sink.parallelism pour des performances d'écriture et une efficacité mémoire optimales. VVR 8.0.10 ou version ultérieure.
sink.file-cached.enable Non false BOOLEAN Activez le mode de cache de fichier lors de l'écriture dans des partitions dynamiques d'une table Delta. Réduit le nombre de petits fichiers écrits sur le serveur, mais augmente la latence d'écriture. Activez cette option lorsque la destination présente un parallélisme élevé. VVR 8.0.10 ou version ultérieure.
sink.file-cached.writer.num Non 16 INTEGER Nombre de threads de téléchargement simultanés par tâche en mode de cache de fichier. Évitez de définir cette valeur trop élevée, car l'écriture simultanée dans de nombreuses partitions peut provoquer des erreurs de mémoire insuffisante (OOM). Effective uniquement lorsque sink.file-cached.enable=true. VVR 8.0.10+.
sink.bucket.check-interval Non 60000 INTEGER Intervalle de vérification de la taille des fichiers en mode de cache de fichier, en millisecondes. Effectif uniquement lorsque sink.file-cached.enable=true. VVR 8.0.10+.
sink.file-cached.rolling.max-size Non 16 MB MEMORYSIZE Taille maximale d'un fichier mis en cache. Lorsque cette limite est dépassée, les données sont téléchargées vers le serveur. Effectif uniquement lorsque sink.file-cached.enable=true. VVR 8.0.10+.
sink.file-cached.memory Non 64 MB MEMORYSIZE Mémoire hors tas (off-heap) maximale pour les écritures de fichiers en mode de cache de fichier. Effectif uniquement lorsque sink.file-cached.enable=true. VVR 8.0.10+.
sink.file-cached.memory.segment-size Non 128 KB MEMORYSIZE Taille du segment de tampon pour les écritures de fichiers en mode de cache de fichier. Effectif uniquement lorsque sink.file-cached.enable=true. VVR 8.0.10+.
sink.file-cached.flush.always Non true BOOLEAN Indique si le cache doit être utilisé lors de l'écriture de fichiers en mode de cache de fichier. Effectif uniquement lorsque sink.file-cached.enable=true. VVR 8.0.10+.
sink.file-cached.write.max-retries Non 3 INTEGER Nombre de tentatives pour les téléchargements de données en mode de cache de fichier. Effectif uniquement lorsque sink.file-cached.enable=true. VVR 8.0.10+.
upsert.writer.max-retries Non 3 INTEGER Nombre maximal de tentatives d'écriture dans un bucket lors d'une session Upsert Writer. VVR 8.0.10 ou version ultérieure.
upsert.writer.buffer-size Non 64 MB MEMORYSIZE Taille totale du tampon pour tous les buckets lors d'une session Upsert Writer. Les données sont vidées lorsque le total atteint ce seuil. Augmentez cette valeur pour une meilleure efficacité d'écriture ; réduisez-la si l'écriture dans de nombreuses partitions provoque des erreurs OOM. VVR 8.0.10 ou version ultérieure.
upsert.writer.bucket.buffer-size Non 1 MB MEMORYSIZE Taille du tampon par bucket lors d'une session Upsert Writer. Réduisez cette valeur si la mémoire du serveur Flink est insuffisante. VVR 8.0.10 ou version ultérieure.
upsert.write.bucket.num Oui INTEGER Nombre de buckets pour la table Delta cible. Doit correspondre à la valeur de write.bucket.num configurée sur la table Delta. VVR 8.0.10 ou version ultérieure.
upsert.write.slot-num Non 1 INTEGER Slots de tunnel par session upsert. VVR 8.0.10 ou version ultérieure.
upsert.commit.max-retries Non 3 INTEGER Nombre maximal de tentatives pour les validations de session upsert. VVR 8.0.10 ou version ultérieure.
upsert.commit.thread-num Non 16 INTEGER Parallélisme pour les validations de session upsert. Évitez les valeurs élevées, car des validations simultanées excessives augmentent la consommation de ressources et peuvent entraîner des problèmes de performance. VVR 8.0.10 ou version ultérieure.
upsert.commit.timeout Non 600 INTEGER Délai d'expiration pour les validations de session upsert, en secondes. VVR 8.0.10 ou version ultérieure.
upsert.flush.concurrent Non 2 INTEGER Nombre maximal de vidages de bucket simultanés par partition. Chaque vidage de bucket occupe un slot de tunnel. VVR 8.0.10 ou version ultérieure.
insert.commit.thread-num Non 16 INTEGER Parallélisme pour les validations de session d'insertion. VVR 8.0.10 ou version ultérieure.
insert.arrow-writer.enable Non false BOOLEAN Utilisez le format Arrow pour les insertions. VVR 8.0.10 ou version ultérieure.
insert.arrow-writer.batch-size Non 512 INTEGER Nombre maximal de lignes par lot au format Arrow. VVR 8.0.10 ou version ultérieure.
insert.arrow-writer.flush-interval Non 100000 INTEGER Intervalle de vidage de l'outil d'écriture en millisecondes. VVR 8.0.10 ou version ultérieure.
insert.writer.buffer-size Non 64 MB MEMORYSIZE Taille du cache pour l'outil d'écriture tamponné. VVR 8.0.10 ou version ultérieure.
upsert.partial-column.enable Non false BOOLEAN Mettez à jour uniquement les colonnes spécifiées (mise à jour partielle des colonnes). S'applique uniquement aux destinations de table Delta. Consultez la page Mettre à jour les données dans des colonnes spécifiques. Lorsque true : si un enregistrement avec la même clé primaire existe, les champs non nuls spécifiés sont écrasés ; si aucun enregistrement correspondant n'existe, un nouvel enregistrement est inséré avec les nouvelles valeurs pour les colonnes spécifiées et null pour toutes les colonnes non spécifiées. VVR 8.0.11 ou version ultérieure.

Options de table de dimension

Lorsqu'un déploiement démarre, la table de dimension charge toutes les données de la partition spécifiée par partition. L'option partition prend en charge la fonction max_pt(). Lors du rechargement du cache, la dernière partition est relue. Définissez partition sur max_two_pt() pour charger les données de deux partitions.

Les tables de dimension nécessitent cache=ALL . Augmentez la mémoire du nœud de jointure à au moins quatre fois la taille des données de la table distante. Pour les grandes tables de dimension, utilisez l'indication SHUFFLE_HASH pour distribuer les données uniformément. Pour les tables extrêmement volumineuses qui provoquent des garbage collections (GC) fréquentes de la machine virtuelle Java (JVM), convertissez-les en table de dimension clé-valeur avec une politique de cache LRU (Least Recently Used) — par exemple, une table de dimension ApsaraDB for HBase.
Option Obligatoire Valeur par défaut Type Description
cache Oui STRING Politique de cache. Doit être définie sur ALL et explicitement déclarée dans l'instruction DDL. Toutes les données de la table de dimension sont chargées dans le cache avant l'exécution du déploiement. Les recherches ultérieures interrogent uniquement le cache. Le cache est rechargé après l'expiration des entrées.
cacheSize Non 100000 LONG Nombre maximal de lignes à mettre en cache. Si cette limite est dépassée, l'erreur "Row count of table <table-name> partition <partition-name> exceeds maxRowCount limit" s'affiche. Les caches volumineux consomment beaucoup de mémoire tas JVM et ralentissent le démarrage et l'actualisation du cache ; augmentez cette valeur uniquement si votre charge de travail l'exige.
cacheTTLMs Non Long.MAX_VALUE LONG Délai d'expiration du cache en millisecondes.
cacheReloadTimeBlackList Non STRING Périodes durant lesquelles le cache n'est pas actualisé. Utilisez cette option pendant les pics de trafic (comme les événements promotionnels) pour éviter l'instabilité du déploiement due aux actualisations du cache. Consultez la section « Comment configurer cacheReloadTimeBlackList ? » de la rubrique FAQ sur le stockage amont et aval.
maxLoadRetries Non 10 INTEGER Nombre maximal de tentatives pour le chargement initial du cache au démarrage du déploiement. Si le nombre de tentatives est épuisé, le déploiement échoue.

Mappages de types de données

Pour la liste complète des types de données MaxCompute, consultez la page Système de types de données MaxCompute version 2.0.

Type MaxCompute Type Flink
BOOLEAN BOOLEAN
TINYINT TINYINT
SMALLINT SMALLINT
INT INTEGER
BIGINT BIGINT
FLOAT FLOAT
DOUBLE DOUBLE
DECIMAL(precision, scale) DECIMAL(precision, scale)
CHAR(n) CHAR(n)
VARCHAR(n) VARCHAR(n)
STRING STRING
BINARY BYTES
DATE DATE
DATETIME TIMESTAMP(3)
TIMESTAMP TIMESTAMP(9)
TIMESTAMP_NTZ TIMESTAMP(9)
ARRAY ARRAY
MAP MAP
STRUCT ROW
JSON STRING
Important

Si une table physique MaxCompute contient des champs de type composite imbriqué (ARRAY, MAP ou STRUCT) ou un champ JSON, définissez tblproperties('columnar.nested.type'='true') lors de la création de la table pour permettre à Realtime Compute for Apache Flink de lire et d'écrire les données correctement.

Flink CDC (aperçu public)

Le connecteur MaxCompute peut ingérer des données Change Data Capture (CDC) en tant que destination dans les jobs basés sur YAML. Nécessite VVR 11.1 ou version ultérieure.

Syntaxe

source:
  type: xxx

sink:
  type: maxcompute
  name: MaxComputeSink
  access-id: ${your_accessId}
  access-key: ${your_accessKey}
  endpoint: ${your_maxcompute_endpoint}
  project: ${your_project}
  buckets-num: 8

Options de configuration

Option Obligatoire Valeur par défaut Type Description
type Oui String Définissez cette option sur maxcompute.
name Non String Nom de la destination.
access-id Oui String ID AccessKey de votre compte Alibaba Cloud ou utilisateur RAM. Obtenez-le depuis la console Resource Access Management (RAM).
access-key Oui String Clé secrète AccessKey.
endpoint Oui String Endpoint MaxCompute. Configurez-le en fonction de la région et de la méthode de connexion réseau. Consultez la page Endpoints.
project Oui String Nom du projet MaxCompute. Pour le trouver : connectez-vous à la console MaxCompute, accédez à Workspace > Projects et copiez le nom du projet.
tunnel.endpoint Non String Endpoint du tunnel MaxCompute. Généralement déduit automatiquement. Requis dans des environnements réseau spéciaux, tels qu'avec un serveur proxy.
quota.name Non String Nom du quota pour un groupe de ressources exclusif. Si non spécifié, un groupe de ressources partagé est utilisé. Ce paramètre prend effet uniquement lorsque endpoint est défini sur une adresse VPC. Si un endpoint public est utilisé ou si tunnel.endpoint est spécifié, ce paramètre n'a aucun effet.
sts-token Non String Jeton Security Token Service (STS) pour l'authentification par rôle RAM. Requis lors de l'accès à MaxCompute avec un rôle RAM.
buckets-num Non 16 Integer Nombre de buckets pour une table Delta MaxCompute créée automatiquement. Consultez la page Entrepôt de données en quasi temps réel.
compress.algorithm Non zlib String Algorithme de compression des données. Valeurs valides : raw (aucune compression), zlib, snappy.
total.buffer-size Non 64 MB String Taille du tampon en mémoire. Pour les tables partitionnées : s'applique par partition. Pour les tables non partitionnées : s'applique par table. Les tampons pour différentes partitions ou tables sont indépendants. Les données sont vidées lorsque le tampon est plein.
bucket.buffer-size Non 4 MB String Taille du tampon par bucket. S'applique uniquement lors de l'écriture dans des tables Delta MaxCompute.
commit.thread-num Non 16 Integer Nombre maximal de partitions ou de tables validées simultanément lors de la création d'un point de contrôle.
flush.concurrent-num Non 4 Integer Nombre maximal de buckets vidés simultanément. S'applique uniquement lors de l'écriture dans des tables Delta MaxCompute.

Mappages d'emplacement de table

Lorsque le connecteur crée automatiquement des tables dans MaxCompute, les emplacements sont mappés comme suit :

Important

Si la fonctionnalité de schéma est désactivée pour votre projet MaxCompute, le connecteur ignore tableId.namespace. Dans ce cas, une seule base de données (ou son équivalent logique) est ingérée dans MaxCompute — par exemple, une seule base de données MySQL lors de l'ingestion depuis MySQL.

Emplacement MySQL Abstraction Flink CDC Emplacement MaxCompute
N/A Projet (issu de la configuration) Projet
Base de données TableId.namespace Schéma (ignoré si la fonctionnalité de schéma est désactivée)
Table TableId.tableName Table

Mappages de types de données

Type Flink CDC Type MaxCompute
CHAR STRING
VARCHAR STRING
BOOLEAN BOOLEAN
BINARY/VARBINARY BINARY
DECIMAL DECIMAL
TINYINT TINYINT
SMALLINT SMALLINT
INTEGER INTEGER
BIGINT BIGINT
FLOAT FLOAT
DOUBLE DOUBLE
TIME_WITHOUT_TIME_ZONE STRING
DATE DATE
TIMESTAMP_WITHOUT_TIME_ZONE TIMESTAMP_NTZ
TIMESTAMP_WITH_LOCAL_TIME_ZONE (précision > 3) TIMESTAMP
TIMESTAMP_WITH_LOCAL_TIME_ZONE (précision <= 3) DATETIME
TIMESTAMP_WITH_TIME_ZONE (précision > 3) TIMESTAMP
TIMESTAMP_WITH_TIME_ZONE (précision <= 3) DATETIME
ARRAY ARRAY
MAP MAP
ROW STRUCT

Exemples

API SQL

Table source

Lire toutes les données d'une partition

Lisez toutes les données de la partition spécifiée par partition :

CREATE TEMPORARY TABLE odps_source (
  cid VARCHAR,
  rt DOUBLE
) WITH (
  'connector' = 'odps',
  'endpoint' = '<yourEndpointName>',
  'project' = '<yourProjectName>',
  'tableName' = '<yourTableName>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}',
  'partition' = 'ds=201809*'
);

CREATE TEMPORARY TABLE blackhole_sink (
  cid VARCHAR,
  invoke_count BIGINT
) WITH (
  'connector' = 'blackhole'
);

INSERT INTO blackhole_sink
SELECT
   cid,
   COUNT(*) AS invoke_count
FROM odps_source GROUP BY cid;

Lire des données incrémentielles

Lisez les données à partir de la partition spécifiée par startPartition et surveillez continuellement les nouvelles partitions :

CREATE TEMPORARY TABLE odps_source (
  cid VARCHAR,
  rt DOUBLE
) WITH (
  'connector' = 'odps',
  'endpoint' = '<yourEndpointName>',
  'project' = '<yourProjectName>',
  'tableName' = '<yourTableName>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}',
  'startPartition' = 'yyyy=2018,MM=09,dd=05' -- Start reading from the 20180905 partition.
);

CREATE TEMPORARY TABLE blackhole_sink (
  cid VARCHAR,
  invoke_count BIGINT
) WITH (
  'connector' = 'blackhole'
);

INSERT INTO blackhole_sink
SELECT cid, COUNT(*) AS invoke_count
FROM odps_source GROUP BY cid;

Table de destination

Écrire dans une partition statique

Écrivez dans la partition spécifiée par partition :

CREATE TEMPORARY TABLE datagen_source (
  id INT,
  len INT,
  content VARCHAR
) WITH (
  'connector' = 'datagen'
);

CREATE TEMPORARY TABLE odps_sink (
  id INT,
  len INT,
  content VARCHAR
) WITH (
  'connector' = 'odps',
  'endpoint' = '<yourEndpoint>',
  'project' = '<yourProjectName>',
  'tableName' = '<yourTableName>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}',
  'partition' = 'ds=20180905' -- Write to partition 20180905.
);

INSERT INTO odps_sink
SELECT
  id, len, content
FROM datagen_source;

Écrire dans des partitions dynamiques

Écrivez des données dans des partitions déterminées au moment de l'exécution par les valeurs de la colonne ds :

CREATE TEMPORARY TABLE datagen_source (
  id INT,
  len INT,
  content VARCHAR,
  c TIMESTAMP
) WITH (
  'connector' = 'datagen'
);

CREATE TEMPORARY TABLE odps_sink (
  id  INT,
  len INT,
  content VARCHAR,
  ds VARCHAR -- Dynamic partition column.
) WITH (
  'connector' = 'odps',
  'endpoint' = '<yourEndpoint>',
  'project' = '<yourProjectName>',
  'tableName' = '<yourTableName>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}',
  'partition' = 'ds' -- Omit the value; data is routed to partitions based on the ds field.
);

INSERT INTO odps_sink
SELECT
   id,
   len,
   content,
   DATE_FORMAT(c, 'yyMMdd') as ds
FROM datagen_source;

Table de dimension

Clé à valeur unique

Spécifiez une clé primaire lorsque chaque clé correspond à exactement une ligne :

CREATE TEMPORARY TABLE datagen_source (
  k INT,
  v VARCHAR
) WITH (
  'connector' = 'datagen'
);

CREATE TEMPORARY TABLE odps_dim (
  k INT,
  v VARCHAR,
  PRIMARY KEY (k) NOT ENFORCED  -- Specify the primary key.
) WITH (
  'connector' = 'odps',
  'endpoint' = '<yourEndpoint>',
  'project' = '<yourProjectName>',
  'tableName' = '<yourTableName>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}',
  'partition' = 'ds=20180905',
  'cache' = 'ALL'
);

CREATE TEMPORARY TABLE blackhole_sink (
  k VARCHAR,
  v1 VARCHAR,
  v2 VARCHAR
) WITH (
  'connector' = 'blackhole'
);

INSERT INTO blackhole_sink
SELECT k, s.v, d.v
FROM datagen_source AS s
INNER JOIN odps_dim FOR SYSTEM_TIME AS OF PROCTIME() AS d ON s.k = d.k;

Clé à valeurs multiples

Omettez la clé primaire lorsqu'une clé peut correspondre à plusieurs lignes :

CREATE TEMPORARY TABLE datagen_source (
  k INT,
  v VARCHAR
) WITH (
  'connector' = 'datagen'
);

CREATE TEMPORARY TABLE odps_dim (
  k INT,
  v VARCHAR
  -- No primary key needed for multi-value lookups.
) WITH (
  'connector' = 'odps',
  'endpoint' = '<yourEndpoint>',
  'project' = '<yourProjectName>',
  'tableName' = '<yourTableName>',
  'accessId' = '${secret_values.ak_id}',
  'accessKey' = '${secret_values.ak_secret}',
  'partition' = 'ds=20180905',
  'cache' = 'ALL'
);

CREATE TEMPORARY TABLE blackhole_sink (
  k VARCHAR,
  v1 VARCHAR,
  v2 VARCHAR
) WITH (
  'connector' = 'blackhole'
);

INSERT INTO blackhole_sink
SELECT k, s.v, d.v
FROM datagen_source AS s
INNER JOIN odps_dim FOR SYSTEM_TIME AS OF PROCTIME() AS d ON s.k = d.k;

API DataStream

Important
  • Pour utiliser l'API DataStream avec MaxCompute, configurez un connecteur DataStream. Consultez la rubrique Intégrer des connecteurs DataStream.

  • VVR 6.0.6 ou version ultérieure prend en charge le débogage local des programmes DataStream avec le connecteur MaxCompute pendant une durée maximale de 30 minutes. Les sessions dépassant 30 minutes sont interrompues avec une erreur. Consultez la page Déboguer les connecteurs localement.

  • La lecture à partir d'une table Delta MaxCompute (une table créée avec une clé primaire et transactional=true) n'est pas prise en charge.

Déclarez la table MaxCompute à l'aide de SQL, puis accédez-y via l'API Table ou l'API DataStream.

Se connecter à la table source

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tEnv = StreamTableEnvironment.create(env);
tEnv.executeSql(String.join(
    "\n",
    "CREATE TEMPORARY TABLE IF NOT EXISTS odps_source (",
    "  cid VARCHAR,",
    "  rt DOUBLE",
    ") WITH (",
    "  'connector' = 'odps',",
    "  'endpoint' = '<yourEndpointName>',",
    "  'project' = '<yourProjectName>',",
    "  'tableName' = '<yourTableName>',",
    "  'accessId' = '<yourAccessId>',",
    "  'accessKey' = '<yourAccessPassword>',",
    "  'partition' = 'ds=201809*'",
    ")");
DataStream<Row> source = tEnv.toDataStream(tEnv.from("odps_source"));
source.print();
env.execute("odps source");

Se connecter à la destination

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tEnv = StreamTableEnvironment.create(env);
tEnv.executeSql(String.join(
    "\n",
    "CREATE TEMPORARY TABLE IF NOT EXISTS odps_sink (",
    "  cid VARCHAR,",
    "  rt DOUBLE",
    ") WITH (",
    "  'connector' = 'odps',",
    "  'endpoint' = '<yourEndpointName>',",
    "  'project' = '<yourProjectName>',",
    "  'tableName' = '<yourTableName>',",
    "  'accessId' = '<yourAccessId>',",
    "  'accessKey' = '<yourAccessPassword>',",
    "  'partition' = 'ds=20180905'",
    ")");
DataStream<Row> data = env.fromElements(
    Row.of("id0", 3.),
    Row.of("id1", 4.));
tEnv.fromDataStream(data).insertInto("odps_sink").execute();

Dépendance Maven

Ajoutez le connecteur DataStream MaxCompute à votre projet. Toutes les versions sont disponibles dans le dépôt central Maven.

<dependency>
    <groupId>com.alibaba.ververica</groupId>
    <artifactId>ververica-connector-odps</artifactId>
    <version>${vvr-version}</version>
</dependency>

Étapes suivantes