Os runtime filters reduzem o tempo de consultas com join ao podar a tabela da esquerda antes da execução do hash join. Em vez de verificar todas as linhas, o ApsaraDB for SelectDB gera um filtro em tempo de execução a partir da tabela da direita e o envia para a camada de varredura. Assim, apenas as linhas correspondentes entram no join.
Os runtime filters estão ativados por padrão. O ApsaraDB for SelectDB gera automaticamente predicados IN e bloom filters com base na consulta e nas estatísticas das tabelas. Use variáveis de sessão para ajustar o comportamento em consultas específicas.
Como funciona
Em um hash join, o sistema carrega primeiro a tabela da direita para construir uma tabela hash. Um runtime filter captura os valores-chave dessa tabela hash e envia o filtro para o OlapScanNode que lê a tabela da esquerda. O OlapScanNode descarta as linhas não correspondentes antes que elas cheguem ao HashJoinNode.
O exemplo a seguir usa duas tabelas: T1 (uma tabela de fatos com 1.000.000 de linhas) e T2 (uma tabela de dimensões com 200 linhas).
Sem runtime filter — todas as 1.000.000 de linhas da T1 fluem até o join:
HashJoinNode
| |
| 1,000,000 | 200
| |
OlapScanNode OlapScanNode
^ ^
| 1,000,000 | 200
T1 (fact) T2 (dimension)
Com um runtime filter aplicado na camada de varredura — apenas 6.000 linhas chegam ao join:
HashJoinNode
| |
| 6,000 | 200
| |
OlapScanNode OlapScanNode
^ ^
| 1,000,000 | 200
T1 (fact) T2 (dimension)
Com o filtro enviado ao mecanismo de armazenamento — os índices podam os dados antes da leitura:
HashJoinNode
| |
| 6,000 | 200
| |
OlapScanNode OlapScanNode
^ ^
| 6,000 | 200
T1 (fact) T2 (dimension)
Diferentemente do predicate pushdown ou da poda de partição, a condição do filtro não é conhecida durante o planejamento da consulta. O sistema a calcula a partir dos dados reais da tabela da direita durante a execução e depois a transmite ao OlapScanNode que lê a tabela da esquerda.
Conceitos principais
Tabela da esquerda — a tabela no lado esquerdo de um join; usada para a operação de probe. A reordenação de join pode ajustar qual tabela fica em cada lado.
Tabela da direita — a tabela no lado direito de um join; usada para construir a tabela hash. A reordenação de join também pode ajustar essa posição.
Fragment — uma unidade de execução de consulta. O nó frontend (FE) divide uma instrução SQL em fragments e os distribui para os nós backend (BE) no cluster distribuído.
Quando os runtime filters são úteis
Os runtime filters são mais eficazes quando:
A tabela da esquerda é significativamente maior que a tabela da direita. Gerar um filtro tem custo de memória e computação; isso só compensa em grande escala.
O resultado do join é muito menor que a tabela da esquerda — o que significa que o filtro consegue descartar a maioria das linhas da tabela da esquerda.
Se a tabela da direita for grande ou se o resultado do join tiver tamanho próximo ao da tabela da esquerda, os runtime filters podem adicionar sobrecarga sem trazer benefícios.
Tipos de filtro
O ApsaraDB for SelectDB oferece suporte a cinco tipos de runtime filter:
|
Tipo |
Como funciona |
Mais indicado para |
Limitações |
|
Predicado IN |
Cria um HashSet com todos os valores-chave da tabela da direita; filtra a tabela da esquerda com |
Tabelas da direita pequenas em broadcast joins |
Funciona apenas em broadcast joins; torna-se inválido quando as linhas da tabela da direita excedem |
|
Bloom filter |
Cria uma estrutura probabilística a partir da tabela hash; possui uma pequena taxa de falsos positivos |
Tabelas da direita grandes e maioria dos tipos de dados |
Maior sobrecarga de criação e aplicação; pode prejudicar o desempenho se a taxa de filtragem for baixa ou se a tabela da esquerda for pequena |
|
Filtro MinMax |
Extrai o intervalo mínimo/máximo da tabela da direita; filtra linhas fora desse intervalo |
Colunas-chave numéricas com intervalos não sobrepostos |
Ineficaz em colunas não numéricas (por exemplo, VARCHAR); sem benefício quando os intervalos se sobrepõem completamente |
|
IN_OR_BLOOM_FILTER |
Seleciona automaticamente o predicado IN ou bloom filter com base na contagem de linhas da tabela da direita em tempo de execução |
Uso geral (padrão) |
Limiar controlado por |
|
Filtro Bitmap |
Filtra usando uma coluna bitmap retornada por uma subconsulta IN |
Consultas com subconsultas bitmap IN |
Suportado apenas no mecanismo vetorizado |
Configure runtime filters
Use as variáveis de sessão a seguir para ajustar o comportamento dos runtime filters em consultas específicas.
Parâmetros
|
Parâmetro |
Padrão |
Descrição |
|
|
|
Controla até onde os filtros se propagam entre os fragments de execução. Valores: |
|
|
|
Tipo(s) de filtro a gerar. Valores: |
|
|
|
Tempo máximo de espera (ms) do OlapScanNode antes de iniciar uma varredura. Aplica-se por filtro: três filtros significam até 3.000 ms no total. |
|
|
|
Número máximo de bloom filters por consulta. O sistema mantém os bloom filters com maior seletividade quando o limite é excedido. |
|
|
|
Tamanho mínimo do bloom filter em bytes. |
|
|
|
Tamanho máximo do bloom filter em bytes. |
|
|
|
Tamanho padrão do bloom filter quando a cardinalidade da tabela da direita não está disponível. |
|
|
|
Limiar de contagem de linhas acima do qual nenhum predicado IN é gerado. Também controla a alternância entre IN e bloom filter no modo |
runtime_filter_mode
Controla o escopo de propagação do filtro entre os fragments de execução da consulta.
**
LOCAL** — o produtor do filtro (HashJoinNode) e o consumidor (OlapScanNode) devem estar no mesmo fragment. Adequado para broadcast joins. Menor sobrecarga.**
GLOBAL** — os filtros são mesclados e transferidos pela rede entre as fronteiras dos fragments. Necessário para shuffle joins em que o produtor e o consumidor estão em fragments diferentes.
O modo GLOBAL cobre todos os cenários atendidos pelo LOCAL, além dos shuffle joins. Mude para LOCAL se a sobrecarga de mesclar e transmitir filtros superar a economia na varredura para um shuffle join específico. Defina como OFF para desativar completamente os runtime filters.
Para obter detalhes sobre o design técnico da mesclagem de filtros entre fragments, consulte ISSUE 6116.
runtime_filter_type
Especifique um tipo ou combine vários tipos:
-- By name (comma-separated, quoted)
SET runtime_filter_type = "BLOOM_FILTER,IN,MIN_MAX";
-- Equivalent numeric form (1=IN, 2=BLOOM_FILTER, 4=MIN_MAX)
SET runtime_filter_type = 7;
Comportamento do predicado IN:
Usa uma estratégia de mesclagem para execução distribuída.
Quando o IN e outros tipos de filtro são especificados simultaneamente e a tabela da direita permanece dentro de
runtime_filter_max_in_num, o sistema descarta os outros filtros. Como o predicado IN é exato, filtros adicionais não trazem benefícios. Essa otimização aplica-se apenas quando o produtor e o consumidor estão no mesmo fragment.
Comportamento do Bloom filter:
Possui uma taxa de falsos positivos diferente de zero, o que significa que pode deixar passar algumas linhas não correspondentes. Os resultados permanecem corretos; apenas a eficiência da filtragem é afetada.
Pode ser enviado ao mecanismo de armazenamento, mas apenas para colunas-chave na tabela da esquerda. Sem o pushdown na camada de armazenamento, o desempenho geralmente piora.
Evite usar bloom filters em colunas com baixa seletividade ou em tabelas da esquerda pequenas.
Comportamento do filtro MinMax:
Mais eficaz quando os intervalos de valores das tabelas da direita e da esquerda não se sobrepõem (por exemplo, o máximo da tabela da direita está abaixo do mínimo da tabela da esquerda, ou vice-versa).
Em colunas numéricas (INT, BIGINT, DOUBLE), intervalos sobrepostos reduzem a eficácia a zero.
Em colunas não numéricas (VARCHAR), normalmente degrada o desempenho.
Comportamento do IN_OR_BLOOM_FILTER:
Usa o predicado IN quando as linhas da tabela da direita forem < 102.400; muda para bloom filter acima desse limiar.
Ajuste o limiar com
runtime_filter_max_in_num.
Comportamento do filtro Bitmap:
Aplicável apenas quando a subconsulta IN retorna uma coluna bitmap.
Requer o mecanismo vetorizado.
runtime_filter_wait_time_ms
O OlapScanNode aguarda até essa duração por cada runtime filter atribuído antes de iniciar uma varredura. Com três filtros, a espera máxima é três vezes esse valor.
Os filtros que chegam dentro da janela de espera são enviados ao mecanismo de armazenamento. Os filtros que chegam após o início da varredura são aplicados como filtros de expressão nos dados já verificados — eficazes, mas menos eficientes que o pushdown na camada de armazenamento.
Orientações de ajuste:
Clusters ocupados executando joins longos — aumente este valor para que os filtros tenham tempo de ser construídos e transmitidos antes do início da varredura.
Clusters leves executando muitas consultas pequenas — diminua este valor para evitar adicionar latência desnecessária a consultas que terminam em segundos.
Se o cluster estiver ocupado com muitas consultas intensivas em recursos ou demoradas, prolongue a duração da espera para evitar que consultas complexas percam oportunidades de otimização. Se a carga do cluster for leve e houver muitas consultas pequenas que levam apenas alguns segundos, reduza a duração da espera para evitar um aumento de 1 s de latência por consulta.
runtime_filters_max_num
Limita o número de bloom filters por consulta. Predicados IN e filtros MinMax não são contabilizados.
Quando o limite é excedido, o sistema retém os bloom filters com maior seletividade — aqueles que devem filtrar mais linhas:
Selectivity = HashJoinNode cardinality / HashJoinNode left child cardinality
A cardinalidade estimada pelo FE pode ser imprecisa; portanto, a classificação baseada em seletividade pode não refletir perfeitamente a eficácia real da filtragem.
Ajuste este parâmetro apenas ao otimizar consultas de join lentas entre tabelas grandes.
Parâmetros de tamanho do bloom filter
O FE calcula o comprimento do bloom filter no momento do planejamento da consulta. Todos os bloom filters do HashJoinNode em uma consulta devem ter o mesmo comprimento para permitir a mesclagem.
Se a cardinalidade da tabela da direita estiver disponível nas estatísticas, o FE estima o tamanho ideal e arredonda para cima até a potência de 2 mais próxima.
Se a cardinalidade não estiver disponível, o FE usa
runtime_bloom_filter_sizecomo padrão.runtime_bloom_filter_min_sizeeruntime_bloom_filter_max_sizelimitam o tamanho final, independentemente da estimativa.
Bloom filters maiores lidam com colunas de alta cardinalidade com mais precisão, mas consomem mais memória. Se a precisão da filtragem for insuficiente para uma coluna de alta cardinalidade (milhões de valores distintos), aumente runtime_bloom_filter_size e faça benchmarks do resultado.
Ajuste o tamanho do bloom filter por consulta, não globalmente.
Verifique se os runtime filters foram aplicados
Execute EXPLAIN para visualizar o plano de consulta e confirme se os filtros foram gerados e consumidos nas colunas esperadas.
Lado do join (filtro gerado):
runtime filters: RF000[in] <- table.columnLado da varredura (filtro aplicado):
runtime filters: RF000[in] -> table.column
Exemplo:
CREATE TABLE test (t1 INT) DISTRIBUTED BY HASH (t1) BUCKETS 2;
INSERT INTO test VALUES (1), (2), (3), (4);
CREATE TABLE test2 (t2 INT) DISTRIBUTED BY HASH (t2) BUCKETS 2;
INSERT INTO test2 VALUES (3), (4), (5);
EXPLAIN SELECT t1 FROM test JOIN test2 WHERE test.t1 = test2.t2;
+-------------------------------------------------------------------+
| Explain String |
+-------------------------------------------------------------------+
| PLAN FRAGMENT 0 |
| OUTPUT EXPRS:`t1` |
| |
| 4:EXCHANGE |
| |
| PLAN FRAGMENT 1 |
| OUTPUT EXPRS: |
| PARTITION: HASH_PARTITIONED: `default_cluster:ssb`.`test`.`t1` |
| |
| 2:HASH JOIN |
| | join op: INNER JOIN (BUCKET_SHUFFLE) |
| | equal join conjunct: `test`.`t1` = `test2`.`t2` |
| | runtime filters: RF000[in] <- `test2`.`t2` |
| | |
| |----3:EXCHANGE |
| | |
| 0:OlapScanNode |
| TABLE: test |
| runtime filters: RF000[in] -> `test`.`t1` |
| |
| PLAN FRAGMENT 2 |
| OUTPUT EXPRS: |
| PARTITION: HASH_PARTITIONED: `default_cluster:ssb`.`test2`.`t2` |
| |
| 1:OlapScanNode |
| TABLE: test2 |
+-------------------------------------------------------------------+
Lado do join — 2:HASH JOIN no PLAN FRAGMENT 1 gera um predicado IN (RF000) a partir de test2.t2. Os valores só são conhecidos em tempo de execução.
Lado da varredura — 0:OlapScanNode aplica o RF000 para filtrar test.t1 antes que as linhas cheguem ao join.
A execução da consulta retorna [3, 4] — apenas as duas linhas em que t1 = t2.
Verifique a eficácia do filtro com um profile:
Ative o profiling e execute a consulta:
SET enable_profile = true;
No profile, procure a seção RuntimeFilter em OLAP_SCAN_NODE:
RuntimeFilter:in:
HasPushDownToEngine: true -- filter reached the storage engine
AWaitTimeCost: 0ns -- no wait; filter arrived before scan started
EffectTimeCost: 2.76ms -- time spent applying the filter
E o resultado da filtragem:
RowsVectorPredFiltered: 9,320,008 -- rows discarded by the filter
VectorPredEvalTime: 364.39ms -- time spent on filter evaluation
Uma contagem alta de RowsVectorPredFiltered confirma que o filtro é eficaz. HasPushDownToEngine: true confirma que ele chegou à camada de armazenamento.
Regras de planejamento
Os runtime filters são gerados e aplicados de acordo com as regras a seguir. Violações dessas regras fazem com que os filtros sejam ignorados ou produzam resultados incorretos.
Regras de geração:
Os filtros são gerados apenas para condições de igualdade nas cláusulas JOIN ON. A igualdade segura para NULL (
<=>) é excluída porque valores nulos na tabela da esquerda poderiam ser filtrados incorretamente.O tipo da expressão de source não pode ser
HLLouBITMAP.As expressões de source e de destino não podem ser constantes.
As expressões de source e de destino não podem ser a mesma expressão.
Os tipos das expressões de source e de destino devem corresponder (os bloom filters são baseados em hash). Se os tipos diferirem, o sistema tenta converter a expressão de destino para o tipo da source.
Os filtros de
PlanNode.Conjunctsnão sofrem pushdown — estes podem produzir resultados incorretos. Por exemplo, quando uma subconsultaINé reescrita como um join, o sistema armazena a condição JOIN gerada automaticamente emPlanNode.Conjuncts; aplicar um runtime filter ali pode causar a remoção de linhas dos resultados.
Regras de pushdown:
Os filtros só podem ser enviados para o OlapScanNode. Outros tipos de nós de varredura não são suportados.
Os filtros não podem ser enviados para a tabela da esquerda em left outer joins, full outer joins ou anti-joins.
A expressão de destino deve referenciar uma coluna existente na tabela base original.
A expressão de destino não pode conter expressões de verificação de nulo, como
COALESCE,IFNULLouCASE.
Regras de condução de coluna:
Se a cláusula JOIN ON contiver
A.k = B.k AND B.k = C.k, o filtro paraC.kpoderá ser enviado apenas paraB.k— não se propagando ainda mais paraA.k.Se a cláusula JOIN ON contiver
A.a + B.b = C.ceA.afor equivalente aB.a,A.apoderá ser substituído porB.ae o filtro enviado paraB. SeA.aeB.anão forem equivalentes, o filtro não poderá ser enviado paraBporque a expressão de destino deve estar vinculada a uma única tabela da esquerda.