Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Requêtes

Dernière mise à jour :Aug 09, 2026

Realtime Compute for Apache Flink est compatible avec les requêtes SQL natives d'Apache Flink. La grammaire Backus-Naur (BNF) suivante définit le sur-ensemble des requêtes SQL de streaming et par lot prises en charge.

query:
    values
  | WITH withItem [ , withItem ]* query
  | {
        select
      | selectWithoutFrom
      | query UNION [ ALL ] query
      | query EXCEPT query
      | query INTERSECT query
    }
    [ ORDER BY orderItem [, orderItem ]* ]
    [ LIMIT { count | ALL } ]
    [ OFFSET start { ROW | ROWS } ]
    [ FETCH { FIRST | NEXT } [ count ] { ROW | ROWS } ONLY]

withItem:
    name
    [ '(' column [, column ]* ')' ]
    AS '(' query ')'

orderItem:
    expression [ ASC | DESC ]

select:
    SELECT [ ALL | DISTINCT ]
    { * | projectItem [, projectItem ]* }
    FROM tableExpression
    [ WHERE booleanExpression ]
    [ GROUP BY { groupItem [, groupItem ]* } ]
    [ HAVING booleanExpression ]
    [ WINDOW windowName AS windowSpec [, windowName AS windowSpec ]* ]

selectWithoutFrom:
    SELECT [ ALL | DISTINCT ]
    { * | projectItem [, projectItem ]* }

projectItem:
    expression [ [ AS ] columnAlias ]
  | tableAlias . *

tableExpression:
    tableReference [, tableReference ]*
  | tableExpression [ NATURAL ] [ LEFT | RIGHT | FULL ] JOIN tableExpression [ joinCondition ]

joinCondition:
    ON booleanExpression
  | USING '(' column [, column ]* ')'

tableReference:
    tablePrimary
    [ matchRecognize ]
    [ [ AS ] alias [ '(' columnAlias [, columnAlias ]* ')' ] ]

tablePrimary:
    [ TABLE ] tablePath [ dynamicTableOptions ] [systemTimePeriod] [[AS] correlationName]
  | LATERAL TABLE '(' functionName '(' expression [, expression ]* ')' ')'
  | [ LATERAL ] '(' query ')'
  | UNNEST '(' expression ')'

tablePath:
    [ [ catalogName . ] databaseName . ] tableName

systemTimePeriod:
    FOR SYSTEM_TIME AS OF dateTimeExpression

dynamicTableOptions:
    /*+ OPTIONS(key=val [, key=val]*) */

key:
    stringLiteral

val:
    stringLiteral

values:
    VALUES expression [, expression ]*

groupItem:
    expression
  | '(' ')'
  | '(' expression [, expression ]* ')'
  | CUBE '(' expression [, expression ]* ')'
  | ROLLUP '(' expression [, expression ]* ')'
  | GROUPING SETS '(' groupItem [, groupItem ]* ')'

windowRef:
    windowName
  | windowSpec

windowSpec:
    [ windowName ]
    '('
    [ ORDER BY orderItem [, orderItem ]* ]
    [ PARTITION BY expression [, expression ]* ]
    [
        RANGE numericOrIntervalExpression {PRECEDING}
      | ROWS numericExpression {PRECEDING}
    ]
    ')'

matchRecognize:
    MATCH_RECOGNIZE '('
    [ PARTITION BY expression [, expression ]* ]
    [ ORDER BY orderItem [, orderItem ]* ]
    [ MEASURES measureColumn [, measureColumn ]* ]
    [ ONE ROW PER MATCH ]
    [ AFTER MATCH
      ( SKIP TO NEXT ROW
      | SKIP PAST LAST ROW
      | SKIP TO FIRST variable
      | SKIP TO LAST variable
      | SKIP TO variable )
    ]
    PATTERN '(' pattern ')'
    [ WITHIN intervalLiteral ]
    DEFINE variable AS condition [, variable AS condition ]*
    ')'

measureColumn:
    expression AS alias

pattern:
    patternTerm [ '|' patternTerm ]*

patternTerm:
    patternFactor [ patternFactor ]*

patternFactor:
    variable [ patternQuantifier ]

patternQuantifier:
    '*'
  | '*?'
  | '+'
  | '+?'
  | '?'
  | '??'
  | '{' { [ minRepeat ], [ maxRepeat ] } '}' ['?']
  | '{' repeat '}'

Identifiants

Les règles d'identification SQL de Flink diffèrent de Java sur un point important : les caractères non alphanumériques sont autorisés. À tous autres égards, les règles sont identiques à celles de Java :

  • La définition des identifiants respecte la casse, qu'ils soient ou non entourés d'accent grave.

  • La correspondance des identifiants respecte la casse.

Exemple avec un identifiant contenant des caractères non alphanumériques :

SELECT a AS `my field` FROM t

Constantes de chaîne

