Todos os produtos
Search
Central de documentação

Hologres:Acelere junções de múltiplas tabelas com Runtime Filters

Última atualização: Jun 28, 2026

Em consultas com junção de múltiplas tabelas, dados que não correspondem à condição de junção aumentam a sobrecarga de I/O e degradam o desempenho da consulta. Nesse cenário, os runtime filters geram automaticamente filtros leves que eliminam dados não correspondentes durante a fase de varredura. Essa abordagem reduz a sobrecarga de I/O e melhora o desempenho da consulta. O recurso oferece suporte a Hash Join desde o Hologres V2.0 e foi estendido para cenários de Cross Join na versão V4.2.

Contexto

Casos de uso

O Hologres oferece suporte a runtime filters desde a versão V2.0. Esse recurso é comum em cenários de Hash Join envolvendo duas ou mais tabelas, especialmente quando uma tabela grande se junta a uma pequena. Nenhuma configuração manual é necessária. O otimizador e o mecanismo de execução otimizam automaticamente o comportamento de filtragem de junção no momento da consulta, o que reduz a sobrecarga de I/O e melhora o desempenho da junção.

A partir da V4.2, a capacidade de runtime filter se estende a cenários de Cross Join. Ela otimiza especificamente o padrão SQL comum em que uma subconsulta escalar calcula um resultado agregado e o utiliza como condição de filtro em uma tabela grande. Para obter mais informações, consulte Suporte a Cross Join para runtime filters (Novidade na V4.2).

Funcionamento

Runtime filters para hash join

Ao unir duas tabelas, o sistema carrega os dados de uma delas em uma tabela hash e compara os dados da outra tabela com essa estrutura. O processo de junção possui dois lados:

  • Build side: lado que constrói a tabela hash, correspondente ao nó Hash no plano de execução.

  • Probe side: lado que lê os dados e os compara com a tabela hash do build side.

Geralmente, a tabela pequena atua como build side e a tabela grande como probe side.

Os runtime filters funcionam construindo um filtro leve a partir da distribuição de dados do build side e enviando-o ao probe side para eliminar dados. Isso reduz a quantidade de dados processados pelo probe side no Hash Join e minimiza o tráfego de rede, melhorando assim o desempenho da junção. Consequentemente, esse recurso é mais eficaz para junções entre tabelas grandes e pequenas com disparidade significativa de tamanho, oferecendo ganhos de desempenho superiores aos de uma junção padrão.

Runtime filters para cross join (Novidade na V4.2)

Antes da V4.2, os runtime filters cobriam apenas Hash Join. No entanto, usuários frequentemente usam subconsultas escalares para calcular valores agregados, como min ou max, e depois usam o resultado como condição de filtro em uma tabela grande. No plano de execução, esse padrão SQL produz um Cross Join:

  • Build side: resultado da subconsulta escalar, contendo exatamente uma linha.

  • Probe side: varredura da tabela grande.

A V4.2 introduz um novo tipo de filtro, o ScalarFilter. O otimizador identifica automaticamente esse padrão e envia o valor escalar avaliado do build side para o ScanNode do probe side. Isso permite a filtragem no nível de linha e de RowGroup durante a fase de varredura, evitando uma varredura completa da tabela seguida de comparações linha por linha.

Limitações e condições de ativação

Limitações

  • Os runtime filters são suportados apenas no Hologres V2.0 e versões posteriores.

  • Em cenários de Hash Join, a V2.0 suporta runtime filters apenas quando a condição de junção contém um único campo. A partir da V2.1, há suporte para múltiplos campos.

  • O runtime filter TopN é suportado apenas na V4.0 e versões posteriores, sendo utilizado para melhorar o desempenho em cálculos TopN de tabela única.

  • O runtime filter Cross Join (ScalarFilter) é suportado apenas na V4.2 e versões posteriores.

Condições de ativação

Cenários de Hash Join

O mecanismo ativa um runtime filter automaticamente quando todas as seguintes condições são atendidas:

  • O probe side possui 100.000 linhas ou mais.

  • A proporção de dados verificados do build side em relação ao probe side é de 0,1 ou menos. Quanto menor a proporção, maior a probabilidade de ativação do filtro.

  • A proporção de dados de saída da junção em relação aos dados do probe side é de 0,1 ou menos. Quanto menor a proporção, maior a probabilidade de ativação do filtro.

Cenários de Cross Join (V4.2+)

O otimizador gera automaticamente um ScalarFilter quando as seguintes condições são atendidas:

  • O build side do Cross Join tem uma contagem de linhas estatística de 1 e sua distribuição é Replicated.

  • Os predicados de filtro no probe side contêm condições de comparação com expressões do build side, como >, <, >=, <= ou BETWEEN.

Tipos de runtime filters

