Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Fenêtre Top-N

Dernière mise à jour :Aug 09, 2026

La fenêtre Top-N doit respecter simultanément les règles de modification des fonctions de table (TVF) fenêtrées et des requêtes Top-N. Comme ces deux ensembles de règles s'appliquent, la fenêtre Top-N accepte moins de modifications compatibles qu'une TVF fenêtrée ou une requête Top-N utilisée isolément. Cette rubrique décrit les modifications apportées à une requête de fenêtre Top-N qui sont compatibles avec les données d'état existantes, ainsi que celles qui nécessitent une réinitialisation complète de l'état avant le redémarrage du job.

Pourquoi certaines modifications rompent la compatibilité de l'état

Lorsque vous modifiez une requête de fenêtre Top-N, l'optimiseur de requêtes de Flink peut générer un plan d'exécution différent, ce qui modifie la topologie des opérateurs ou le schéma d'état d'un opérateur intermédiaire. Si le nouveau plan ne correspond plus au point de sauvegarde (savepoint) existant, le job ne peut pas reprendre à partir de son état précédent.

Toute modification qui affecte les attributs de la fenêtre, les champs GROUP BY, les clés de partitionnement, les champs ORDER BY, la valeur de N ou les champs d'agrégation déclenche une replanification et entraîne une incompatibilité d'état.

Structure syntaxique de la fenêtre Top-N

Comprendre quelle clause chaque modification cible vous aide à appliquer les règles de compatibilité ci-dessous. La structure d'une requête de fenêtre Top-N est la suivante :

-- Outer query: select the top-ranked rows
SELECT [column_list]
FROM (
  -- Middle layer: assign row numbers using ROW_NUMBER()
  SELECT [column_list],
    ROW_NUMBER() OVER (
      PARTITION BY window_start, window_end [, partition_key...]
      ORDER BY col [ASC|DESC] [, col [ASC|DESC]...]
    ) AS rk
  FROM (
    -- Inner query: window aggregation using a window TVF
    SELECT [agg_fields], window_start, window_end
    FROM TABLE(TUMBLE(TABLE source_table, DESCRIPTOR(ts), INTERVAL '1' MINUTE))
    GROUP BY [group_keys], window_start, window_end
  )
)
WHERE rk < N;

Les règles de compatibilité ci-dessous correspondent à des clauses spécifiques de cette structure.

Modifications compatibles

Les modifications suivantes sont entièrement compatibles avec les données d'état existantes. Vous pouvez les effectuer sans réinitialiser l'état du job.

Ajouter ou supprimer un champ d'attribut de fenêtre du résultat de la requête

L'ajout ou la suppression de window_start ou de window_end dans la liste SELECT externe n'affecte pas la structure interne de l'état.

-- Original: selects window_start only
SELECT a, b, c, window_start FROM (
  SELECT *,
    ROW_NUMBER() OVER (PARTITION BY b, window_start, window_end ORDER BY c) AS rk
  FROM (
    SELECT a, sum(b) AS b, max(c) AS c, window_start, window_end
    FROM TABLE(tumble(TABLE MyTable, DESCRIPTOR(ts), INTERVAL '1' MINUTE))
    GROUP BY a, window_start, window_end
  )
) WHERE rk < 3;

-- Compatible: add window_end to the output
SELECT a, b, c, window_start, window_end FROM (
  SELECT *,
    ROW_NUMBER() OVER (PARTITION BY b, window_start, window_end ORDER BY c) AS rk
  FROM (
    SELECT a, sum(b) AS b, max(c) AS c, window_start, window_end
    FROM TABLE(tumble(TABLE MyTable, DESCRIPTOR(ts), INTERVAL '1' MINUTE))
    GROUP BY a, window_start, window_end
  )
) WHERE rk < 3;

Inclure ou exclure le champ de position de classement du résultat de la requête

L'inclusion ou l'exclusion du champ rk (la sortie de ROW_NUMBER) dans la clause SELECT externe est entièrement compatible.

-- Original: rk is not in the output
SELECT a, b, c, window_start FROM (
  SELECT *,
    ROW_NUMBER() OVER (PARTITION BY b, window_start, window_end ORDER BY c) AS rk
  FROM (
    SELECT a, sum(b) AS b, max(c) AS c, window_start, window_end
    FROM TABLE(tumble(TABLE MyTable, DESCRIPTOR(ts), INTERVAL '1' MINUTE))
    GROUP BY a, window_start, window_end
  )
) WHERE rk < 3;

-- Compatible: include rk in the output
SELECT a, b, c, window_start, rk FROM (
  SELECT *,
    ROW_NUMBER() OVER (PARTITION BY b, window_start, window_end ORDER BY c) AS rk
  FROM (
    SELECT a, sum(b) AS b, max(c) AS c, window_start, window_end
    FROM TABLE(tumble(TABLE MyTable, DESCRIPTOR(ts), INTERVAL '1' MINUTE))
    GROUP BY a, window_start, window_end
  )
) WHERE rk < 3;

Réorganiser les clés de partitionnement dans la clause OVER

La modification de l'ordre des clés de partitionnement dans PARTITION BY ne change ni la logique de partitionnement ni la structure de l'état.

-- Original: PARTITION BY a, b, window_start, window_end
SELECT a, b, c, window_start FROM (
  SELECT *,
    ROW_NUMBER() OVER (PARTITION BY a, b, window_start, window_end ORDER BY c) AS rk
  FROM (
    SELECT a, sum(b) AS b, max(c) AS c, window_start, window_end
    FROM TABLE(tumble(TABLE MyTable, DESCRIPTOR(ts), INTERVAL '1' MINUTE))
    GROUP BY a, b, window_start, window_end
  )
) WHERE rk < 3;

