O Realtime Compute for Apache Flink é compatível com consultas SQL nativas do Apache Flink. A gramática BNF (Backus-Naur Form) a seguir define o superconjunto de consultas SQL em streaming e em lote compatíveis.
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 '}'
Identificadores
As regras de identificadores do Flink SQL diferem do Java em um aspecto importante: a permissão de caracteres não alfanuméricos. Nos demais aspectos, as regras são idênticas às do Java:
A definição de identificadores diferencia maiúsculas de minúsculas, independentemente do uso de crases.
A correspondência de identificadores diferencia maiúsculas de minúsculas.
Exemplo com identificador não alfanumérico:
SELECT a AS `my field` FROM t
Constantes de string
Constantes de string devem usar aspas simples ('), e não aspas duplas (").
SELECT 'Hello World'
Para escapar uma aspa simples dentro de uma string, duplique-a:
Flink SQL> SELECT 'Hello World', 'It''s me';
+-------------+---------+
| EXPR$0 | EXPR$1 |
+-------------+---------+
| Hello World | It's me |
+-------------+---------+
1 row in set
Sequências de escape Unicode
Para incluir valores Unicode em uma constante de string, use o prefixo U& seguido de uma sequência de escape.
|
Método |
Sintaxe |
Exemplo |
|
Escape padrão (barra invertida) |
|
|
|
Caractere de escape personalizado |
|
|
Consultas compatíveis
A tabela a seguir lista todas as consultas compatíveis com o Apache Flink 1.15. Para consultar a referência de outra versão do Flink, alterne as versões no site do Apache Flink.
|
Consulta |
Referência |
|
Hints |
|
|
Cláusula WITH |
|
|
Cláusulas SELECT e WHERE |
|
|
SELECT DISTINCT |
|
|
Funções de janela |
|
|
Agregação de janela |
|
|
Agregação de grupo |
|
|
Agregação Over |
|
|
Join |
|
|
Window join |
|
|
Operações de conjunto |
|
|
Cláusula ORDER BY |
|
|
Cláusula LIMIT |
|
|
Top-N |
|
|
Window Top-N |
|
|
Deduplicação |
|
|
Deduplicação de janela |
|
|
Reconhecimento de padrões |
Execução de consultas
No modo streaming, os fluxos de entrada dividem-se em duas categorias:
Fluxos sem atualização: contêm apenas eventos do tipo INSERT.
Fluxos com atualização: contêm outros tipos de eventos. Fontes de captura de dados de alteração (CDC) produzem fluxos com atualização. Certas operações do Flink, como agregação de grupo e Top-N, também geram eventos de atualização internamente.
A maioria das operações que produzem eventos de atualização depende de operadores com estado, os quais utilizam estado gerenciado para rastrear atualizações. Contudo, nem todos os operadores com estado aceitam fluxos com atualização como entrada. Por exemplo, a agregação Over e o Interval Join não aceitam entradas de fluxo com atualização.
A tabela a seguir descreve as características de execução de cada consulta compatível. As informações aplicam-se ao Ververica Runtime (VVR) 6,0.X e versões posteriores.
| Consulta | Operador de runtime | Usa dados de estado | Consome fluxos com atualização | Gera eventos de atualização | Observações |
|---|---|---|---|---|---|
| SELECT e WHERE | Calc | Não | Sim | Não | — |
| Lookup Join | LookupJoin | Não* | Sim | Não | Para VVR 8.0.1 e posterior: defina table.optimizer.non-deterministic-update.strategy como TRY_RESOLVE ativa a resolução automática baseada em estado para problemas de atualização não determinística. Defina como IGNORE para desativar o uso de estado. Alterar esse parâmetro pode causar incompatibilidade e exigir a reexecução da consulta. |
| Função de tabela | Correlate | Não | Sim | Não | — |
| SELECT DISTINCT | GroupAggregate | Sim | Sim | Sim | — |
| Agregação de grupo | GroupAggregate / LocalGroupAggregate / GlobalGroupAggregate / IncrementalGroupAggregate | Sim* | Sim | Sim | O operador de pré-agregação LocalGroupAggregate não utiliza dados de estado. |
| Agregação Over | OverAggregate | Sim | Não | Não | — |
| Agregação de janela | GroupWindowAggregate / WindowAggregate / LocalWindowAggregate / GlobalWindowAggregate | Sim* | Sim* | Não* | LocalWindowAggregate não utiliza dados de estado. O suporte a fluxos com atualização difere entre VVR e Apache Flink. Consulte a seção "Comparação do suporte para fluxos com atualização" no tópico Agregação de janela para obter detalhes. Se o recurso de disparo antecipado ou tardio (experimental) estiver ativado, o sistema gerará eventos de atualização. |
| Join (junção regular) | Join | Sim | Sim | Sim* | Junções externas (LEFT JOIN, RIGHT JOIN, FULL OUTER JOIN) geram eventos de atualização. |
| Interval Join | IntervalJoin | Sim | Não | Não | — |
| Temporal Join | TemporalJoin | Sim | Sim | Não | — |
| Window join | WindowJoin | Sim | Não | Não | — |
| Top-N | Rank | Sim | Sim | Sim | Top-N não aceita classificação por tempo de processamento. Use funções integradas como CURRENT_TIMESTAMP. Aviso Especificar um campo de tempo de processamento na cláusula ORDER BY pode causar erros de dados. Esse problema não é relatado durante verificações de sintaxe no VVR 8.0.7 e versões anteriores. |
| Window Top-N | WindowRank | Sim | Não | Não | — |
| Deduplicação | Deduplicate | Sim | Não | Sim* | Usar a política Deduplicate Keep FirstRow com tempo de processamento (Proctime) não gera eventos de atualização. |
| Deduplicação de janela | WindowDeduplicate | Sim | Não | Não | — |
Operadores sem estado repassam os tipos de evento sem modificação; os eventos de saída mantêm o mesmo tipo dos eventos de entrada. Operadores sem estado nunca geram eventos de atualização, mesmo que a entrada seja um fluxo com atualização.