Os runtime filters podem ser classificados nas duas dimensões a seguir.

Por escopo de shuffle (para hash join)

Tipo

Versões suportadas

Cenários

Local

V2.0+

Utilizado quando os dados do probe side não precisam passar por shuffle. Um runtime filter Local pode ser usado se as chaves de junção do build side e do probe side compartilharem a mesma distribuição, se os dados do build side forem transmitidos (broadcast) para o probe side, ou se os dados do build side passarem por shuffle para corresponder à distribuição do probe side. Este tipo reduz apenas a quantidade de dados verificados e processados pelo Hash Join.

Global

V2.2+

Aplicado quando os dados do probe side precisam passar por shuffle. O runtime filter é aplicado antes do shuffle de dados, o que reduz o tráfego de rede.

Nota

Não é necessário especifique o tipo. O mecanismo o selecione adaptativamente.

Por tipo de filtro

Tipo

Versões suportadas

Descrição

Bloom filter

V2.0+

Filtro probabilístico sujeito a falsos positivos, o que significa que alguns dados podem não ser filtrados. Contudo, possui ampla aplicabilidade e mantém alta eficiência de filtragem mesmo quando o build side contém uma grande quantidade de dados.

In filter

V2.0+

Recomendado quando o build side possui um NDV (número de valores distintos) baixo. Constrói um HashSet a partir dos dados do build side e o envia ao probe side para filtragem. Filtra com precisão todos os dados necessários e pode ser usado em conjunto com um índice bitmap.

MinMax filter

V2.0+

Envia os valores mínimo e máximo dos dados do build side para o probe side para filtragem. Pode aproveitar metadados para ignorar arquivos inteiros ou lotes de dados, reduzindo custos de I/O.

ScalarFilter

V4.2+

Projetado especificamente para cenários de Cross Join. Quando o build side tem exatamente uma linha, o mecanismo envia o valor escalar para o ScanNode do probe side para filtragem durante a fase de varredura.

Nota

Não é necessário especifique o tipo de filtro. O Hologres selecione um tipo adaptativamente com base nas condições de junção em tempo de execução.

Suporte a Cross Join para runtime filters (Novidade na V4.2)

Motivação

Usuários frequentemente calculam valores agregados com subconsultas escalares e usam os resultados como condições de filtro em uma tabela grande. Antes da V4.2, esse padrão SQL apresentava os seguintes problemas no plano de execução:

  • O ScanNode do probe side não conseguia aproveitar o valor do build side para filtragem antecipada e precisava realizar uma varredura completa da tabela.

  • Após a varredura, o Cross Join realizava comparação linha por linha, causando desperdício significativo de I/O.

  • O runtime filter existente cobria apenas Hash Join e não suportava Cross Join.

A V4.2 introduz o ScalarFilter para resolver esse problema. Esse recurso está habilitado por padrão e não requer ação do usuário.

Comparação de SQL típico e plano de execução

SQL típico:

-- t1 is a large table; the aggregated min/max result from t2 is exactly 1 row.
SELECT * FROM t1
WHERE a >= (SELECT min(a) FROM t2 WHERE b BETWEEN 0 AND 1);

Antes da otimização (sem runtime filter):

 Cross Join
   ->  Seq Scan on t1             -- Full table scan, no early filtering
   ->  Aggregate                  -- Build side min/max, 1-row result
         ->  Seq Scan on t2

O probe side precisa verificar todos os dados e, em seguida, o Cross Join compara as linhas uma a uma, causando desperdício significativo de I/O.

Após a otimização (V4.2 com ScalarFilter):

 Cross Join
   Runtime Filter Build Expr: (min(t2.a)), (max(t2.a))
   ->  Seq Scan on t1             -- Receives ScalarFilter, filters non-matching rows and RowGroups during scan
         Runtime Filter Target Expr: (t1.a >= ${1}) AND (t1.a <= ${2})
   ->  Aggregate
         ->  Seq Scan on t2

Depois que o build side é avaliado, os valores reais são enviados para o ScanNode do probe side, permitindo que a filtragem seja concluída durante a fase de varredura.

Mais cenários típicos

-- Scenario 1: Single table + scalar subquery
SELECT * FROM t1
WHERE a >= (SELECT min(a) FROM t2 WHERE b BETWEEN 0 AND 1);
-- Scenario 2: Multi-table + scalar subquery
SELECT * FROM t1, t3
WHERE (SELECT min(a) FROM t2) <= t1.a
  AND t3.a <= (SELECT max(a) FROM t2);
-- Scenario 3: CTE + Cross Join
WITH r AS (SELECT MIN(a) AS lo, MAX(a) AS hi FROM t2)
SELECT t1.* FROM t1, r
WHERE t1.a >= r.lo AND t1.a <= r.hi;

