Tous les produits
Search
Centre de documentation

MaxCompute:Near real-time partial column updates to Delta tables with Flink

Dernière mise à jour :Aug 10, 2026

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

upsert.partial-column.enable

Active les mises à jour partielles de colonnes. Lorsque upsert.partial-column.name n'est pas défini, le mode dynamique est utilisé.

upsert.partial-column.name

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