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ête | Opérateur d'exécution | Utilise des données d'état | Consomme des flux de mise à jour | Génère des événements de mise à jour | Notes |
|---|---|---|---|---|---|
| SELECT et WHERE | Calc | Non | Oui | Non | — |
| Lookup Join | LookupJoin | Non* | Oui | Non | Pour 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 table | Correlate | Non | Oui | Non | — |
| SELECT DISTINCT | GroupAggregate | Oui | Oui | Oui | — |
| Agrégation de groupe | GroupAggregate / LocalGroupAggregate / GlobalGroupAggregate / IncrementalGroupAggregate | Oui* | Oui | Oui | L'opérateur de pré-agrégation LocalGroupAggregate n'utilise pas de données d'état. |
| Agrégation Over | OverAggregate | Oui | Non | Non | — |
| Agrégation de fenêtre | GroupWindowAggregate / WindowAggregate / LocalWindowAggregate / GlobalWindowAggregate | Oui* | 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) | Join | Oui | Oui | Oui* | Les jointures externes (LEFT JOIN, RIGHT JOIN, FULL OUTER JOIN) génèrent des événements de mise à jour. |
| Interval Join | IntervalJoin | Oui | Non | Non | — |
| Temporal Join | TemporalJoin | Oui | Oui | Non | — |
| Jointure de fenêtre | WindowJoin | Oui | Non | Non | — |
| Top-N | Rank | Oui | Oui | Oui | Top-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-N | WindowRank | Oui | Non | Non | — |
| Déduplication | Deduplicate | Oui | Non | Oui* | 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être | WindowDeduplicate | Oui | Non | Non | — |
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.