Funcionamento

  1. Identificação pelo otimizador: quando o build side de um Cross Join tem uma contagem de linhas estatística de 1 e distribuição Replicated, o otimizador extrai expressões de comparação dos predicados (expandindo BETWEEN para >= e <=) e gera candidatos a ScalarFilter.

  2. Pushdown em tempo de execução: durante a fase Open, o Cross Join consome os dados do build side, avalia a build_expr para produzir um valor escalar, constrói um ScalarFilter e o publica no ScanNode do probe side.

  3. Filtragem antecipada no ScanNode: após o NiagaraScan do probe side receber o ScalarFilter, ele substitui os espaços reservados na target_expr pelos valores reais e gera dois níveis de filtragem:

    • Filtragem no nível de linha (avaliação conjuntiva).

    • Filtragem no nível de RowGroup (o mecanismo de armazenamento ignora blocos de dados não correspondentes).

Limitações

  • O build side deve ter exatamente uma linha. O otimizador gera um ScalarFilter apenas quando o build side tem uma contagem de linhas estatística de 1 e distribuição Replicated. Se o build side não produzir exatamente uma linha em tempo de execução, o Hologres relata um erro RT_CHECK.

  • Se o valor do build side for NULL, o Hologres publica um filtro FILTER_ALL (que filtra todas as linhas), limpa os dados do build side para interromper o Cross Join e retorna um conjunto de resultados vazio.

  • Colunas de dicionário não são suportadas. Se a coluna do probe side referenciada pela target_expr for uma coluna codificada por dicionário (tipo DICTIONARY), o ScalarFilter não será aplicado.

  • Os tipos de dados suportados incluem INT8, UINT8, INT16, UINT16, INT32, UINT32, INT64, UINT64, DATE32, TIMESTAMP, FLOAT, DOUBLE e STRING. Se um tipo de dados não for suportado, o ScalarFilter não é enviado via pushdown e nenhum erro é relatado.

  • Os operadores de comparação suportados incluem >, <, >= e <=. O operador BETWEEN é expandido em duas comparações de intervalo.

Parâmetros GUC

Parâmetro

Tipo

Padrão

Nível

Descrição

hg_experimental_generate_runtime_scalar_filter

bool

true

PGC_USERSET

Controla a geração de ScalarFilter para Cross Joins. Esse recurso está habilitado por padrão e não requer ação do usuário.

Para desabilitar esse recurso, execute o seguinte comando:

SET hg_experimental_generate_runtime_scalar_filter = off;

Novos campos na saída EXPLAIN

Quando o ScalarFilter está habilitado, o nó Cross Join exibe as seguintes informações:

  • Runtime Filter Build Expr: expressão do build side, como (min(t2.a)) e (max(t2.a)).

  • Runtime Filter Target Expr: expressão de filtro alvo no probe side, como (t1.a >= ${1}) AND (t1.a <= ${2}). Nesta expressão, ${N} é um espaço reservado substituído em tempo de execução pelo valor do build side correspondente ao filter_id.

Benefícios de desempenho

No conjunto de dados TPC-DS de 10 TB, o seguinte padrão SQL é comum:

DELETE FROM inventory
WHERE inv_date_sk >= (SELECT min(d_date_sk) FROM date_dim WHERE d_date BETWEEN 'INV_S_1' AND 'INV_E_1')
  AND inv_date_sk <= (SELECT max(d_date_sk) FROM date_dim WHERE d_date BETWEEN 'INV_S_1' AND 'INV_E_1');

Com a otimização ScalarFilter da V4.2, o tempo de consulta cai de 2.672 segundos para 0,552 segundos, representando um aumento de desempenho de aproximadamente 4,8x.

Verificando runtime filters

Os exemplos a seguir demonstram como verifique os efeitos dos runtime filters em diferentes cenários.

Exemplo 1: Condição de junção de coluna única (tipo local)

BEGIN;
CREATE TABLE test1 (x int, y int);
CALL set_table_property('test1', 'distribution_key', 'x');
CREATE TABLE test2 (x int, y int);
CALL set_table_property('test2', 'distribution_key', 'x');
END;
INSERT INTO test1 SELECT t, t FROM generate_series(1, 100000) t;
INSERT INTO test2 SELECT t, t FROM generate_series(1, 1000) t;
ANALYZE test1;
ANALYZE test2;
EXPLAIN ANALYZE SELECT * FROM test1 JOIN test2 ON test1.x = test2.x;

Plano de execução:

