O Flink SQL oferece suporte a junções regulares (também chamadas de junções de fluxo duplo) em tabelas dinâmicas. Diferentemente das junções em lote, os resultados são atualizados continuamente conforme novos dados chegam em qualquer um dos fluxos. Isso garante que a saída final seja consistente com o resultado de uma junção em lote sobre os mesmos dados.
Contexto
As junções regulares armazenam em buffer todos os registros históricos de ambos os fluxos de entrada no estado do Flink. Assim, qualquer linha recebida pode ser comparada com todas as linhas vistas anteriormente no outro lado. Consequentemente, o estado cresce indefinidamente enquanto o job estiver em execução. Em jobs de longa duração com fluxos de alto volume, o crescimento ilimitado do estado pode esgotar a memória e causar falhas.
Para limitar o tamanho do estado, defina um tempo de vida (TTL) de estado usando o parâmetro table.exec.state.ttl ou a dica JOIN_STATE_TTL por fluxo (consulte Dicas). A evicção do estado pode afetar a precisão dos resultados, pois linhas expiradas não corresponderão mais aos novos dados recebidos.
Caso o estado ilimitado seja uma preocupação, avalie também o uso de junções de intervalo, junções de consulta e junções de janela como alternativas.
Sintaxe
tableReference [, tableReference ]*
| tableExpression [ NATURAL | INNER ] [ { LEFT | RIGHT | FULL } [ OUTER ] ] JOIN tableExpression [ joinCondition ]
| tableExpression CROSS JOIN tableExpression
| tableExpression [ CROSS | OUTER ] APPLY tableExpression
joinCondition:
ON booleanExpression
| USING '(' column [, column ]* ')'
tableReference: nome da tabela.tableExpression: expressão da tabela.joinCondition: condição de junção.
Liste primeiro a tabela com menor frequência de atualização e, por último, a tabela com maior frequência. Essa prática reduz a quantidade de estado necessária para correspondência a cada linha recebida.
Tipos de junção
INNER JOIN
Retorna apenas as linhas em que a condição de junção é satisfeita em ambos os lados.
SELECT id, productName, status
FROM orders o
JOIN shipments s
ON o.id = s.orderId;
LEFT JOIN
Retorna todas as linhas da tabela à esquerda. Linhas sem correspondência no lado direito geram valores NULL nas colunas do lado direito.
SELECT o.id, o.productName, s.status
FROM orders o
LEFT JOIN shipments s
ON o.id = s.orderId;
RIGHT JOIN
Retorna todas as linhas da tabela à direita. Linhas sem correspondência no lado esquerdo geram valores NULL nas colunas do lado esquerdo.
SELECT o.id, o.productName, s.status
FROM orders o
RIGHT JOIN shipments s
ON o.id = s.orderId;
FULL OUTER JOIN
Retorna todas as linhas de ambas as tabelas. Linhas sem correspondência em qualquer um dos lados produzem valores NULL nas colunas não correspondidas.
SELECT o.id, o.productName, s.status
FROM orders o
FULL OUTER JOIN shipments s
ON o.id = s.orderId;
CROSS JOIN
Retorna o produto cartesiano de ambas as tabelas, combinando cada linha do lado esquerdo com todas as linhas do lado direito.
SELECT *
FROM orders o
CROSS JOIN shipments s;
Dicas
No Ververica Runtime (VVR) 8.0.1 e versões posteriores, utilize dicas para especificar valores diferentes de tempo de vida (TTL) para os estados dos fluxos esquerdo e direito. Isso reduz a quantidade de dados de estado mantidos.
A dica JOIN_STATE_TTL aplica-se exclusivamente a junções regulares. Ela não oferece suporte a junções de consulta, de intervalo ou de janela. As dicas para junções regulares são experimentais e a sintaxe pode mudar em versões futuras.
Sintaxe
-- VVR 8.0.1 and later
SELECT /*+ JOIN_STATE_TTL('tableReference1' = 'ttl1' [, 'tableReference2' = 'ttl2']*) */ ...
-- VVR 8.0.7 and later also support the Apache Flink syntax
SELECT /*+ STATE_TTL('tableReference1' = 'ttl1' [, 'tableReference2' = 'ttl2']*) */ ...
Defina tableReference como um nome de tabela, nome de view ou alias. Se você especificar um alias para uma tabela, use esse alias na dica.
Ao especificar um TTL para apenas um fluxo, o outro fluxo utilizará o TTL no nível de implantação definido pelo parâmetro table.exec.state.ttl. O valor padrão é 1,5 dias. Para mais detalhes, consulte Configurações básicas.
Exemplos
Todos os exemplos usam JOIN_STATE_TTL para definir TTLs independentes para cada fluxo. O VVR 8.0.7 e versões posteriores também aceitam a sintaxe equivalente STATE_TTL mostrada nos comentários.
Uso de alias
SELECT /*+ JOIN_STATE_TTL('o' = '3d', 'p' = '1d') */
o.rowtime, o.productid, o.orderid, o.units, p.name, p.unitprice
FROM Orders AS o
JOIN Products AS p
ON o.productid = p.productid;
-- VVR 8.0.7 and later
SELECT /*+ STATE_TTL('o' = '3d', 'p' = '1d') */
o.rowtime, o.productid, o.orderid, o.units, p.name, p.unitprice
FROM Orders AS o
JOIN Products AS p
ON o.productid = p.productid;
Uso de nome de tabela
SELECT /*+ JOIN_STATE_TTL('Orders' = '3d', 'Products' = '1d') */ *
FROM Orders
JOIN Products
ON Orders.productid = Products.productid;
-- VVR 8.0.7 and later
SELECT /*+ STATE_TTL('Orders' = '3d', 'Products' = '1d') */ *
FROM Orders
JOIN Products
ON Orders.productid = Products.productid;
Uso de nome de view
CREATE TEMPORARY VIEW v AS
SELECT id, ...
FROM (
SELECT id, ROW_NUMBER() OVER (PARTITION BY ... ORDER BY ..) AS rn
FROM src1
WHERE ...
) tmp
WHERE rn = 1;
SELECT /*+ JOIN_STATE_TTL('v' = '1d', 'b' = '3d') */ v.* , b.*
FROM v
LEFT JOIN src2 AS b ON v.id = b.id;
-- VVR 8.0.7 and later
SELECT /*+ STATE_TTL('v' = '1d', 'b' = '3d') */ v.* , b.*
FROM v
LEFT JOIN src2 AS b ON v.id = b.id;
Exemplos
Exemplo 1: Junção das tabelas Orders e Shipments
Dados de teste
Tabela 1. Orders
|
id |
productName |
ordertime |
|
1 |
phone_a |
2025-05-01 10:00:00.0 |
|
2 |
notebook_x |
2025-05-01 10:02:00.0 |
|
3 |
phone_b |
2025-05-01 10:03:00.0 |
|
4 |
pad_m |
2025-05-01 10:05:00.0 |
Tabela 2. Shipments
|
shipId |
orderId |
status |
shipTime |
|
101 |
1 |
shipped |
2025-05-01 11:00:00.0 |
|
102 |
2 |
delivered |
2025-05-01 17:00:00.0 |
|
103 |
3 |
shipped |
2025-05-01 12:00:00.0 |
|
104 |
4 |
shipped |
2025-05-01 11:30:00.0 |
Instrução de teste
SELECT id, productName, status
FROM orders o
JOIN shipments s
ON o.id = s.orderId;
Resultados do teste
|
id |
productName |
status |
|
1 |
phone_a |
shipped |
|
2 |
notebook_x |
delivered |
|
3 |
phone_b |
shipped |
|
4 |
pad_m |
shipped |
Exemplo 2: Junção das tabelas datahub_stream1 e datahub_stream2
Dados de teste
Tabela 3. datahub_stream1
|
a (BIGINT) |
b (BIGINT) |
c (VARCHAR) |
|
0 |
10 |
test11 |
|
1 |
10 |
test21 |
Tabela 4. datahub_stream2
|
a (BIGINT) |
b (BIGINT) |
c (VARCHAR) |
|
0 |
10 |
test11 |
|
1 |
10 |
test21 |
|
0 |
10 |
test31 |
|
1 |
10 |
test41 |
Instrução de teste
SELECT s1.c, s2.c
FROM datahub_stream1 AS s1
JOIN datahub_stream2 AS s2
ON s1.a = s2.a
WHERE s1.a = 0;
Resultados do teste
|
s1.c (VARCHAR) |
s2.c (VARCHAR) |
|
test11 |
test11 |
|
test11 |
test31 |