Query parallele in PostgreSQL

Query parallele in PostgreSQL
Nei moderni CPU ci sono molti core. Da anni le applicazioni inviano richieste ai database in modo parallelo. Se si tratta di una richiesta di report su molte righe in una tabella, viene eseguita più rapidamente quando utilizza più CPU, e in PostgreSQL questo è possibile a partire dalla versione 9.6.

Ci sono voluti 3 anni per implementare la funzione delle richieste parallele — è stato necessario riscrivere il codice in diverse fasi dell'esecuzione delle richieste. Con PostgreSQL 9.6 è stata introdotta un'infrastruttura per migliorare ulteriormente il codice. Nelle versioni successive anche altri tipi di richieste possono essere eseguiti in modo parallelo.

Limitazioni

  • Non attivare l'esecuzione parallela se tutti i core sono già occupati, altrimenti altre richieste subiranno ritardi.
  • La cosa più importante è che l'elaborazione parallela con valori elevati di WORK_MEM utilizza molta memoria — ogni operazione di join hash o ordinamento richiede memoria in un volume pari a work_mem.
  • Le richieste OLTP con bassa latenza non possono essere accelerate tramite esecuzione parallela. E se una richiesta restituisce una sola riga, l'elaborazione parallela la rallenterà.
  • Gli sviluppatori amano utilizzare il benchmark TPC-H. Magari hai delle query simili per un'esecuzione parallela ideale.
  • Solo le query SELECT senza blocco predicativo vengono eseguite in parallelo.
  • A volte un'indicizzazione corretta è migliore della scansione sequenziale della tabella in modalità parallela.
  • Le sospensioni delle query e i cursori non sono supportati.
  • Le funzioni di finestra e le funzioni aggregate degli insiemi ordinati non sono parallele.
  • Non guadagni nulla nella carico di lavoro di input-output.
  • Non esistono algoritmi di ordinamento paralleli. Ma le query con ordinamenti possono essere eseguite in parallelo in alcuni aspetti.
  • Sostituisci CTE (WITH …) con un SELECT annidato per includere l'elaborazione parallela.
  • Le wrapper dei dati di terze parti non supportano ancora l'elaborazione parallela (ma potrebbero farlo!)
  • FULL OUTER JOIN non è supportato.
  • max_rows disattiva l'elaborazione parallela.
  • Se la query contiene una funzione non contrassegnata come PARALLEL SAFE, sarà monothread.
  • Il livello di isolamento della transazione SERIALIZABLE disattiva l'elaborazione parallela.

Ambiente di test

Gli sviluppatori di PostgreSQL hanno cercato di ridurre il tempo di risposta delle query del benchmark TPC-H. Scarica il benchmark e adattalo a PostgreSQL. Questo è un uso non ufficiale del benchmark TPC-H — non per confrontare database o hardware.

  1. Scarica TPC-H_Tools_v2.17.3.zip (o una versione più recente) dal sito ufficiale TPC.
  2. Rinomina makefile.suite in Makefile e modifica come descritto qui: https://github.com/tvondra/pg_tpch . Compila il codice con il comando make.
  3. Genera i dati: ./dbgen -s 10 crea un database di 23 GB. Questo è sufficiente per vedere la differenza nelle prestazioni tra query parallele e non parallele.
  4. Converti i file tbl in csv in for e sed.
  5. Clona il repository pg_tpch e copia i file csv in pg_tpch/dss/data.
  6. Crea le query con il comando qgen.
  7. Carica i dati nel database con il comando ./tpch.sh.

Scansione sequenziale parallela

Può essere più veloce non a causa della lettura parallela, ma perché i dati sono distribuiti su molti core CPU. Nei moderni sistemi operativi, i file di dati di PostgreSQL vengono ben messi in cache. Con la lettura anticipata, è possibile recuperare dallo storage un blocco più grande di quello richiesto dal demone PG. Pertanto, le prestazioni della query non sono limitate dall'I/O del disco. Consuma cicli CPU per:

  • leggere le righe una alla volta dalle pagine della tabella;
  • confrontare i valori delle righe e le condizioni DOVE.

