Why PostgreSQL Uses Two Hash Operations for JOIN and GROUP BY

You might have noticed PostgreSQL handling queries having Join and Group by clauses. If you have not, this is how PostgreSQL does the operation: first join a big fact table to a small dimension table, then GROUP BY the join key. And most of the reporting queries have this query shape.

Suppose this is the query we have.

SELECT aml.journal_id, aj.name, count(*), sum(aml.debit)
FROM account_move_line aml
LEFT JOIN account_journal aj ON aml.journal_id = aj.id
GROUP BY aml.journal_id, aj.name;

And postgrs executes this as two separate operations stacked on top of each other. That is a hash join followed by a HashAggregate. Each builds its own hash table, over the same key. Now, we have a new operator that does these two operations together as one called the Group Join.

Why do we have two separate Operators?

This is the actual plan: a query on a table with 91516 account move lines joined against 76 journals:

HashAggregate (actual time=736.631..736.679 rows=14.00 loops=1)
  Group Key: aml.journal_id, aj.name
  Batches: 1  Memory Usage: 40kB
  Buffers: shared hit=2629
  ->  Hash Left Join (actual time=0.418..508.233 rows=91516.00 loops=1)
        Hash Cond: (aml.journal_id = aj.id)
        Buffers: shared hit=2629
        ->  Seq Scan on account_move_line aml (actual time=0.023..158.383 rows=91516.00 loops=1)
              Buffers: shared hit=2627
        ->  Hash (actual time=0.343..0.349 rows=76.00 loops=1)
              Buckets: 1024  Batches: 1  Memory Usage: 14kB
              Buffers: shared hit=2
              ->  Seq Scan on account_journal aj (actual time=0.008..0.177 rows=76.00 loops=1)
                    Buffers: shared hit=2
Execution Time: 736.946 ms

Read this from bottom up; the hash join builds a hash table over the left side (account_journal with 76 rows) and then streams all 91516 account_move_linerows through it, which produces a joined row for each match. Then, those rows go through a hash aggregate, which creates a second hash table, which is keyed on the exact same field, journal_id, in order to get the count(*) and sum() to 14 groups.

The idea: one hash table doing both jobs

The structure of the hash join’s hash table is such that it can be used as an aggregate’s accumulator table too. It consists of one henry per build side row, and the reason why the same rows end up in the same bucket is because of their key.

Instead of creating an entry that holds only the matching row, what the new group join operator does is allocate extra memory for each entry to hold a set of aggregate states, including a running sum, running count, and whatever the aggregate requires. In this way, no materialization of the joined rows is required. No additional hash table needs to be created.

This is the basic concept, and all other steps are taken care by the planner to prove that the group join operation is possible under certain conditions.

The conditions to satisfy the Group Join

In order to reduce one hash entry and one group into one same thing, it requires proving that they can have the same partition of data. This is the check required by the planner before moving forward with this plan.

There must be a Hash Join on which the fusion should take place. The query cannot have any DISTINCT, ORDER BY sensitive and grouping set aggregate. Since they require sorted data in the hash table.

The GROUP BY key columns must exactly match the key columns of the hash join. That means for each and every hash join key column, GROUP BY must reference either the build column or the probe column of it. Any irrelevant probe side column is unacceptable, as it might change between two rows in the same bucket; hence, splitting the hash entry into two groups is impossible.

All join keys on the build side have to be provably NOT NULL – a unique index allows NULLs. In joins with the matching condition, NULL will never be equal to anything, so there is nothing to worry about. The GROUP BY does the opposite - it aggregates all NULLs into one group. A unique key containing NULLs is not safe – two NULL rows on the build side could be individually unique but end up in the same GROUP BY group.

The join type has to allow accumulators on the build side – inner join and the more common right join are always OK. However, in the case of a left join, we need the planner to prove that there is a validated foreign key constraint between the probe column and the build column (and that there are no additional join conditions). So the only case when a left join can result in a NULL value in the probe column is when there is a NULL in the build column, and that will use one reserved accumulator.

The build side has to be unique on the hash-join columns – those are the columns that determine whether the rows have the same hash entry.

What does the executor actually do?

Once the planner confirms all of the above and selects this plan as a cost-effective join method, execution goes like this:

  1. Build – the usual hash table build on the smaller side with an additional pass to initialize transition states for each bucket entry to their initial value.
  2. Probe - for each row on the larger side, probe its corresponding bucket entry and update the transition states of that entry. In this phase, we only update the rows, so the number of output rows will be zero.
  3. Emit – iterate over all bucket entries, not just those which were matched, finalize their transition states, apply the HAVING clause, and emit the resulting row.

When the build side does not fit into the work_mem, it will spill to batches in the same way as regular hash joins do, and this is safe in our case, because the batch number of the row is a pure function of its hash value calculated using the join key, which is equal to the grouping key. Thus, two rows of one group will not fall into different batches, and this allows us to finalize them before the next batch arrives.

Where does the win actually come from?

Same query, same data; only the new group join option is enabled here.

Hash Group Left Join (actual time=279.776..279.813 rows=14.00 loops=1)
  Hash Cond: (aml.journal_id = aj.id)
  Group Key: aj.id, aj.name
  Buffers: shared hit=2629
  ->  Seq Scan on account_move_line aml (actual time=0.015..126.657 rows=91516.00 loops=1)
        Buffers: shared hit=2627
  ->  Hash (actual time=0.300..0.310 rows=76.00 loops=1)
        Buckets: 1024  Batches: 1  Memory Usage: 16kB
        Buffers: shared hit=2
        ->  Seq Scan on account_journal aj (actual time=0.009..0.157 rows=76.00 loops=1)
              Buffers: shared hit=2
Execution Time: 279.970 ms

In other words, one plan node is used instead of two here, and therefore 736.9 ms becomes 280.0 ms. In addition, note that the shared buffer hit (2629) is also the same for two cases, and scans of both tables have also been changed. But the difference is that the line of hashAggregate with its memory consumption (40 KB) and build phase is now missing. This optimization is not only related to the optimization of aggregation, but the whole operation has disappeared.

As the data set increases, the additional hash table becomes even more expensive. While the regular plan may need to spill it to disk several times, Group Join never needs to construct this second hash table at all. This saves time and resources.

There is nothing wrong with the operator plan postgresql creates for the "join and then group by the join key" query – it just pays for the same information twice. The Group Join realization is an understanding that a hash join's hash table and a hash aggregate's hash table are the same. For this shape, where the same table wears two hats – one of hash entry and one of group with accumulator – we can

WhatsApp