QUERY PLAN
Gather  (cost=0.00..10.20 rows=1000 width=16)
[40:1 id=100002 dop=1 time=9/9/9ms rows=1000(1000/1000/1000) mem=16/16/16KB open=2/2/2ms get_next=7/7/7ms]
  -> Hash Join  (cost=0.00..10.16 rows=1000 width=16)
      Hash Cond: (test1.x = test2.x)
      Runtime Filter Cond: (test1.x = test2.x)
      [id=8 dop=40 time=9/4/2ms rows=1000(39/25/17) mem=6/5/5KB open=4/1/0ms get_next=6/2/0ms]
      -> Local Gather  (cost=0.00..5.11 rows=1000000 width=8)
          [id=3 dop=40 time=6/1/0ms rows=1000(39/25/17) mem=600/600/600B open=1/0/0ms get_next=6/1/0ms local_dop=1/1/1]
          -> Seq Scan on test1  (cost=0.00..5.10 rows=1000000 width=8)
              Runtime Filter Target Expr: test1.x
              [id=2 split_count=40 time=11/7/7ms rows=1000(39/25/17) mem=41/41/41KB open=11/7/7ms get_next=0/0/0ms scan_rows=1000000 25270/25000/24697))]
      -> Hash  (cost=5.00..5.00 rows=1000 width=8)
          [id=7 dop=40 time=3/1/0ms rows=1000(39/25/17) mem=396/396/396KB open=3/1/0ms get_next=1/0/0ms rehash=1/1/1 hash_mem=384/384/384KB]
          -> Local Gather  (cost=0.00..5.00 rows=1000 width=8)
              [id=5 dop=40 time=1/0/0ms rows=1000(39/25/17) mem=0/0/0B open=1/0/0ms get_next=1/0/0ms local_dop=0/0/0]
              -> Seq Scan on test2  (cost=0.00..5.00 rows=1000 width=8)
                  [id=4 split_count=40 time=1/0/0ms rows=1000(39/25/17) mem=528/528/528B open=1/0/0ms get_next=1/0/0ms scan_rows=1000(39/25/17)]
  • A tabela test2 tem 1.000 linhas e a tabela test1 tem 100.000 linhas. A proporção de tamanho de dados entre build side e probe side é 0,01 (menos que 0,1), o que atende à condição padrão de ativação para runtime filters.

  • A varredura do probe side em test1 mostra Runtime Filter Target Expr, indicando que um runtime filter foi enviado via pushdown.

  • No probe side, scan_rows é 100.000 (número de linhas lidas do armazenamento), enquanto rows é 1.000 (número de linhas após filtragem). Essa diferença demonstra o efeito da filtragem.

Exemplo 2: Condição de junção de múltiplas colunas (V2.1+, tipo local)

DROP TABLE IF EXISTS test1, test2;
BEGIN;
CREATE TABLE test1 (x int, y int);
CREATE TABLE test2 (x int, y int);
END;
INSERT INTO test1 SELECT t, t FROM generate_series(1, 1000000) t;
INSERT INTO test2 SELECT t, t FROM generate_series(1, 1000) t;
ANALYZE test1;
ANALYZE test2;
EXPLAIN ANALYZE SELECT * FROM test1 JOIN test2 ON test1.x = test2.x AND test1.y = test2.y;

Plano de execução:

QUERY PLAN
Gather  (cost=0.00..10.46 rows=1000 width=16)
[40:1 id=100003 dop=1 time=6/6/6ms rows=1000(1000/1000/1000) mem=600/600/600B open=0/0/0ms get_next=6/6/6ms]
  -> Hash Join  (cost=0.00..10.43 rows=1000 width=16)
       Hash Cond: ((test1.x = test2.x) AND (test1.y = test2.y))
       Runtime Filter Cond: ((test1.x = test2.x) AND (test1.y = test2.y))
       [id=8 dop=40 time=5/3/3ms rows=1000(1000/25/0) mem=40/2/1KB open=1/0/0ms get_next=4/3/3ms]
       -> Local Gather  (cost=0.00..5.11 rows=1000000 width=8)
            [id=5 dop=40 time=4/3/3ms rows=1000(1000/25/0) mem=600/600/600B open=0/0/0ms get_next=4/3/3ms local_dop=1/1/1]
            -> Seq Scan on test1  (cost=0.00..5.10 rows=1000000 width=8)
                 Runtime Filter Target Expr: (test1.x AND test1.y)
                 [id=4 split_count=40 time=7/5/5ms rows=1000(1000/25/0) mem=49/11/9KB open=7/5/5ms get_next=1/0/0ms scan_rows=1000000(32768/25000/24576)]
       -> Hash  (cost=5.02..5.02 rows=40000 width=8)
            [id=7 dop=40 time=1/0/0ms rows=40000(1000/1000/1000) mem=417/417/417KB open=1/0/0ms get_next=1/0/0ms rehash=1/1/1 hash_mem=384/384/384KB]
            -> Broadcast  (cost=0.00..5.02 rows=40000 width=8)
                 [40:40 id=100002 dop=40 time=1/0/0ms rows=40000(1000/1000/1000) mem=0/0/0B open=1/0/0ms get_next=0/0/0ms * ]
                 -> Local Gather  (cost=0.00..5.00 rows=1000 width=8)
                      [id=3 dop=40 time=1/0/0ms rows=1000(1000/25/0) mem=600/600/600B open=1/0/0ms get_next=0/0/0ms local_dop=1/1/1]
                      -> Seq Scan on test2  (cost=0.00..5.00 rows=1000 width=8)
                           [id=2 split_count=40 time=1/0/0ms rows=1000(1000/25/0) mem=528/528/528B open=0/0/0ms get_next=1/0/0ms scan_rows=1000(1000/1000/1000)]
  • A condição de junção contém múltiplas colunas, e o runtime filter também é gerado para múltiplas colunas.

  • Os dados do build side são transmitidos (broadcast), portanto, um runtime filter Local é utilizado.

