Lorsque plusieurs flux Flink mettent à jour différentes colonnes d'une même ligne, une opération UPSERT standard écrase la ligne entière. Le dernier flux efface ainsi silencieusement les modifications du précédent. Les mises à jour partielles de colonnes permettent à chaque flux d'écrire uniquement dans les colonnes qui lui sont attribuées, préservant les mises à jour concurrentes de flux indépendants au sein d'une même ligne sans conflit.
Quand utiliser cette fonctionnalité
Flux indépendants, même ligne : deux flux mettent à jour simultanément différentes colonnes d'un même enregistrement. Sans les mises à jour partielles de colonnes, l'écriture ultérieure écrase les colonnes définies par l'écriture précédente, entraînant une perte de données.
Écritures éparses : un flux fournit des valeurs uniquement pour un sous-ensemble de colonnes. Sans les mises à jour partielles de colonnes, les colonnes non spécifiées passent à NULL, détruisant les données existantes.
Fonctionnement
L'opération UPSERT combine INSERT et UPDATE en une seule action. Chaque enregistrement traité par UPSERT doit inclure les colonnes de clé primaire.
Si aucun enregistrement correspondant à la clé primaire donnée n'existe, UPSERT insère un nouvel enregistrement.
Si un enregistrement correspondant à la clé primaire donnée existe déjà, UPSERT le met à jour avec les nouvelles valeurs.
Dans le traitement de flux impliquant plusieurs jointures, différents flux peuvent mettre à jour diverses colonnes d'une même ligne. L'UPSERT standard remplace la ligne entière ; les mises à jour d'un flux peuvent donc écraser celles d'un autre. Les mises à jour partielles de colonnes résolvent ce problème en limitant l'écriture de chaque flux aux seules colonnes qui lui sont attribuées.
Modes de mise à jour partielle des colonnes
MaxCompute propose deux modes de mise à jour partielle des colonnes pour le connecteur Flink.
|
Mode |
Fonctionnement |
Cas d'utilisation |
|
Mode dynamique |
Détecte automatiquement les colonnes non NULL dans chaque enregistrement et met à jour uniquement celles-ci. Les colonnes NULL restent inchangées. |
L'ensemble des colonnes est inconnu à l'avance ou varie selon les enregistrements |
|
Mode statique |
Met à jour uniquement les colonnes que vous spécifiez explicitement. Toutes les autres colonnes conservent leurs valeurs existantes. |
L'ensemble des colonnes est fixe et connu à l'avance |
Le tableau suivant illustre comment la même séquence d'écritures produit des résultats différents selon le mode utilisé. La colonne a constitue la clé primaire.
|
Mode |
Données initiales |
Après écriture (a, b, c) |
Après écriture (a, d, NULL) |
Après écriture (a, NULL, e) |
|
Mode standard |
(null, null, null) |
(a, b, c) |
(a, d, null) |
(a, null, e) |
|
Mode dynamique |
(null, null, null) |
(a, b, c) |
(a, d, c) |
(a, d, e) |
|
Mode statique (deuxième colonne spécifiée) |
(null, null, null) |
(a, b, null) |
(a, d, null) |
(a, null, null) |
Prérequis
Avant de commencer, assurez-vous de disposer des éléments suivants :
Un projet MaxCompute avec une table Delta activée pour les mises à jour partielles de colonnes
Un connecteur Flink configuré pour se connecter à MaxCompute
Activer les mises à jour partielles de colonnes sur une table Delta
Définissez acid.partial.fields.update.enable=true dans TBLPROPERTIES lors de la création de la table Delta. Pour obtenir la liste complète des paramètres des tables Delta, consultez la rubrique Paramètres des tables Delta.
CREATE TABLE IF NOT EXISTS partial_upsert_test
(pk INT NOT NULL,
c1 STRING,
c2 STRING,
c3 STRING,
PRIMARY KEY(pk))
TBLPROPERTIES('transactional'='true', 'acid.partial.fields.update.enable'='true');
Configurer le connecteur Flink pour les mises à jour partielles de colonnes
Paramètres
|
Paramètre |
Description |
|
|
Active les mises à jour partielles de colonnes. Lorsque |
|
|
Spécifie les colonnes à mettre à jour (mode statique). Les colonnes de clé primaire sont toujours incluses automatiquement. |
Le paramètre upsert.partial-column.name doit utiliser les noms de colonnes de la table MaxCompute, et non ceux de la table interne Flink.
Les noms des colonnes de clé de partition ne peuvent pas être ajoutés au paramètre upsert.partial-column.name .
Exemple de mode dynamique
La définition de table Flink suivante active les mises à jour partielles de colonnes en mode dynamique. Le système détecte automatiquement les colonnes contenant des valeurs non NULL et met à jour uniquement celles-ci.
CREATE TABLE partialtable (
pk INT,
c1 STRING,
c2 STRING,
c3 STRING,
PRIMARY KEY(pk) NOT ENFORCED
) WITH (
'connector' = 'maxcompute',
'odps.end.point' = 'https://service.cn-hangzhou-vpc.maxcompute.aliyun-inc.com/api', -- VPC endpoint
'odps.project.name' = 'project_name',
'odps.namespace.schema' = 'true', -- Enable the three-layer model
'table.name' = 'project.schema.tablename',
'sink.operation' = 'upsert',
'upsert.write.bucket.num' = '1',
'upsert.partial-column.enable' = 'true', -- Enable dynamic mode
'odps.access.id' = '<your-access-key-id>',
'odps.access.key' = '<your-access-key-secret>'
);
L'exemple ci-dessous montre comment le mode dynamique préserve les valeurs existantes lors de trois écritures séquentielles.
Étape 1 : Insérez l'enregistrement initial. Données après cette étape : [1, a, b, c].
INSERT INTO partialtable VALUES (1, 'a', 'b', 'c');
Étape 2 : Mettez à jour uniquement c2. Étant donné que c1 et c3 ne sont pas spécifiés, ils conservent leurs valeurs. Données après cette étape : [1, a, d, c].
INSERT INTO partialtable(pk, c2) VALUES (1, 'd');
Étape 3 : Mettez à jour uniquement c3. Étant donné que c1 et c2 ne sont pas spécifiés, ils conservent leurs valeurs. Données après cette étape : [1, a, d, e].
INSERT INTO partialtable(pk, c3) VALUES (1, 'e');
Exemple de mode statique
La définition de table Flink suivante limite les mises à jour à la colonne c2. Toutes les écritures dans cette table affectent c2 quelles que soient les autres valeurs fournies.
CREATE TABLE PartialTable2 (
pk INT,
c1 STRING,
c2 STRING,
c3 STRING,
PRIMARY KEY(pk) NOT ENFORCED
) WITH (
'connector' = 'maxcompute',
'odps.end.point' = 'https://service.cn-hangzhou-vpc.maxcompute.aliyun-inc.com/api', -- VPC endpoint
'odps.project.name' = 'project_name',
'odps.namespace.schema' = 'true', -- Enable the three-layer model
'table.name' = 'project.schema.tablename',
'sink.operation' = 'upsert',
'upsert.write.bucket.num' = '1',
'upsert.partial-column.enable' = 'true',
'upsert.partial-column.name' = 'c2', -- Only update c2; primary key is included by default
'odps.access.id' = '<your-access-key-id>',
'odps.access.key' = '<your-access-key-secret>'
);
Étapes suivantes
Pour en savoir plus sur l'écriture de données dans des tables Delta avec Flink, consultez la documentation Utilisation de Flink pour écrire des données dans une table Delta.
Pour obtenir une référence complète sur la fonctionnalité de mise à jour partielle des colonnes dans les tables Delta, consultez la documentation Mise à jour des données dans des colonnes spécifiques.