Eseguiamo una semplice query select:

tpch=# explain analyze select l_quantity as sum_qty from lineitem where l_shipdate <= date '1998-12-01' - interval '105' day;
QUERY PLAN
--------------------------------------------------------------------------------------------------------------------------
Seq Scan on lineitem (cost=0.00..1964772.00 rows=58856235 width=5) (actual time=0.014..16951.669 rows=58839715 loops=1)
Filter: (l_shipdate <= '1998-08-18 00:00:00'::timestamp without time zone)
Rows Removed by Filter: 1146337
Planning Time: 0.203 ms
Execution Time: 19035.100 ms

La scansione sequenziale restituisce troppe righe senza aggregazione, quindi la query viene eseguita da un solo core della CPU.

Se aggiungi SUM(), si nota che due processi di lavoro possono aiutare ad accelerare la query:

explain analyze select sum(l_quantity) as sum_qty from lineitem where l_shipdate  Gather (cost=1589701.91..1589702.12 rows=2 width=32) (actual time=8553.241..8555.067 rows=3 loops=1)
Workers Planned: 2
Workers Launched: 2
-> Partial Aggregate (cost=1588701.91..1588701.92 rows=1 width=32) (actual time=8547.546..8547.546 rows=1 loops=3)
-> Parallel Seq Scan on lineitem (cost=0.00..1527393.33 rows=24523431 width=5) (actual time=0.038..5998.417 rows=19613238 loops=3)
Filter: (l_shipdate <= '1998-08-18 00:00:00'::timestamp without time zone)
Rows Removed by Filter: 382112
Planning Time: 0.241 ms
Execution Time: 8555.131 ms

Aggregazione parallela

Il nodo «Parallel Seq Scan» genera righe per l'aggregazione parziale. Il nodo «Partial Aggregate» riduce queste righe utilizzando SUM(). Alla fine, il contatore SUM di ogni processo di lavoro viene raccolto dal nodo «Gather».

Il risultato finale è calcolato dal nodo «Finalize Aggregate». Se hai funzioni di aggregazione personalizzate, non dimenticare di contrassegnarle come «parallel safe».

Numero di processi di lavoro

Il numero di processi di lavoro può essere aumentato senza riavviare il server:

explain analyze select sum(l_quantity) as sum_qty from lineitem where l_shipdate  Gather (cost=1589701.91..1589702.12 rows=2 width=32) (actual time=8553.241..8555.067 rows=3 loops=1)
Workers Planned: 2
Workers Launched: 2
-> Partial Aggregate (cost=1588701.91..1588701.92 rows=1 width=32) (actual time=8547.546..8547.546 rows=1 loops=3)
-> Parallel Seq Scan on lineitem (cost=0.00..1527393.33 rows=24523431 width=5) (actual time=0.038..5998.417 rows=19613238 loops=3)
Filter: (l_shipdate <= '1998-08-18 00:00:00'::timestamp without time zone)
Rows Removed by Filter: 382112
Planning Time: 0.241 ms
Execution Time: 8555.131 ms

Cosa sta succedendo qui? Il numero di processi di lavoro è raddoppiato, mentre la query è diventata solo 1,6599 volte più veloce. I calcoli sono interessanti. Inizialmente avevamo 2 processi di lavoro e 1 leader. Dopo la modifica, abbiamo 4+1.

Il nostro massimo di accelerazione dalla lavorazione parallela: 5/3 = 1,66(6) volte.

Come funziona?

Processi

L'esecuzione della query inizia sempre con il processo leader. Il leader gestisce tutto il lavoro non parallelo e parte del lavoro parallelo. Gli altri processi che eseguono le stesse query sono chiamati processi di lavoro. L'elaborazione parallela utilizza l'infrastruttura di processi di lavoro in background dinamici (dalla versione 9.4). Poiché le altre parti di PostgreSQL utilizzano processi invece di thread, una query con 3 processi di lavoro potrebbe essere 4 volte più veloce dell'elaborazione tradizionale.

Interazione