Exemplo 3: Tipo Global (V2.2+, shuffle join)

SET hg_experimental_enable_result_cache = OFF;
DROP TABLE IF EXISTS test1, test2;
BEGIN;
CREATE TABLE test1 (x int, y int);
CREATE TABLE test2 (x int, y int);
END;
INSERT INTO test1 SELECT t, t FROM generate_series(1, 100000) t;
INSERT INTO test2 SELECT t, t FROM generate_series(1, 1000) t;
ANALYZE test1;
ANALYZE test2;
EXPLAIN ANALYZE SELECT * FROM test1 JOIN test2 ON test1.x = test2.x;

Plano de execução:

QUERY PLAN
   -> Hash Join  (cost=0.00..10.08 rows=1000 width=16)
      Hash Cond: (test1.x = test2.x)
      Runtime Filter Cond: (test1.x = test2.x)
      [id=9 dop=40 time=10/8/8ms rows=1000(34/25/13) mem=6/6/5KB open=2/1/1ms get_next=8/7/7ms]
      -> Redistribution  (cost=0.00..5.07 rows=100000 width=8)
         Hash Key: test1.x
         [40:40 id=100002 dop=40 time=8/7/7ms rows=1289(46/32/20) mem=512/432/0B open=0/0/0ms get_next=8/7/7ms * ]
      -> Local Gather  (cost=0.00..5.01 rows=100000 width=8)
         [id=3 dop=40 time=9/2/0ms rows=1289(1042/32/0) mem=600/600/600B open=0/0/0ms get_next=9/2/0ms local_dop=1/1/1]
         -> Seq Scan on test1  (cost=0.00..5.01 rows=100000 width=8)
            Runtime Filter Target Expr: test1.x
            [id=2 split_count=40 time=11/3/0ms rows=1289(1042/32/0) mem=50512/11849/528B open=11/3/0ms get_next=1/0/0ms scan_rows=100000(8192/7692/1696)]
      -> Hash  (cost=5.00..5.00 rows=1000 width=8)
         [id=8 dop=40 time=2/1/1ms rows=1000(34/25/13) mem=396/396/396KB open=2/1/1ms get_next=0/0/0ms rehash=1/1/1 hash_mem=384/384/384KB]
         -> Redistribution  (cost=0.00..5.00 rows=1000 width=8)
            Hash Key: test2.x
            [40:40 id=100003 dop=40 time=2/1/1ms rows=1000(34/25/13) mem=0/0/0B open=0/0/0ms get_next=2/1/1ms * ]
         -> Local Gather  (cost=0.00..5.00 rows=1000 width=8)
            [id=5 dop=40 time=1/0/0ms rows=1000(1000/25/0) mem=600/600/600B open=1/0/0ms get_next=1/0/0ms local_dop=1/1/1]
            -> Seq Scan on test2  (cost=0.00..5.00 rows=1000 width=8)
               [id=4 split_count=40 time=0/0/0ms rows=1000(1000/25/0) mem=528/528/528B open=0/0/0ms get_next=0/0/0ms scan_rows=1000(1000/1000/1000)]

Os dados do probe side passam por shuffle para o operador Hash Join. O mecanismo usa automaticamente um runtime filter Global para acelerar a consulta.

Exemplo 4: In filter com índice bitmap (V2.2+)

SET hg_experimental_enable_result_cache = OFF;
DROP TABLE IF EXISTS test1, test2;
BEGIN;
CREATE TABLE test1 (x text, y text);
CALL set_table_property('test1', 'distribution_key', 'x');
CALL set_table_property('test1', 'bitmap_columns', 'x');
CALL set_table_property('test1', 'dictionary_encoding_columns', '');
CREATE TABLE test2 (x text, y text);
CALL set_table_property('test2', 'distribution_key', 'x');
END;
INSERT INTO test1 SELECT t::text, t::text FROM generate_series(1, 10000000) t;
INSERT INTO test2 SELECT t::text, t::text FROM generate_series(1, 50) t;
ANALYZE test1;
ANALYZE test2;
EXPLAIN ANALYZE SELECT * FROM test1 JOIN test2 ON test1.x = test2.x;