Les constantes de chaîne doivent utiliser des guillemets simples ('), et non des guillemets doubles (").

SELECT 'Hello World'

Pour échapper un guillemet simple au sein d'une chaîne, doublez-le :

Flink SQL> SELECT 'Hello World', 'It''s me';
+-------------+---------+
|      EXPR$0 |  EXPR$1 |
+-------------+---------+
| Hello World | It's me |
+-------------+---------+
1 row in set

Séquences d'échappement Unicode

Pour inclure des valeurs Unicode dans une constante de chaîne, utilisez le préfixe U& suivi d'une séquence d'échappement.

Méthode Syntaxe Exemple
Échappement par défaut (barre oblique inverse) U&'\<unicode>' SELECT U&'\263A'
Caractère d'échappement personnalisé U&'<char><unicode>' UESCAPE '<char>' SELECT U&'#263A' UESCAPE '#' — utilise # comme caractère d'échappement

Requêtes prises en charge

Le tableau suivant répertorie toutes les requêtes prises en charge par Apache Flink 1.15. Pour consulter la référence d'une autre version de Flink, changez de version sur le site Web d'Apache Flink.

Requête Référence
Indices Indices SQL
Clause WITH Clause WITH
Clauses SELECT et WHERE Clauses SELECT et WHERE
SELECT DISTINCT SELECT DISTINCT
Fonctions de fenêtrage Fonctions de table de fenêtrage (Windowing TVFs)
Agrégation de fenêtre Agrégation de fenêtre
Agrégation de groupe Agrégation de groupe
Agrégation Over Agrégation Over
Jointure Jointures
Jointure de fenêtre Jointure de fenêtre
Opérations sur les ensembles Opérations sur les ensembles
Clause ORDER BY Clause ORDER BY
Clause LIMIT Clause LIMIT
Top-N Top-N
Window Top-N Window Top-N
Déduplication Déduplication
Déduplication de fenêtre Déduplication de fenêtre
Reconnaissance de modèle Reconnaissance de modèle

Exécution des requêtes

En mode streaming, les flux d'entrée se divisent en deux catégories :

  • Les flux sans mise à jour contiennent uniquement des événements de type INSERT.

  • Les flux de mise à jour contiennent d'autres types d'événements. Les sources de capture des données modifiées (CDC) produisent des flux de mise à jour. Certaines opérations Flink, telles que l'agrégation de groupe et Top-N, génèrent également des événements de mise à jour en interne.

La plupart des opérations qui produisent des événements de mise à jour s'appuient sur des opérateurs avec état, qui utilisent un état géré pour suivre les mises à jour. Toutefois, tous les opérateurs avec état n'acceptent pas les flux de mise à jour en entrée. Par exemple, l'agrégation Over et la jointure d'intervalle ne prennent pas en charge les flux de mise à jour en entrée.

Le tableau suivant décrit les caractéristiques d'exécution de chaque requête prise en charge. Ces informations s'appliquent à Ververica Runtime (VVR) 6.0.X et versions ultérieures.

RequêteOpérateur d'exécutionUtilise des données d'étatConsomme des flux de mise à jourGénère des événements de mise à jourNotes
SELECT et WHERECalcNonOuiNon
Lookup JoinLookupJoinNon*OuiNonPour VVR 8.0.1 et versions ultérieures : la définition de table.optimizer.non-deterministic-update.strategy sur TRY_RESOLVE permet la résolution automatique basée sur l'état des problèmes de mise à jour non déterministe. Définissez-la sur IGNORE pour désactiver l'utilisation de l'état. La modification de ce paramètre peut provoquer une incompatibilité et nécessiter la réexécution de la requête.
Fonction de tableCorrelateNonOuiNon
SELECT DISTINCTGroupAggregateOuiOuiOui
Agrégation de groupeGroupAggregate / LocalGroupAggregate / GlobalGroupAggregate / IncrementalGroupAggregateOui*OuiOuiL'opérateur de pré-agrégation LocalGroupAggregate n'utilise pas de données d'état.
Agrégation OverOverAggregateOuiNonNon
Agrégation de fenêtreGroupWindowAggregate / WindowAggregate / LocalWindowAggregate / GlobalWindowAggregateOui*Oui*Non*LocalWindowAggregate n'utilise pas de données d'état. La prise en charge des flux de mise à jour diffère entre VVR et Apache Flink — consultez la section « Comparaison de la prise en charge des flux de mise à jour » de la rubrique Agrégation de fenêtre pour plus de détails. Si la fonctionnalité de déclenchement anticipé ou tardif (expérimentale) est activée, des événements de mise à jour sont générés.
Jointure (jointure régulière)JoinOuiOuiOui*Les jointures externes (LEFT JOIN, RIGHT JOIN, FULL OUTER JOIN) génèrent des événements de mise à jour.
Interval JoinIntervalJoinOuiNonNon
Temporal JoinTemporalJoinOuiOuiNon
Jointure de fenêtreWindowJoinOuiNonNon
Top-NRankOuiOuiOuiTop-N ne prend pas en charge le classement par temps de traitement. Utilisez plutôt des fonctions intégrées telles que CURRENT_TIMESTAMP.
Avertissement

La spécification d'un champ de temps de traitement dans la clause ORDER BY peut entraîner des erreurs de données. Ce problème n'est pas signalé lors des vérifications de syntaxe dans VVR 8.0.7 et versions antérieures.

Window Top-NWindowRankOuiNonNon
DéduplicationDeduplicateOuiNonOui*L'utilisation de la stratégie Deduplicate Keep FirstRow avec le temps de traitement (Proctime) ne génère pas d'événements de mise à jour.
Déduplication de fenêtreWindowDeduplicateOuiNonNon
Remarque

Les opérateurs sans état transmettent les types d'événements sans modification : les événements de sortie sont du même type que les événements d'entrée. Les opérateurs sans état ne génèrent jamais d'événements de mise à jour, que l'entrée soit ou non un flux de mise à jour.