Vous pouvez utiliser le connecteur MaxCompute Flink pour écrire les données de Flink dans des tables standard et delta de MaxCompute, ce qui simplifie l'ingestion des données. Cette rubrique décrit les fonctionnalités du connecteur et présente les procédures d'écriture des données.
Contexte
-
Modes d'écriture pris en charge
Le connecteur Flink prend en charge deux modes d'écriture :
upsertetinsert. En modeupsert, vous pouvez regrouper les flux de données de l'une des manières suivantes :Regroupement par clé primaire
-
Regroupement par champ de partition
Bien que cette méthode convienne à un grand nombre de partitions, le regroupement par champ de partition peut entraîner une répartition inégale des données (data skew).
Pour la procédure d'écriture
upsertet les paramètres recommandés, consultez la rubrique Ingestion de données en temps réel dans un entrepôt de données.Spécifiez le mode d'écriture à l'aide des paramètres du connecteur Flink. Pour obtenir la liste complète des paramètres du connecteur, consultez la section Annexe : Paramètres du connecteur Flink.
Définissez l'intervalle de point de contrôle (checkpoint) pour les travaux d'écriture Flink upsert sur au moins 3 minutes. Des intervalles plus courts peuvent réduire l'efficacité de l'écriture et générer un grand nombre de petits fichiers.
-
Le tableau suivant établit la correspondance entre les types de données de Realtime Compute for Apache Flink et ceux de MaxCompute.
Type de données Flink
Type de données MaxCompute
CHAR(p)
CHAR(p)
VARCHAR(p)
VARCHAR(p)
STRING
STRING
BOOLEAN
BOOLEAN
TINYINT
TINYINT
SMALLINT
SMALLINT
INT
INT
BIGINT
BIGINT
FLOAT
FLOAT
DOUBLE
DOUBLE
DECIMAL(p, s)
DECIMAL(p, s)
DATE
DATE
TIMESTAMP(9) WITHOUT TIME ZONE, TIMESTAMP_LTZ(9)
TIMESTAMP
TIMESTAMP(3) WITHOUT TIME ZONE, TIMESTAMP_LTZ(3)
DATETIME
BYTES
BINARY
ARRAY<T>
ARRAY<T>
MAP<K, V>
MAP<K, V>
ROW
STRUCT
RemarqueLe type de données TIMESTAMP de Flink n'inclut pas d'informations de fuseau horaire, contrairement au type de données TIMESTAMP de MaxCompute. Cette différence peut entraîner un décalage horaire de 8 heures. Pour aligner les horodatages, utilisez TIMESTAMP_LTZ(9).
-- Flink SQL CREATE TEMPORARY TABLE odps_source( id BIGINT NOT NULL COMMENT 'ID', created_time TIMESTAMP NOT NULL COMMENT 'Creation time', updated_time TIMESTAMP_LTZ(9) NOT NULL COMMENT 'Update time', PRIMARY KEY (id) NOT ENFORCED ) WITH ( 'connector' = 'maxcompute', ... );
Écrire des données depuis un cluster Flink autonome
-
Prérequis : Créez une table MaxCompute.
Vous devez d'abord créer une table MaxCompute vers laquelle Flink écrira les données. L'exemple suivant illustre ce processus en créant deux tables (une table delta non partitionnée et une table partitionnée). Pour plus d'informations sur la configuration des propriétés des tables, consultez les rubriques relatives aux paramètres des tables delta.
-- Create a non-partitioned delta table. CREATE TABLE mf_flink_tt ( id BIGINT not null, name STRING, age INT, status BOOLEAN, primary key (id) ) tblproperties ("transactional"="true", "write.bucket.num" = "64", "acid.data.retain.hours"="12") ; --Create a partitioned delta table. CREATE TABLE mf_flink_tt_part ( id BIGINT not null, name STRING, age INT, status BOOLEAN, primary key (id) ) partitioned by (dd string, hh string) tblproperties ("transactional"="true", "write.bucket.num" = "64", "acid.data.retain.hours"="12") ; -
Configurez un cluster Flink open source. Le connecteur prend en charge les versions 1.13, 1.15, 1.16 et 1.17 de Flink. Sélectionnez le connecteur Flink correspondant à votre version de Flink :
RemarqueLe connecteur Flink pour Flink 1.16 est compatible avec Flink 1.17.
Cette rubrique utilise le connecteur Flink pour Flink 1.13 à titre d'exemple. Téléchargez et décompressez le package.
-
Téléchargez le connecteur Flink et ajoutez-le au package de votre cluster Flink.
Téléchargez le package JAR du connecteur Flink dans votre environnement local.
-
Ajoutez le package JAR du connecteur Flink au répertoire lib du package d'installation Flink décompressé.
mv flink-connector-odps-1.13-shaded.jar $FLINK_HOME/lib/flink-connector-odps-1.13-shaded.jar
-
Démarrez le service d'instance Flink.
cd $FLINK_HOME/bin ./start-cluster.sh -
Démarrez le client Flink.
cd $FLINK_HOME/bin ./sql-client.sh -
Créez une table Flink et configurez les paramètres du connecteur Flink.
Vous pouvez créer une table Flink et configurer ses paramètres soit via Flink SQL, soit via l'API DataStream. Les sections suivantes fournissent des exemples essentiels pour ces deux approches.
Flink SQL
-
Dans l'éditeur Flink SQL, exécutez les commandes suivantes pour créer une table et configurer les paramètres.
-- Register a non-partitioned table in Flink SQL. CREATE TABLE mf_flink ( id BIGINT, name STRING, age INT, status BOOLEAN, PRIMARY KEY(id) NOT ENFORCED ) WITH ( 'connector' = 'maxcompute', 'table.name' = 'mf_flink_tt', 'sink.operation' = 'upsert', 'odps.access.id'='LTAI****************', 'odps.access.key'='********************', 'odps.end.point'='http://service.cn-beijing.maxcompute.aliyun.com/api', 'odps.project.name'='mf_mc_bj' ); -- Register a partitioned table in Flink SQL. CREATE TABLE mf_flink_part ( id BIGINT, name STRING, age INT, status BOOLEAN, dd STRING, hh STRING, PRIMARY KEY(id) NOT ENFORCED ) PARTITIONED BY (`dd`,`hh`) WITH ( 'connector' = 'maxcompute', 'table.name' = 'mf_flink_tt_part', 'sink.operation' = 'upsert', 'odps.access.id'='LTAI****************', 'odps.access.key'='********************', 'odps.end.point'='http://service.cn-beijing.maxcompute.aliyun.com/api', 'odps.project.name'='mf_mc_bj' ); -
Écrivez des données dans la table Flink et interrogez la table MaxCompute pour vérifier le résultat.
-- Insert data into the non-partitioned table in the Flink SQL client. INSERT INTO mf_flink VALUES (1,'Danny',27, false); -- Query result in MaxCompute. SELECT * FROM mf_flink_tt; +------------+------+------+--------+ | id | name | age | status | +------------+------+------+--------+ | 1 | Danny | 27 | false | +------------+------+------+--------+ -- Insert data into the non-partitioned table in the Flink SQL client to update the record. INSERT INTO mf_flink VALUES (1,'Danny',28, false); -- Query result in MaxCompute. SELECT * FROM mf_flink_tt; +------------+------+------+--------+ | id | name | age | status | +------------+------+------+--------+ | 1 | Danny | 28 | false | +------------+------+------+--------+ -- Insert data into the partitioned table in the Flink SQL client. INSERT INTO mf_flink_part VALUES (1,'Danny',27, false, '01','01'); -- Query result in MaxCompute. SELECT * FROM mf_flink_tt_part WHERE dd=01 AND hh=01; +------------+------+------+--------+----+----+ | id | name | age | status | dd | hh | +------------+------+------+--------+----+----+ | 1 | Danny | 27 | false | 01 | 01 | +------------+------+------+--------+----+----+ -- Insert data into the partitioned table in the Flink SQL client to update the record. INSERT INTO mf_flink_part VALUES (1,'Danny',30, false, '01','01'); -- Query result in MaxCompute. SELECT * FROM mf_flink_tt_part WHERE dd=01 AND hh=01; +------------+------+------+--------+----+----+ | id | name | age | status | dd | hh | +------------+------+------+--------+----+----+ | 1 | Danny | 30 | false | 01 | 01 | +------------+------+------+--------+----+----+
DataStream API
-
Pour utiliser l'API DataStream, ajoutez la dépendance suivante.
<dependency> <groupId>com.aliyun.odps</groupId> <artifactId>flink-connector-maxcompute</artifactId> <version>xxx</version> <scope>system</scope> <systemPath>${mvn_project.basedir}/lib/flink-connector-maxcompute-xxx-shaded.jar</systemPath> </dependency>RemarqueRemplacez « xxx » par le numéro de version réel.
-
L'exemple de code suivant montre comment créer une table et configurer les paramètres.
package com.aliyun.odps.flink.examples; import org.apache.flink.configuration.Configuration; import org.apache.flink.odps.table.OdpsOptions; import org.apache.flink.odps.util.OdpsConf; import org.apache.flink.odps.util.OdpsPipeline; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.table.api.Table; import org.apache.flink.table.api.bridge.java.StreamTableEnvironment; import org.apache.flink.table.data.RowData; public class Examples { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.enableCheckpointing(120 * 1000); StreamTableEnvironment streamTableEnvironment = StreamTableEnvironment.create(env); Table source = streamTableEnvironment.sqlQuery("SELECT * FROM source_table"); DataStream<RowData> input = streamTableEnvironment.toAppendStream(source, RowData.class); Configuration config = new Configuration(); config.set(OdpsOptions.SINK_OPERATION, "upsert"); config.set(OdpsOptions.UPSERT_COMMIT_THREAD_NUM, 8); config.set(OdpsOptions.UPSERT_MAJOR_COMPACT_MIN_COMMITS, 100); OdpsConf odpsConfig = new OdpsConf("accessid", "accesskey", "endpoint", "project", "tunnel endpoint"); OdpsPipeline.Builder builder = OdpsPipeline.builder(); builder.projectName("sql2_isolation_2a") .tableName("user_ledger_portfolio") .partition("") .configuration(config) .odpsConf(odpsConfig) .sink(input, false); env.execute(); } }
-
Écrire des données depuis Fully Managed Flink
-
Prérequis : Créez une table MaxCompute.
Vous devez créer une table MaxCompute cible pour les données Flink. L'exemple suivant montre comment créer une table delta.
SET odps.sql.type.system.odps2=true; DROP TABLE mf_flink_upsert; CREATE TABLE mf_flink_upsert ( c1 int not null, c2 string, gt timestamp, primary key (c1) ) PARTITIONED BY (ds string) tblproperties ("transactional"="true", "write.bucket.num" = "64", "acid.data.retain.hours"="12") ; Le connecteur Flink est préchargé sur Fully Managed Flink, aucune installation manuelle n'est donc requise. Vous pouvez consulter les détails du connecteur dans la console Realtime Compute for Apache Flink.
-
Créez une table Flink, construisez des données Flink en temps réel à l'aide d'un travail Flink SQL, puis déployez le travail après son développement.
Sur la page de développement des travaux Flink, créez et modifiez un travail SQL. L'exemple suivant définit une table source qui génère des données aléatoires, une table de résultats qui se connecte à MaxCompute et une instruction INSERT pour transférer les données. Pour plus d'informations sur le développement d'un travail SQL, consultez la carte de développement des travaux.
-- Create a Flink source table. CREATE TEMPORARY TABLE fake_src_table ( c1 int, c2 VARCHAR, gt AS CURRENT_TIMESTAMP ) WITH ( 'connector' = 'faker', 'fields.c2.expression' = '#{superhero.name}', 'rows-per-second' = '100', 'fields.c1.expression' = '#{number.numberBetween ''0'',''1000''}' ); -- Create a temporary result table in Flink. CREATE TEMPORARY TABLE test_c_d_g ( c1 int, c2 VARCHAR, gt TIMESTAMP, ds varchar, PRIMARY KEY(c1) NOT ENFORCED ) PARTITIONED BY(ds) WITH ( 'connector' = 'maxcompute', 'table.name' = 'mf_flink_upsert', 'sink.operation' = 'upsert', 'odps.access.id'='LTAI****************', 'odps.access.key'='********************', 'odps.end.point'='http://service.cn-beijing.maxcompute.aliyun.com/api', 'odps.project.name'='mf_mc_bj', 'upsert.write.bucket.num'='64' ); -- Flink computation logic. INSERT INTO test_c_d_g SELECT c1 AS c1, c2 AS c2, gt AS gt, date_format(gt, 'yyyyMMddHH') AS ds FROM fake_src_table;Paramètres :
odps.end.point: Utilisez le point de terminaison du réseau interne de la région correspondante.upsert.write.bucket.num: Cette valeur doit être identique à la valeur de la propriété write.bucket.num de la table delta créée dans MaxCompute. -
Interrogez la table MaxCompute pour vérifier que les données ont bien été écrites.
SELECT * FROM mf_flink_upsert WHERE ds=2023061517; -- Your results may differ because the source data is randomly generated. +------+----+------+----+ | c1 | c2 | gt | ds | +------+----+------+----+ | 0 | Skaar | 2023-06-16 01:59:41.116 | 2023061517 | | 21 | Supah Century | 2023-06-16 01:59:59.117 | 2023061517 | | 104 | Dark Gorilla Grodd | 2023-06-16 01:59:57.117 | 2023061517 | | 126 | Leader | 2023-06-16 01:59:39.116 | 2023061517 |
Annexe : Paramètres du connecteur Flink
-
Paramètres de base
Paramètre
Obligatoire
Valeur par défaut
Description
connector
Oui
—
Définissez le type de connecteur sur
MaxCompute.odps.project.name
Oui
—
Le nom du projet MaxCompute.
odps.access.id
Oui
—
L'ID AccessKey de votre compte. Consultez la page Paires de clés d'accès.
odps.access.key
Oui
—
La clé secrète AccessKey de votre compte. Consultez la page Paires de clés d'accès.
odps.end.point
Oui
—
Le point de terminaison MaxCompute. Pour obtenir la liste des points de terminaison par région, consultez la rubrique Points de terminaison.
odps.tunnel.end.point
Non
—
Le point de terminaison public du service Tunnel. Par défaut, les requêtes sont automatiquement routées vers le point de terminaison Tunnel approprié. Définissez ce paramètre pour utiliser un point de terminaison spécifique et désactiver le routage automatique.
Pour plus d'informations sur les points de terminaison Tunnel selon les régions et les réseaux, consultez la rubrique Points de terminaison.
odps.tunnel.quota.name
Non
—
Le nom du quota Tunnel utilisé pour accéder à MaxCompute.
table.name
Oui
—
Le nom de la table MaxCompute au format
[project.][schema.]table.odps.namespace.schema
Non
false
Indique s'il faut utiliser le modèle à trois couches. Pour plus d'informations sur le modèle à trois couches, consultez la rubrique Opérations sur les schémas.
sink.operation
Oui
insert
Le type d'écriture. Les valeurs valides sont
insertouupsert.RemarqueLe mode
upsertest pris en charge uniquement pour les tables delta MaxCompute.sink.parallelism
Non
—
Le parallélisme d'écriture. Si ce paramètre n'est pas défini, il prend par défaut la valeur du parallélisme de la source en amont.
RemarqueAssurez-vous que la propriété de table
write.bucket.numest un multiple entier de la valeur de configuration pour optimiser les performances d'écriture et maximiser l'économie de mémoire sur le nœud Sink.sink.meta.cache.time
Non
400
La taille du cache des métadonnées.
sink.meta.cache.expire.time
Non
1200
Le délai d'expiration du cache pour les métadonnées, en secondes.
sink.coordinator.enable
Non
true
Indique s'il faut activer le mode coordinateur.
-
Paramètres de partition
Paramètre
Obligatoire
Valeur par défaut
Description
sink.partition
Non
—
Le nom de la partition vers laquelle les données sont écrites.
Si vous utilisez le partitionnement dynamique, ce paramètre spécifie le nom de la partition parente des partitions dynamiques.
sink.partition.default-value
Non
__DEFAULT_PARTITION__
Le nom de partition par défaut lors de l'utilisation du partitionnement dynamique.
sink.dynamic-partition.limit
Non
100
Le nombre maximal de partitions pouvant être écrites simultanément lors d'un seul point de contrôle (checkpoint) pendant le partitionnement dynamique.
RemarqueUne augmentation significative de cette valeur peut provoquer des erreurs d'épuisement de la mémoire (OOM) sur le nœud sink. Si le nombre de partitions simultanées dépasse cette limite, le travail échouera.
sink.group-partition.enable
Non
false
Indique s'il faut effectuer un regroupement par partition lors de l'utilisation du partitionnement dynamique.
sink.partition.assigner.class
Non
—
La classe d'implémentation
PartitionAssigner. -
Paramètres d'écriture en mode FileCached
Utilisez le mode de cache de fichier pour les travaux comportant un grand nombre de partitions dynamiques. Les paramètres suivants configurent ce mode.
Paramètre
Obligatoire
Valeur par défaut
Description
sink.file-cached.enable
Non
false
Active le mode FileCached. Recommandé pour les travaux comportant un grand nombre de partitions dynamiques.
false : Le mode FileCached est désactivé.
true : Le mode FileCached est activé.
RemarqueLorsque le nombre de partitions dynamiques est élevé, vous pouvez utiliser le mode de cache de fichier.
sink.file-cached.tmp.dirs
Non
./local
Le répertoire de cache de fichier par défaut en mode FileCached.
sink.file-cached.writer.num
Non
16
Le nombre de threads de téléchargement de données simultanés pour une tâche unique en mode FileCached.
RemarqueN'augmentez pas significativement la valeur de ce paramètre. Si un nombre excessif de partitions est écrit simultanément, des erreurs OOM sont susceptibles de se produire.
sink.bucket.check-interval
Non
60000
L'intervalle de vérification de la taille des fichiers en mode FileCached. Unité : millisecondes (ms).
sink.file-cached.rolling.max-size
Non
16 M
La taille maximale d'un fichier de cache unique.
Lorsqu'un fichier dépasse cette taille, il est téléchargé.
sink.file-cached.memory
Non
64 M
La taille maximale de la mémoire hors tas (off-heap) utilisée pour l'écriture de fichiers en mode FileCached.
sink.file-cached.memory.segment-size
Non
128 KB
La taille du tampon utilisée pour l'écriture de fichiers en mode FileCached.
sink.file-cached.flush.always
Non
true
Indique s'il faut utiliser le cache pour l'écriture de fichiers en mode FileCached.
sink.file-cached.write.max-retries
Non
3
Le nombre de tentatives de téléchargement des données en mode FileCached.
-
Paramètres d'écriture
InsertouUpsertParamètres d'écriture Upsert
Paramètre
Obligatoire
Valeur par défaut
Description
upsert.writer.max-retries
Non
3
Le nombre de tentatives après l'échec d'écriture des données dans un bucket par un writer upsert.
upsert.writer.buffer-size
Non
64 MB
La taille du cache pour un writer upsert unique dans Flink.
RemarqueLorsque la somme des tailles de tampon de tous les buckets atteint le seuil prédéfini, le système déclenche automatiquement une opération de vidage (flush) pour mettre à jour les données sur le serveur.
Un writer upsert écrit des données dans plusieurs buckets simultanément. Il est recommandé d'augmenter la valeur de ce paramètre pour améliorer l'efficacité de l'écriture.
Si les données sont écrites dans un grand nombre de partitions, des erreurs OOM peuvent se produire. Dans ce cas, vous pouvez réduire la valeur de ce paramètre.
upsert.writer.bucket.buffer-size
Non
1 MB
La taille du cache pour un bucket unique dans Flink. Si les ressources mémoire du serveur Flink sont insuffisantes, vous pouvez réduire la valeur de ce paramètre.
upsert.write.bucket.num
Oui
—
Le nombre de buckets pour la table de destination doit être identique à la valeur de
write.bucket.num.upsert.write.slot-num
Non
1
Le nombre de slots Tunnel utilisés par une session unique.
upsert.commit.max-retries
Non
3
Le nombre de tentatives pour la validation (commit) d'une session upsert.
upsert.commit.thread-num
Non
16
Le parallélisme de la validation (commit) d'une session upsert.
Ne définissez pas cette valeur trop élevée. Un nombre élevé de validations simultanées entraîne une augmentation de la consommation de ressources, ce qui peut provoquer des problèmes de performance ou une consommation excessive de ressources.
upsert.major-compact.min-commits
Non
100
Le nombre minimal de validations requis pour déclencher une compaction majeure.
upsert.commit.timeout
Non
600
Le délai d'expiration pour la validation (commit) d'une session upsert. Unité : secondes (s).
upsert.major-compact.enable
Non
false
Indique s'il faut activer la compaction majeure.
upsert.flush.concurrent
Non
2
Le nombre maximal de buckets vers lesquels les données peuvent être écrites simultanément dans une partition unique.
RemarqueLors du vidage (flush) des données d'un bucket, un slot Tunnel est occupé.
RemarquePour plus d'informations sur les configurations de paramètres recommandées pour les écritures upsert, consultez la rubrique Configurations de paramètres recommandées pour les écritures upsert.
Paramètres d'écriture Insert
Paramètre
Obligatoire
Valeur par défaut
Description
insert.commit.thread-num
Non
16
Le parallélisme d'une session de validation.
insert.arrow-writer.enable
Non
false
Indique s'il faut utiliser le format Arrow.
insert.arrow-writer.batch-size
Non
512
Le nombre maximal de lignes dans un lot Arrow.
insert.arrow-writer.flush-interval
Non
100000
L'intervalle de vidage du writer. Unité : millisecondes (ms).
insert.writer.buffer-size
Non
64 MB
La taille du cache du writer mis en tampon.