I processi di lavoro comunicano con il leader tramite una coda di messaggi (basata su memoria condivisa). Ogni processo ha 2 code: per gli errori e per le tuple.

Quanti processi di lavoro sono necessari?

Il limite minimo è stabilito dal parametro max_parallel_workers_per_gather. Poi l'esecutore delle query prende i processi di lavoro da un pool limitato dal parametro max_parallel_workers size. Il limite finale è dato da max_worker_processes, cioè il numero totale di processi in background.

Se non è possibile allocare un processo di lavoro, l'elaborazione sarà monoprocessore.

Il pianificatore delle query può ridurre i processi di lavoro in base alla dimensione della tabella o dell'indice. A questo scopo vi sono i parametri min_parallel_table_scan_size e min_parallel_index_scan_size.

imposta min_parallel_table_scan_size='8MB'
8MB tabella => 1 processo di lavoro
24MB tabella => 2 processi di lavoro
72MB tabella => 3 processi di lavoro
x => log(x / min_parallel_table_scan_size) / log(3) + 1 processo di lavoro

Ogni volta che la tabella è 3 volte più grande di min_parallel_(index|table)_scan_size, Postgres aggiunge un processo di lavoro. Il numero di processi di lavoro non si basa sui costi. La dipendenza circolare rende complesse le implementazioni. Invece, il pianificatore utilizza regole semplici.

In pratica, queste regole non sono sempre adatte per la produzione, quindi è possibile modificare il numero di processi di lavoro per una specifica tabella: ALTER TABLE … SET (parallel_workers = N).

Perché non viene utilizzato il parallelo?

Oltre a un lungo elenco di vincoli, ci sono anche controlli sui costi:

parallel_setup_cost — per evitare il parallelismo per le query brevi. Questo parametro stima il tempo necessario per preparare la memoria, avviare il processo e avviare lo scambio iniziale di dati.

parallel_tuple_cost: la comunicazione tra il leader e i lavoratori può allungarsi proporzionalmente al numero di tuple dai processi di lavoro. Questo parametro stima i costi per lo scambio di dati.

Join a ciclo nidificato — Nested Loop Join

PostgreSQL 9.6+ может выполнять вложенные циклы параллельно — это простая операция.

explain (costs off) select c_custkey, count(o_orderkey)
                from    customer left outer join orders on
                                c_custkey = o_custkey and o_comment not like '%special%deposits%'
                group by c_custkey;
                                      QUERY PLAN
--------------------------------------------------------------------------------------
 Finalize GroupAggregate
   Group Key: customer.c_custkey
   ->  Gather Merge
         Workers Planned: 4
         ->  Partial GroupAggregate
               Group Key: customer.c_custkey
               ->  Nested Loop Left Join
                     ->  Parallel Index Only Scan using customer_pkey on customer
                     ->  Index Scan using idx_orders_custkey on orders
                           Index Cond: (customer.c_custkey = o_custkey)
                           Filter: ((o_comment)::text !~~ '%special%deposits%'::text)

La raccolta avviene nell'ultima fase, quindi il Nested Loop Left Join è un'operazione parallela. Il Parallel Index Only Scan è stato introdotto solo nella versione 10. Funziona in modo simile alla scansione sequenziale parallela. La condizione c_custkey = o_custkey legge un ordine per ogni riga cliente. Quindi non è parallelo.

Hash Join

Ogni processo crea la propria tabella hash fino a PostgreSQL 11. E se ci sono più di quattro processi, le prestazioni non miglioreranno. Nella nuova versione, la tabella hash è condivisa. Ogni processo può utilizzare WORK_MEM per creare la tabella hash.

seleziona
        l_shipmode,
        somma(caso
                quando o_orderpriority = '1-URGENT'
                        o quando o_orderpriority = '2-HIGH'
                        allora 1
                altrimenti 0
        fine) come high_line_count,
        somma(caso
                quando o_orderpriority <> '1-URGENT'
                        e o_orderpriority <> '2-HIGH'
                        allora 1
                altrimenti 0
        fine) come low_line_count
da
        ordini,
        riga_articolo