Plano de execução:

QUERY PLAN
Gather  (cost=0.00..11.70 rows=50 width=14)
[40:1 id=100002 dop=1 time=16/16/16ms rows=50(50/50/50) mem=2/2/2KB open=0/0ms get_next=16/16/16ms]
  -> Hash Join  (cost=0.00..11.70 rows=50 width=14)
       Hash Cond: (test1.x = test2.x)
       Runtime Filter Cond: (test1.x = test2.x)
       [id=7 dop=40 time=15/9/3ms rows=50(3/1/0) mem=5132/3774/264B open=1/0/0ms get_next=14/8/3ms]
       -> Local Gather  (cost=0.00..6.26 rows=10000000 width=12)
            [id=3 dop=40 time=14/8/3ms rows=50(3/1/0) mem=600/600/600B open=1/0/0ms get_next=14/8/3ms local_dop=1/1/1]
            -> Seq Scan on test1  (cost=0.00..6.06 rows=10000000 width=12)
                 Runtime Filter Target Expr: test1.x
                 [id=2 split_count=40 time=16/10/5ms rows=50(3/1/0) mem=67544/48945/528B open=16/9/5ms get_next=1/0/0ms scan_rows=7247692(250875/249920/248982) bitmap_used=50]
       -> Hash  (cost=5.00..5.00 rows=50 width=2)
            [id=6 dop=40 time=1/0/0ms rows=61(3/1/1) mem=534/530/521KB open=1/0/0ms get_next=1/0/0ms rehash=1/1/1 hash_mem=512/512/512KB]
            -> Local Gather  (cost=0.00..5.00 rows=50 width=2)
                 [id=5 dop=40 time=1/0/0ms rows=50(3/1/0) mem=600/600/600B open=0/0/0ms get_next=1/0/0ms local_dop=1/1/1]
                 -> Seq Scan on test2  (cost=0.00..5.00 rows=50 width=2)
                      [id=4 split_count=40 time=1/0/0ms rows=50(3/1/0) mem=528/528/528B open=1/0/0ms get_next=0/0/0ms scan_rows=50(3/1/1)]

O operador de varredura do probe side usa um índice bitmap. O In filter fornece filtragem precisa, deixando apenas 50 linhas. O valor scan_rows no operador de varredura é superior a 7 milhões, inferior aos 10 milhões de linhas originais. Isso ocorre porque o In filter pode ser enviado via pushdown para o mecanismo de armazenamento, reduzindo a sobrecarga de I/O. Combinar um In filter com um índice bitmap proporciona ganhos significativos quando a chave de junção é do tipo STRING.

Exemplo 5: MinMax filter para redução de I/O (V2.2+)

SET hg_experimental_enable_result_cache = OFF;
DROP TABLE IF EXISTS test1, test2;
BEGIN;
CREATE TABLE test1 (x int, y int);
CALL set_table_property('test1', 'distribution_key', 'x');
CREATE TABLE test2 (x int, y int);
CALL set_table_property('test2', 'distribution_key', 'x');
END;
INSERT INTO test1 SELECT t::int, t::int FROM generate_series(1, 10000000) t;
INSERT INTO test2 SELECT t::int, t::int FROM generate_series(1, 100000) t;
ANALYZE test1;
ANALYZE test2;
EXPLAIN ANALYZE SELECT * FROM test1 JOIN test2 ON test1.x = test2.x;

Plano de execução:

QUERY PLAN
 Gather  (cost=0.00..15.68 rows=100000 width=16)
   [40:1 id=100002 dop=1 time=5/5/5ms rows=100000(100000/100000/100000) mem=600/600/600B open=0/0/0ms get_next=5/5/5ms]
   -> Hash Join  (cost=0.00..11.98 rows=100000 width=16)
        Hash Cond: (test1.x = test2.x)
        Runtime Filter Cond: (test1.x = test2.x)
        [id=7 dop=40 time=5/4/4ms rows=100000(2639/2500/2406) mem=97/92/89KB open=1/0/0ms get_next=4/3/3ms]
        -> Local Gather  (cost=0.00..6.14 rows=10000000 width=8)
             [id=3 dop=40 time=5/3/3ms rows=100000(2639/2500/2406) mem=600/600/600B open=1/0/0ms get_next=4/3/3ms local_dop=1/1/1]
             -> Seq Scan on test1  (cost=0.00..6.00 rows=10000000 width=8)
                  Runtime Filter Target Expr: test1.x
                  [id=2 split_count=40 time=6/6/5ms rows=100000(2639/2500/2406) mem=61/60/59KB open=6/5/5ms get_next=0/0/0ms scan_rows=327680(8192/8192/8192)]
        -> Hash  (cost=5.01..5.01 rows=100000 width=8)
             [id=6 dop=40 time=1/0/0ms rows=100000(2639/2500/2406) mem=463/460/458KB open=1/0/0ms get_next=1/0/0ms rehash=1/1/1 hash_mem=384/384/384KB]
             -> Local Gather  (cost=0.00..5.01 rows=100000 width=8)
                  [id=5 dop=40 time=1/0/0ms rows=100000(2639/2500/2406) mem=600/600/600B open=0/0/0ms get_next=1/0/0ms local_dop=1/1/1]
                  -> Seq Scan on test2  (cost=0.00..5.01 rows=100000 width=8)
                       [id=4 split_count=40 time=1/0/0ms rows=100000(2639/2500/2406) mem=528/528/528B open=1/0/0ms get_next=0/0/0ms scan_rows=100000(2639/2500/2406)]

O operador de varredura do probe side lê pouco mais de 320.000 linhas do mecanismo de armazenamento, muito menos que os 10 milhões originais. Isso acontece porque o runtime filter é enviado via pushdown para o mecanismo de armazenamento, que usa os metadados de um lote de dados para filtrar lotes inteiros de uma só vez. Tal comportamento reduz significativamente a sobrecarga de I/O. Este tipo de filtro é mais eficaz quando a chave de junção é numérica e o intervalo de valores do build side é menor que o intervalo do probe side.

Exemplo 6: Runtime filter TopN (V4.0+)

Quando uma instrução SQL inclui um operador topN, o Hologres não calcula todos os resultados. Em vez disso, gera um filtro dinâmico para eliminar dados antecipadamente.

SELECT o_orderkey FROM orders ORDER BY o_orderdate LIMIT 5;

Plano de execução:

QUERY PLAN
Limit  (cost=0.00..116554.70 rows=0 width=8)
  ->  Sort  (cost=0.00..116554.70 rows=100 width=12)
        Sort Key: o_orderdate
      [id=6 dop=1 time=317/317/317ms rows=5(5/5/5) mem=1/1/1KB open=317/317/317ms get_next=0/0/0ms]
        ->  Gather  (cost=0.00..116554.25 rows=100 width=12)
            [20:1 id=100002 dop=1 time=317/317/317ms rows=100(100/100/100) mem=6/6/6KB open=0/0/0ms get_next=317/317/317ms * ]
              ->  Limit  (cost=0.00..116554.25 rows=0 width=12)
                    ->  Sort  (cost=0.00..116554.25 rows=150000000 width=12)
                          Sort Key: o_orderdate
                          Runtime Filter Sort Column: o_orderdate
                        [id=3 dop=20 time=318/282/258ms rows=100(5/5/5) mem=96/96/96KB open=318/282/258ms get_next=1/0/0ms]
                          ->  Local Gather  (cost=0.00..9.59 rows=150000000 width=12)
                              [id=2 dop=20 time=316/280/256ms rows=1372205(68691/68610/68498) mem=0/0/0B open=0/0/0ms get_next=316/280/256ms local_dop=1/1/1 * ]
                                ->  Seq Scan on orders  (cost=0.00..8.24 rows=150000000 width=12)
                                      Runtime Filter Target Expr: o_orderdate
                                    [id=1 split_count=20 time=286/249/222ms rows=1372205(68691/68610/68498) mem=179/179/179KB open=0/0/0ms get_next=286/249/222ms physical_reads=27074(1426/1353/1294) scan_rows=144867963(7324934/7243398/7172304)]
