Use uma junção temporal por tempo de processamento para enriquecer um fluxo de dados com o snapshot mais recente de uma tabela de dimensões. No momento da consulta, cada linha recebida é comparada ao estado atual da tabela de dimensões, ou seja, ao estado no instante em que a linha chega ao Realtime Compute for Apache Flink, e não ao timestamp codificado no próprio evento.
Esse comportamento difere da junção temporal por tempo de evento, que correlaciona as linhas com base no timestamp incorporado ao evento.
Pré-requisitos
Antes de começar, verifique se você tem:
Realtime Compute for Apache Flink com Ververica Runtime (VVR) 8.0.10 ou superior
Uma tabela de dimensões MySQL criada e acessível
Observações de uso
Defina
execution.checkpointing.interval-during-backlog = 0para desativar o checkpoint durante a sincronização completa. O Realtime Compute for Apache Flink não oferece suporte a checkpoints nessa fase. Essa configuração não afeta os checkpoints da sincronização incremental de dados.Defina
table.optimizer.proctime-temporal-join-strategy = TEMPORAL_JOINpara ativar as junções temporais por tempo de processamento.
Sintaxe
A junção temporal por tempo de processamento usa a mesma sintaxe de uma junção padrão com tabela de dimensões:
SELECT column-names
FROM table1 [AS <alias1>]
[LEFT] JOIN table2 FOR SYSTEM_TIME AS OF PROCTIME() [AS <alias2>]
ON table1.column-name1 = table2.key-name1;
A cláusula FOR SYSTEM_TIME AS OF PROCTIME() instrui o Flink a consultar a tabela de dimensões no tempo de processamento de cada linha recebida.
Exemplo
Este exemplo une uma tabela de fatos Kafka a uma tabela de dimensões MySQL para buscar números de telefone pelo nome.
Dados de teste
Tabela 1: kafka_input
|
id (bigint) |
name (varchar) |
age (bigint) |
|
1 |
Lee |
22 |
|
2 |
Harry |
20 |
|
3 |
Liban |
28 |
Tabela 2: phoneNumber
|
name (varchar) |
phoneNumber (bigint) |
|
David |
1390000111 |
|
Brooks |
1390000222 |
|
Liban |
1390000333 |
|
Lee |
1390000444 |
Código de teste
-- Enable processing-time temporal joins.
SET 'table.optimizer.proctime-temporal-join-strategy' = 'TEMPORAL_JOIN';
-- Disable checkpointing during full synchronization.
SET 'execution.checkpointing.interval-during-backlog' = '0';
-- Define the Kafka fact table.
-- proc_time is a computed column that returns the current processing time.
CREATE TEMPORARY TABLE kafka_input (
id BIGINT,
name VARCHAR,
age BIGINT,
proc_time AS PROCTIME()
) WITH (
'connector' = 'kafka',
'topic' = '<yourTopic>',
'properties.bootstrap.servers' = '<yourKafkaBrokers>',
'properties.group.id' = '<yourKafkaConsumerGroupId>',
'format' = 'csv'
);
-- Define the MySQL dimension table.
CREATE TEMPORARY TABLE phoneNumber (
name VARCHAR,
phoneNumber BIGINT,
PRIMARY KEY (name) NOT ENFORCED
) WITH (
'connector' = 'mysql',
'hostname' = '<yourHostname>',
'port' = '3306',
'username' = '<yourUsername>',
'password' = '<yourPassword>',
'database-name' = '<yourDatabaseName>',
'table-name' = '<yourTableName>'
);
-- Define the output table (blackhole discards results; replace with your target sink).
CREATE TEMPORARY TABLE result_infor (
id BIGINT,
phoneNumber BIGINT,
name VARCHAR
) WITH (
'connector' = 'blackhole'
);
-- Join the Kafka stream with the MySQL dimension table at processing time.
INSERT INTO result_infor
SELECT
t.id,
w.phoneNumber,
t.name
FROM kafka_input AS t
JOIN phoneNumber FOR SYSTEM_TIME AS OF t.proc_time AS w
ON t.name = w.name;
Substitua os seguintes espaços reservados pelos valores reais:
|
Espaço reservado |
Descrição |
|
|
Nome do tópico Kafka |
|
|
Endereços dos servidores bootstrap do Kafka |
|
|
ID do grupo de consumidores Kafka |
|
|
Hostname ou endereço IP do servidor MySQL |
|
|
Nome de usuário do MySQL |
|
|
Senha do MySQL |
|
|
Nome do banco de dados MySQL |
|
|
Nome da tabela MySQL |
Resultados do teste
|
id (bigint) |
phoneNumber (bigint) |
name (varchar) |
|
1 |
1390000444 |
Lee |
|
3 |
1390000333 |
Liban |
A linha com id = 2 (Harry) não tem entrada correspondente na tabela de dimensões phoneNumber e foi excluída dos resultados. Use LEFT JOIN em vez de JOIN para incluir linhas sem correspondência com número de telefone null.