dove
        o_orderkey = l_orderkey
        e l_shipmode in ('MAIL', 'AIR')
        e l_commitdate < l_receiptdate
        e l_shipdate < l_commitdate
        e l_receiptdate >= data '1996-01-01'
        e l_receiptdate < data '1996-01-01' + intervallo '1' anno
gruppo per
        l_shipmode
ordine per
        l_shipmode
LIMIT 1;
                                                                                                                                    PIANO QUERY
-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------
 Limite  (costo=1964755.66..1964961.44 righe=1 larghezza=27) (tempo effettivo=7579.592..7922.997 righe=1 cicli=1)
   ->  Finalizza GroupAggregate  (costo=1964755.66..1966196.11 righe=7 larghezza=27) (tempo effettivo=7579.590..7579.591 righe=1 cicli=1)
         Chiave di gruppo: lineitem.l_shipmode
         ->  Raccogli Merge  (costo=1964755.66..1966195.83 righe=28 larghezza=27) (tempo effettivo=7559.593..7922.319 righe=6 cicli=1)
               Lavoratori pianificati: 4
               Lavoratori lanciati: 4
               ->  Gruppo Parziale  (costo=1963755.61..1965192.44 righe=7 larghezza=27) (tempo effettivo=7548.103..7564.592 righe=2 cicli=5)
                     Chiave di gruppo: lineitem.l_shipmode
                     ->  Ordina  (costo=1963755.61..1963935.20 righe=71838 larghezza=27) (tempo effettivo=7530.280..7539.688 righe=62519 cicli=5)
                           Chiave di ordinamento: lineitem.l_shipmode
                           Metodo di ordinamento: unione esterna  Disco: 2304kB
                           Lavoratore 0:  Metodo di ordinamento: unione esterna  Disco: 2064kB
                           Lavoratore 1:  Metodo di ordinamento: unione esterna  Disco: 2384kB
                           Lavoratore 2:  Metodo di ordinamento: unione esterna  Disco: 2264kB
                           Lavoratore 3:  Metodo di ordinamento: unione esterna  Disco: 2336kB
                           ->  Unione Hash Parallela  (costo=382571.01..1957960.99 righe=71838 larghezza=27) (tempo effettivo=7036.917..7499.692 righe=62519 cicli=5)
                                 Condizione Hash: (lineitem.l_orderkey = orders.o_orderkey)
                                 ->  Scansione Sequenziale Parallela su lineitem  (costo=0.00..1552386.40 righe=71838 larghezza=19) (tempo effettivo=0.583..4901.063 righe=62519 cicli=5)
                                       Filtro: ((l_shipmode = QUALSIASI ('{MAIL,AIR}'::bpchar[])) E (l_commitdate < l_receiptdate) E (l_shipdate < l_commitdate) E (l_receiptdate >= '1996-01-01'::data) E (l_receiptdate < '1997-01-01 00:00:00'::timestamp senza fuso orario))
                                       Rigetti rimossi dal filtro: 11934691
                                 ->  Hash Parallelo  (costo=313722.45..313722.45 righe=3750045 larghezza=20) (tempo effettivo=2011.518..2011.518 righe=3000000 cicli=5)
                                       Cestini: 65536  Batch: 256  Utilizzo di memoria: 3840kB
                                       ->  Scansione Sequenziale Parallela su ordini  (costo=0.00..313722.45 righe=3750045 larghezza=20) (tempo effettivo=0.029..995.948 righe=3000000 cicli=5)
 Tempo di pianificazione: 0.977 ms
 Tempo di esecuzione: 7923.770 ms

La Query 12 di TPC-H illustra chiaramente una join hash parallela. Ogni processo lavora per creare una tabella hash comune.

Join per fusione — Merge Join

La merge join è intrinsecamente non parallela. Non preoccuparti se è l'ultimo passaggio della query: può comunque essere eseguita in parallelo.