-- Compatible: reorder to PARTITION BY b, a, window_start, window_end
SELECT a, b, c, window_start FROM (
  SELECT *,
    ROW_NUMBER() OVER (PARTITION BY b, a, window_start, window_end ORDER BY c) AS rk
  FROM (
    SELECT a, sum(b) AS b, max(c) AS c, window_start, window_end
    FROM TABLE(tumble(TABLE MyTable, DESCRIPTOR(ts), INTERVAL '1' MINUTE))
    GROUP BY a, b, window_start, window_end
  )
) WHERE rk < 3;

Modifications incompatibles

Les modifications suivantes sont incompatibles avec les données d'état existantes. Après avoir effectué l'une de ces modifications, vous devez réinitialiser l'état du job avant de le redémarrer.

Modification Détails
Modifier un attribut de fenêtre Inclut le type de fenêtre, la taille de la fenêtre ou l'attribut temporel. Pour des exemples, consultez Modifications entraînant une incompatibilité totale.
Ajouter, supprimer ou modifier des champs dans la clause GROUP BY, ou changer leur logique de calcul Affecte l'état d'agrégation interne. Pour des exemples, consultez Modifications entraînant une incompatibilité totale.
Ajouter, supprimer ou modifier un champ d'agrégation, ou changer l'entrée de la requête Top-N Modifie le schéma des données transmises à la couche ROW_NUMBER. Voir l'exemple ci-dessous.
Ajouter, supprimer ou modifier des clés de partitionnement, ou changer la logique de calcul des champs de clé de partitionnement Affecte l'état ROW_NUMBER de la couche intermédiaire. Pour des exemples, consultez Modifications incompatibles.
Modifier les champs ou l'ordre dans la clause ORDER BY Affecte la logique de classement et l'état. Pour des exemples, consultez Modifications incompatibles.
Modifier la valeur de N N spécifie le nombre de résultats les mieux classés à renvoyer. Pour des exemples, consultez Modifications incompatibles.
Réorganiser uniquement les champs liés à la TVF de fenêtre dans la clause GROUP BY Bien que les champs GROUP BY non liés à la fenêtre restent identiques, la réorganisation de window_start et de window_end modifie le résultat du classement basé sur la fenêtre et est incompatible. Voir l'exemple ci-dessous.
Réorganiser uniquement les champs non liés à la fenêtre dans la clause GROUP BY La réorganisation des champs non liés à la fenêtre (par exemple, l'échange de a et de b) modifie l'état d'agrégation interne et est incompatible. Voir l'exemple ci-dessous.

Exemple : ajout d'un champ d'agrégation

L'ajout de min(d) AS d à l'agrégation interne modifie le schéma d'entrée de la requête Top-N et rompt la compatibilité de l'état.

-- Original
SELECT a, b, c, window_start FROM (
  SELECT *,
    ROW_NUMBER() OVER (PARTITION BY b, window_start, window_end ORDER BY c) AS rk
  FROM (
    SELECT a, sum(b) AS b, max(c) AS c, window_start, window_end
    FROM TABLE(tumble(TABLE MyTable, DESCRIPTOR(ts), INTERVAL '1' MINUTE))
    GROUP BY a, window_start, window_end
  )
) WHERE rk < 3;

-- Incompatible: adds min(d) as d, changing the Top-N input
SELECT a, b, c, d, window_start FROM (
  SELECT *,
    ROW_NUMBER() OVER (PARTITION BY b, window_start, window_end ORDER BY c) AS rk
  FROM (
    SELECT a, sum(b) AS b, max(c) AS c, min(d) AS d, window_start, window_end
    FROM TABLE(tumble(TABLE MyTable, DESCRIPTOR(ts), INTERVAL '1' MINUTE))
    GROUP BY a, window_start, window_end
  )
) WHERE rk < 3;

Exemple : réorganisation des champs GROUP BY

Les deux modifications suivantes sont incompatibles, même si aucun champ n'est ajouté ou supprimé.

-- Original
SELECT a, b, c, window_start FROM (
  SELECT *,
    ROW_NUMBER() OVER (PARTITION BY b, window_start, window_end ORDER BY c) AS rk
  FROM (
    SELECT a, sum(b) AS b, max(c) AS c, window_start, window_end
    FROM TABLE(tumble(TABLE MyTable, DESCRIPTOR(ts), INTERVAL '1' MINUTE))
    GROUP BY a, b, window_start, window_end
  )
) WHERE rk < 3;

-- Incompatible: reorders window TVF fields (window_end before window_start)
-- This changes the window-based ranking result.
SELECT a, b, c, window_start FROM (
  SELECT *,
    ROW_NUMBER() OVER (PARTITION BY b, window_start, window_end ORDER BY c) AS rk
  FROM (
    SELECT a, sum(b) AS b, max(c) AS c, window_start, window_end
    FROM TABLE(tumble(TABLE MyTable, DESCRIPTOR(ts), INTERVAL '1' MINUTE))
    GROUP BY a, b, window_end, window_start
  )
) WHERE rk < 3;

-- Incompatible: reorders non-window fields (b before a)
SELECT a, b, c, window_start FROM (
  SELECT *,
    ROW_NUMBER() OVER (PARTITION BY b, window_start, window_end ORDER BY c) AS rk
  FROM (
    SELECT a, sum(b) AS b, max(c) AS c, window_start, window_end
    FROM TABLE(tumble(TABLE MyTable, DESCRIPTOR(ts), INTERVAL '1' MINUTE))
    GROUP BY b, a, window_start, window_end
  )
) WHERE rk < 3;

Étapes suivantes