Query id:[1001003033996040311]
QE version: 2.0
Query Queue: init_warehouse.default_queue
======================cost======================
Total cost:[343] ms
Optimizer cost:[13] ms
Build execution plan cost:[0] ms
Init execution plan cost:[6] ms
Start query cost:[6] ms
- Queue cost: [0] ms
- Wait schema cost:[0] ms
- Lock query cost:[0] ms
- Create dataset reader cost:[0] ms
- Create split reader cost:[0] ms
Get result cost:[318] ms
- Get the first block cost:[318] ms
====================resource====================
Memory: total 7 MB. Worker stats: max 3 MB, avg 3 MB, min 3 MB, max memory worker id: 189*****.
CPU time: total 5167 ms. Worker stats: max 2610 ms, avg 2583 ms, min 2557 ms, max CPU time worker id: 189*****.
DAG CPU time stats: max 5165 ms, avg 2582 ms, min 0 ms, cnt 2, max CPU time dag id: 1.
Fragment CPU time stats: max 5137 ms, avg 1721 ms, min 0 ms, cnt 3, max CPU time fragment id: 2.
Ec wait time: total 90 ms. Worker stats: max 46 ms, max(max) 2 ms, avg 45 ms, min 44 ms, max ec wait time worker id: 189*****, max(max) ec wait time worker id: 189*****.
Physical read bytes: total 799 MB. Worker stats: max 400 MB, avg 399 MB, min 399 MB, max physical read bytes worker id: 189*****.
Read bytes: total 898 MB. Worker stats: max 450 MB, avg 449 MB, min 448 MB, max read bytes worker id: 189*****.
DAG instance count: total 3. Worker stats: max 2, avg 1, min 1, max DAG instance count worker id: 189*****.
Fragment instance count: total 41. Worker stats: max 21, avg 20, min 20, max fragment instance count worker id: 189*****.

Sem o runtime filter TopN, o ScanNode leria cada bloco de dados da tabela orders e o passaria para o nó TopN, que então usaria uma ordenação por heap para manter as 5 principais linhas vistas até o momento.

Por exemplo, cada bloco de dados contém aproximadamente 8.192 linhas. Após o processamento do primeiro bloco, o TopN conhece a o_orderdate classificada em 5º lugar naquele bloco. Suponha que seja 1995-01-01. Quando o nó Scan lê o segundo bloco, ele usa 1995-01-01 como condição de filtro e envia apenas linhas onde o_orderdate <= 1995-01-01 para o TopN. O limiar é atualizado dinamicamente. Se a o_orderdate classificada em 5º lugar no segundo bloco for menor, o TopN substitui o limiar antigo pelo novo valor.

Use a saída EXPLAIN para visualize o runtime filter TopN gerado pelo otimizador:

->  Limit  (cost=0.00..116554.25 rows=0 width=12)
  ->  Sort  (cost=0.00..116554.25 rows=150000000 width=12)
        Sort Key: o_orderdate
        Runtime Filter Sort Column: o_orderdate
      [id=3 dop=20 time=318/282/258ms rows=100(5/5/5) mem=96/96/96KB open=318/282/258ms get_next=1/0/0ms]

A presença de Runtime Filter Sort Column no nó TopN indica que este nó gera um runtime filter TopN.

Exemplo 7: Cross join com ScalarFilter (V4.2+)

SET hg_experimental_enable_result_cache = OFF;
DROP TABLE IF EXISTS t1, t2;
BEGIN;
CREATE TABLE t1 (a int, b int);
CREATE TABLE t2 (a int, b int);
END;
INSERT INTO t1 SELECT t, t FROM generate_series(1, 1000000) t;
INSERT INTO t2 SELECT t, t FROM generate_series(1, 1000) t;
ANALYZE t1;
ANALYZE t2;
EXPLAIN ANALYZE
SELECT * FROM t1
WHERE a >= (SELECT min(a) FROM t2 WHERE b BETWEEN 0 AND 1)
  AND a <= (SELECT max(a) FROM t2 WHERE b BETWEEN 0 AND 1);

Destaques do plano de execução:

  • O nó Cross Join mostra Runtime Filter Build Expr: (min(t2.a)), (max(t2.a)).

  • O ScanNode em t1 mostra Runtime Filter Target Expr: (t1.a >= ${1}) AND (t1.a <= ${2}).

  • Para t1, scan_rows reflete o número de linhas lidas do armazenamento, enquanto rows cai significativamente após a aplicação do ScalarFilter. Essa diferença verifica o efeito do filtro.

Desabilitando o recurso para comparação:

SET hg_experimental_generate_runtime_scalar_filter = off;
EXPLAIN ANALYZE
SELECT * FROM t1
WHERE a >= (SELECT min(a) FROM t2 WHERE b BETWEEN 0 AND 1)
  AND a <= (SELECT max(a) FROM t2 WHERE b BETWEEN 0 AND 1);

Após desabilitar o recurso, o plano de execução não mostra mais os campos Runtime Filter Build Expr e Target Expr, e a varredura em t1 reverte para uma varredura completa da tabela.

Histórico de versões

Versão

Novos recursos

V2.0

Suporte a runtime filters em Hash Join (tipo Local, incluindo filtros Bloom, In e MinMax).

V2.1

Suporte a runtime filters com condições de junção de múltiplas colunas.

V2.2

Suporte a runtime filters Globais (para shuffle joins); filtros In podem ser combinados com índices bitmap.

V4.0

Suporte a runtime filter TopN.

V4.2

Suporte a runtime filters Cross Join (ScalarFilter) para otimizar cenários que usam subconsulta escalar para filtrar uma tabela grande.