-- Query 2 da TPC-H
spiega (costi esclusi) seleziona s_acctbal, s_name, n_name, p_partkey, p_mfgr, s_address, s_phone, s_comment
da part, supplier, partsupp, nation, region
dove
        p_partkey = ps_partkey
        e s_suppkey = ps_suppkey
        e p_size = 36
        e p_type simile a '%BRASS'
        e s_nationkey = n_nationkey
        e n_regionkey = r_regionkey
        e r_name = 'AMERICA'
        e ps_supplycost = (
                seleziona
                        min(ps_supplycost)
                da partsupp, supplier, nation, region
                dove
                        p_partkey = ps_partkey
                        e s_suppkey = ps_suppkey
                        e s_nationkey = n_nationkey
                        e n_regionkey = r_regionkey
                        e r_name = 'AMERICA'
        )
ordina per s_acctbal discendente, n_name, s_name, p_partkey
LIMIT 100;
                                                PIANO QUERY
----------------------------------------------------------------------------------------------------------
 Limit
   -&gt;  Ordinamento
         Chiave di Ordinamento: supplier.s_acctbal DESC, nation.n_name, supplier.s_name, part.p_partkey
         -&gt;  Merge Join
               Condizione di Merge: (part.p_partkey = partsupp.ps_partkey)
               Filtro di Join: (partsupp.ps_supplycost = (SubPlan 1))
               -&gt;  Raccolta Merge
                     Lavoratori Pianificati: 4
                     -&gt;  Scansione Indice Parallelo usando <strong>part_pkey</strong> sulla parte
                           Filtro: (((p_type)::text ~~ '%BRASS'::text) E (p_size = 36))
               -&gt;  Materializza
                     -&gt;  Ordina
                           Chiave di ordinamento: partsupp.ps_partkey
                           -&gt;  Ciclo annidato
                                 -&gt;  Ciclo annidato
                                       Filtro di join: (nation.n_regionkey = region.r_regionkey)
                                       -&gt;  Scansione sequenziale su region
                                             Filtro: (r_name = 'AMERICA'::bpchar)
                                       -&gt;  Join hash
                                             Condizione hash: (supplier.s_nationkey = nation.n_nationkey)
                                             -&gt;  Scansione sequenziale su supplier
                                             -&gt;  Hash
                                                   -&gt;  Scansione sequenziale su nation
                                 -&gt;  Scansione indice utilizzando idx_partsupp_suppkey su partsupp
                                       Condizione indice: (ps_suppkey = supplier.s_suppkey)
               SottoPiano 1
                 -&gt;  Aggregare
                       -&gt;  Ciclo annidato
                             Filtro di join: (nation_1.n_regionkey = region_1.r_regionkey)
                             -&gt;  Scansione sequenziale su region region_1
                                   Filtro: (r_name = 'AMERICA'::bpchar)
                             -&gt;  Ciclo annidato
                                   -&gt;  Ciclo annidato
                                         -&gt;  Scansione indice utilizzando idx_partsupp_partkey su partsupp partsupp_1
                                               Condizione indice: (part.p_partkey = ps_partkey)
                                         -&gt;  Scansione indice utilizzando supplier_pkey su supplier supplier_1
                                               Condizione indice: (s_suppkey = partsupp_1.ps_suppkey)
                                   -&gt;  Scansione indice utilizzando nation_pkey su nation nation_1
                                         Condizione indice: (n_nationkey = supplier_1.s_nationkey)

Il nodo 'Merge Join' è sopra 'Gather Merge'. Pertanto, la fusione non utilizza l'elaborazione parallela. Tuttavia, il nodo 'Parallel Index Scan' supporta ancora il segmento. part_pkey.

Join per sezioni

in PostgreSQL 11 join per sezioni è disattivato per impostazione predefinita: ha una pianificazione molto costosa. Le tabelle con sezionamenti simili possono essere unite sezione per sezione. In questo modo, Postgres utilizzerà tabelle hash più piccole. Ogni join di sezioni può essere parallelo.

tpch=# set enable_partitionwise_join=t;
tpch=# explain (costs off) select * from prt1 t1, prt2 t2
where t1.a = t2.b and t1.b = 0 and t2.b between 0 and 10000;
                    QUERY PLAN
