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. |
Prérequis
Avant de commencer, assurez-vous de disposer des éléments suivants :
Une table MaxCompute. Consultez la rubrique Créer une table.
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ètrestartPartition.À 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écessitentcache=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'indicationSHUFFLE_HASHpour 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 |
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 :
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
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>