---------------------------------------------------
 Append
   ->  Hash Join
         Hash Cond: (t2.b = t1.a)
         ->  Seq Scan on prt2_p1 t2
               Filter: ((b >= 0) AND (b   Hash
               ->  Seq Scan on prt1_p1 t1
                     Filter: (b = 0)
   ->  Hash Join
         Hash Cond: (t2_1.b = t1_1.a)
         ->  Seq Scan on prt2_p2 t2_1
               Filter: ((b >= 0) AND (b   Hash
               ->  Seq Scan on prt1_p2 t1_1
                     Filter: (b = 0)
tpch=# set parallel_setup_cost = 1;
tpch=# set parallel_tuple_cost = 0.01;
tpch=# explain (costs off) select * from prt1 t1, prt2 t2
where t1.a = t2.b and t1.b = 0 and t2.b between 0 and 10000;
                        QUERY PLAN
-----------------------------------------------------------
 Gather
   Workers Planned: 4
   ->  Parallel Append
         ->  Parallel Hash Join
               Hash Cond: (t2_1.b = t1_1.a)
               ->  Parallel Seq Scan on prt2_p2 t2_1
                     Filter: ((b >= 0) AND (b   Parallel Hash
                     ->  Parallel Seq Scan on prt1_p2 t1_1
                           Filter: (b = 0)
         ->  Parallel Hash Join
               Hash Cond: (t2.b = t1.a)
               ->  Parallel Seq Scan on prt2_p1 t2
                     Filter: ((b >= 0) AND (b   Parallel Hash
                     ->  Parallel Seq Scan on prt1_p1 t1
                           Filter: (b = 0)

È importante notare che la connessione per sezioni può avvenire in parallelo solo se queste sezioni sono abbastanza grandi.

Appendice Parallelo — Parallel Append

Parallel Append può essere utilizzato al posto di diversi blocchi in vari flussi di lavoro. Questo è comune con le query UNION ALL. Lo svantaggio è una riduzione del parallelismo, poiché ogni flusso di lavoro gestisce solo 1 query.

Sono stati avviati 2 flussi di lavoro, anche se sono stati attivati 4.

tpch=# explain (costs off) select sum(l_quantity) as sum_qty from lineitem where l_shipdate <= date '1998-12-01' - interval '105' day union all select sum(l_quantity) as sum_qty from lineitem where l_shipdate <= date '2000-12-01' - interval '105' day;
                                           QUERY PLAN
------------------------------------------------------------------------------------------------
 Gather
   Workers Planned: 2
   ->  Parallel Append
         ->  Aggregate
               ->  Seq Scan on lineitem
                     Filter: (l_shipdate <= '2000-08-18 00:00:00'::timestamp without time zone)
         ->  Aggregate
               ->  Seq Scan on lineitem lineitem_1
                     Filter: (l_shipdate <= '1998-08-18 00:00:00'::timestamp without time zone)

Le variabili più importanti

  • WORK_MEM limita la quantità di memoria per ogni processo, non solo per le query: work_mem processi collegamenti = molta memoria.
  • max_parallel_workers_per_gather — quanti flussi di lavoro utilizzerà il programma per l'elaborazione parallela dal piano.
  • max_worker_processes — adatta il numero totale di flussi di lavoro in base al numero di core CPU nel server.
  • max_parallel_workers — lo stesso, ma per processi di lavoro paralleli.

Risultati

A partire dalla versione 9.6, l'elaborazione parallela può migliorare notevolmente le prestazioni delle query complesse che scansionano molte righe o indici. In PostgreSQL 10, l'elaborazione parallela è abilitata di default. Ricordate di disattivarla sui server con carichi di lavoro OLTP elevati. Le scansioni sequenziali o le scansioni degli indici consumano molte risorse. Se non state eseguendo un report su un intero set di dati, le query possono essere rese più efficienti semplicemente aggiungendo gli indici mancanti o utilizzando una corretta partizionamento.

Link

Fonte: habr.com

Acquista un hosting affidabile per siti web con protezione DDoS, VPS VDS server 🔥 Acquista un hosting affidabile per siti web con protezione DDoS, VPS VDS server | ProHoster