agora inbox for pgsql-hackers@postgresql.org  
help / color / mirror / Atom feed
From: Jehan-Guillaume de Rorthais <jgdr@dalibo.com>
To: Tomas Vondra <tomas.vondra@enterprisedb.com>
Cc: Melanie Plageman <melanieplageman@gmail.com>
Cc: pgsql-hackers@lists.postgresql.org
Subject: Re: Memory leak from ExecutorState context?
Date: Sat, 8 Apr 2023 02:01:19 +0200
Message-ID: <20230408020119.32a0841b@karst> (raw)
In-Reply-To: <20230331140611.3695d30d@karst>
References: <3013398b-316c-638f-2a73-3783e8e2ef02@enterprisedb.com>
	<20230302001827.66e95dc3@karst>
	<41c5766d-ed71-b70c-bbbc-d3396c462d62@enterprisedb.com>
	<20230302130838.717e888d@karst>
	<77a96d42-00cb-2448-465a-aa1e92d00cac@enterprisedb.com>
	<20230302191530.781909fe@karst>
	<dbae24d7-0dda-18aa-5e08-8138ac1caef9@enterprisedb.com>
	<20230310195114.6d0c5406@karst>
	<20230317091834.22e97642@karst>
	<ae017eef-79d5-fcd6-b865-b7e55ac5290b@enterprisedb.com>
	<ZBdjJ8l3CNmBZUg0@telsasoft.com>
	<455abe0e-91b2-f428-6f4c-b95c7c8dfb52@enterprisedb.com>
	<20230320151234.38b2235e@karst>
	<20230327231323.08277083@karst>
	<f548a8d6-9ac2-46a9-1197-585ad79c1797@enterprisedb.com>
	<20230328151745.0f6061f8@karst>
	<f532d33d-e4c1-ed00-e047-94f3aafa49db@enterprisedb.com>
	<20230331140611.3695d30d@karst>

On Fri, 31 Mar 2023 14:06:11 +0200
Jehan-Guillaume de Rorthais <jgdr@dalibo.com> wrote:

> > [...]
> > >> Hmmm, not sure is WARNING is a good approach, but I don't have a better
> > >> idea at the moment.  
> > > 
> > > I stepped it down to NOTICE and added some more infos.
> > > 
> > > [...]
> > >   NOTICE:  Growing number of hash batch to 32768 is exhausting allowed
> > > memory (137494656 > 2097152)
> > [...]
> > 
> > OK, although NOTICE that may actually make it less useful - the default
> > level is WARNING, and regular users are unable to change the level. So
> > very few people will actually see these messages.
[...]
> Anyway, maybe this should be added in the light of next patch, balancing
> between increasing batches and allowed memory. The WARNING/LOG/NOTICE message
> could appears when we actually break memory rules because of some bad HJ
> situation.

So I did some more minor editions to the memory context patch and start working
on the balancing memory patch. Please, find in attachment the v4 patch set:

* 0001-v4-Describe-hybrid-hash-join-implementation.patch:
  Adds documentation written by Melanie few years ago
* 0002-v4-Allocate-hash-batches-related-BufFile-in-a-dedicated.patch:
  The batches' BufFile dedicated memory context patch
* 0003-v4-Add-some-debug-and-metrics.patch:
  A pure debug patch I use to track memory in my various tests
* 0004-v4-Limit-BufFile-memory-explosion-with-bad-HashJoin.patch
  A new and somewhat different version of the balancing memory patch, inspired
  from Tomas work.

After rebasing Tomas' memory balancing patch, I did some memory measures
to answer some of my questions. Please, find in attachment the resulting charts
"HJ-HEAD.png" and "balancing-v3.png" to compare memory consumption between HEAD
and Tomas' patch. They shows an alternance of numbers before/after calling
ExecHashIncreaseNumBatches (see the debug patch). I didn't try to find the
exact last total peak of memory consumption during the join phase and before
all the BufFiles are destroyed. So the last number might be underestimated.

Looking at Tomas' patch, I was quite surprised to find that data+bufFile
actually didn't fill memory up to spaceAllowed before splitting the batches and
rising the memory limit. This is because the patch assume the building phase
consume inner and outer BufFiles equally, where only the inner side is really
allocated. That's why the peakMB value is wrong compared to actual bufFileMB
measured.

So I worked on the v4 patch were BufFile are accounted in spaceUsed. Moreover,
instead of rising the limit and splitting the batches in the same step, the
patch first rise the memory limit if needed, then split in a later call if we
have enough room. The "balancing-v4.png" chart shows the resulting memory
activity. We might need to discuss the proper balancing between memory
consumption and batches.

Note that the patch now log a message when breaking the work_mem. Eg.:

  WARNING:  Hash Join node must grow outside of work_mem
  DETAIL:  Rising memory limit from 4194304 to 6291456
  HINT:  You might need to ANALYZE your table or tune its statistics collection.

Regards,

Attachments:

  [image/png] balancing-v3.png (63.8K, ../20230408020119.32a0841b@karst/2-balancing-v3.png)
  download | view image

  [application/octet-stream] balancing-v3.data (1.3K, ../20230408020119.32a0841b@karst/3-balancing-v3.data)
  download

  [image/png] balancing-v4.png (61.4K, ../20230408020119.32a0841b@karst/4-balancing-v4.png)
  download | view image

  [application/octet-stream] balancing-v4.data (1.6K, ../20230408020119.32a0841b@karst/5-balancing-v4.data)
  download

  [image/png] HJ-HEAD.png (47.5K, ../20230408020119.32a0841b@karst/6-HJ-HEAD.png)
  download | view image

  [application/octet-stream] HJ-HEAD.data (1.5K, ../20230408020119.32a0841b@karst/7-HJ-HEAD.data)
  download

  [application/sql] test.sql (711B, ../20230408020119.32a0841b@karst/8-test.sql)
  download

  [text/x-patch] 0001-v4-Describe-hybrid-hash-join-implementation.patch (2.8K, ../20230408020119.32a0841b@karst/9-0001-v4-Describe-hybrid-hash-join-implementation.patch)
  download | inline diff:
From 46de733d095453615b69a67ba457c86551891d02 Mon Sep 17 00:00:00 2001
From: Melanie Plageman <melanieplageman@gmail.com>
Date: Thu, 30 Apr 2020 07:16:28 -0700
Subject: [PATCH 1/4] Describe hybrid hash join implementation

This is just a draft to spark conversation on what a good comment might
be like in this file on how the hybrid hash join algorithm is
implemented in Postgres. I'm pretty sure this is the accepted term for
this algorithm https://en.wikipedia.org/wiki/Hash_join#Hybrid_hash_join
---
 src/backend/executor/nodeHashjoin.c | 36 +++++++++++++++++++++++++++++
 1 file changed, 36 insertions(+)

diff --git a/src/backend/executor/nodeHashjoin.c b/src/backend/executor/nodeHashjoin.c
index 52ed05c6f5..5454afbff2 100644
--- a/src/backend/executor/nodeHashjoin.c
+++ b/src/backend/executor/nodeHashjoin.c
@@ -10,6 +10,42 @@
  * IDENTIFICATION
  *	  src/backend/executor/nodeHashjoin.c
  *
+ *   HYBRID HASH JOIN
+ *
+ *  If the inner side tuples of a hash join do not fit in memory, the hash join
+ *  can be executed in multiple batches.
+ *
+ *  If the statistics on the inner side relation are accurate, planner chooses a
+ *  multi-batch strategy and estimates the number of batches.
+ *
+ *  The query executor measures the real size of the hashtable and increases the
+ *  number of batches if the hashtable grows too large.
+ *
+ *  The number of batches is always a power of two, so an increase in the number
+ *  of batches doubles it.
+ *
+ *  Serial hash join measures batch size lazily -- waiting until it is loading a
+ *  batch to determine if it will fit in memory. While inserting tuples into the
+ *  hashtable, serial hash join will, if that tuple were to exceed work_mem,
+ *  dump out the hashtable and reassign them either to other batch files or the
+ *  current batch resident in the hashtable.
+ *
+ *  Parallel hash join, on the other hand, completes all changes to the number
+ *  of batches during the build phase. If it increases the number of batches, it
+ *  dumps out all the tuples from all batches and reassigns them to entirely new
+ *  batch files. Then it checks every batch to ensure it will fit in the space
+ *  budget for the query.
+ *
+ *  In both parallel and serial hash join, the executor currently makes a best
+ *  effort. If a particular batch will not fit in memory, it tries doubling the
+ *  number of batches. If after a batch increase, there is a batch which
+ *  retained all or none of its tuples, the executor disables growth in the
+ *  number of batches globally. After growth is disabled, all batches that would
+ *  have previously triggered an increase in the number of batches instead
+ *  exceed the space allowed.
+ *
+ *  TODO: should we discuss that tuples can only spill forward?
+ *
  * PARALLELISM
  *
  * Hash joins can participate in parallel query execution in several ways.  A
-- 
2.39.2

  [text/x-patch] 0002-v4-Allocate-hash-batches-related-BufFile-in-a-dedicated.patch (8.9K, ../20230408020119.32a0841b@karst/10-0002-v4-Allocate-hash-batches-related-BufFile-in-a-dedicated.patch)
  download | inline diff:
From fe050ef3a451660f90a3801e759aa4257a9abfd1 Mon Sep 17 00:00:00 2001
From: Jehan-Guillaume de Rorthais <jgdr@dalibo.com>
Date: Mon, 27 Mar 2023 15:54:39 +0200
Subject: [PATCH 2/4] Allocate hash batches related BufFile in a dedicated
 context

---
 src/backend/executor/nodeHash.c     | 35 +++++++++++++++++++++--------
 src/backend/executor/nodeHashjoin.c | 18 ++++++++++-----
 src/include/executor/hashjoin.h     | 15 ++++++++++---
 src/include/executor/nodeHashjoin.h |  2 +-
 4 files changed, 52 insertions(+), 18 deletions(-)

diff --git a/src/backend/executor/nodeHash.c b/src/backend/executor/nodeHash.c
index a45bd3a315..4544296391 100644
--- a/src/backend/executor/nodeHash.c
+++ b/src/backend/executor/nodeHash.c
@@ -484,7 +484,7 @@ ExecHashTableCreate(HashState *state, List *hashOperators, List *hashCollations,
 	 *
 	 * The hashtable control block is just palloc'd from the executor's
 	 * per-query memory context.  Everything else should be kept inside the
-	 * subsidiary hashCxt or batchCxt.
+	 * subsidiary hashCxt, batchCxt or fileCxt.
 	 */
 	hashtable = palloc_object(HashJoinTableData);
 	hashtable->nbuckets = nbuckets;
@@ -538,6 +538,10 @@ ExecHashTableCreate(HashState *state, List *hashOperators, List *hashCollations,
 												"HashBatchContext",
 												ALLOCSET_DEFAULT_SIZES);
 
+	hashtable->fileCxt = AllocSetContextCreate(hashtable->hashCxt,
+								 "HashBatchFiles",
+								 ALLOCSET_DEFAULT_SIZES);
+
 	/* Allocate data that will live for the life of the hashjoin */
 
 	oldcxt = MemoryContextSwitchTo(hashtable->hashCxt);
@@ -570,15 +574,21 @@ ExecHashTableCreate(HashState *state, List *hashOperators, List *hashCollations,
 
 	if (nbatch > 1 && hashtable->parallel_state == NULL)
 	{
+		MemoryContext oldctx;
+
 		/*
 		 * allocate and initialize the file arrays in hashCxt (not needed for
 		 * parallel case which uses shared tuplestores instead of raw files)
 		 */
+		oldctx = MemoryContextSwitchTo(hashtable->fileCxt);
+
 		hashtable->innerBatchFile = palloc0_array(BufFile *, nbatch);
 		hashtable->outerBatchFile = palloc0_array(BufFile *, nbatch);
 		/* The files will not be opened until needed... */
 		/* ... but make sure we have temp tablespaces established for them */
 		PrepareTempTablespaces();
+
+		MemoryContextSwitchTo(oldctx);
 	}
 
 	MemoryContextSwitchTo(oldcxt);
@@ -913,7 +923,6 @@ ExecHashIncreaseNumBatches(HashJoinTable hashtable)
 	int			oldnbatch = hashtable->nbatch;
 	int			curbatch = hashtable->curbatch;
 	int			nbatch;
-	MemoryContext oldcxt;
 	long		ninmemory;
 	long		nfreed;
 	HashMemoryChunk oldchunks;
@@ -934,13 +943,16 @@ ExecHashIncreaseNumBatches(HashJoinTable hashtable)
 		   hashtable, nbatch, hashtable->spaceUsed);
 #endif
 
-	oldcxt = MemoryContextSwitchTo(hashtable->hashCxt);
-
 	if (hashtable->innerBatchFile == NULL)
 	{
+		MemoryContext oldcxt = MemoryContextSwitchTo(hashtable->fileCxt);
+
 		/* we had no file arrays before */
 		hashtable->innerBatchFile = palloc0_array(BufFile *, nbatch);
 		hashtable->outerBatchFile = palloc0_array(BufFile *, nbatch);
+
+		MemoryContextSwitchTo(oldcxt);
+
 		/* time to establish the temp tablespaces, too */
 		PrepareTempTablespaces();
 	}
@@ -951,8 +963,6 @@ ExecHashIncreaseNumBatches(HashJoinTable hashtable)
 		hashtable->outerBatchFile = repalloc0_array(hashtable->outerBatchFile, BufFile *, oldnbatch, nbatch);
 	}
 
-	MemoryContextSwitchTo(oldcxt);
-
 	hashtable->nbatch = nbatch;
 
 	/*
@@ -1022,9 +1032,11 @@ ExecHashIncreaseNumBatches(HashJoinTable hashtable)
 			{
 				/* dump it out */
 				Assert(batchno > curbatch);
+
 				ExecHashJoinSaveTuple(HJTUPLE_MINTUPLE(hashTuple),
 									  hashTuple->hashvalue,
-									  &hashtable->innerBatchFile[batchno]);
+									  &hashtable->innerBatchFile[batchno],
+									  hashtable->fileCxt);
 
 				hashtable->spaceUsed -= hashTupleSize;
 				nfreed++;
@@ -1681,9 +1693,11 @@ ExecHashTableInsert(HashJoinTable hashtable,
 		 * put the tuple into a temp file for later batches
 		 */
 		Assert(batchno > hashtable->curbatch);
+
 		ExecHashJoinSaveTuple(tuple,
 							  hashvalue,
-							  &hashtable->innerBatchFile[batchno]);
+							  &hashtable->innerBatchFile[batchno],
+							  hashtable->fileCxt);
 	}
 
 	if (shouldFree)
@@ -2663,8 +2677,11 @@ ExecHashRemoveNextSkewBucket(HashJoinTable hashtable)
 		{
 			/* Put the tuple into a temp file for later batches */
 			Assert(batchno > hashtable->curbatch);
+
 			ExecHashJoinSaveTuple(tuple, hashvalue,
-								  &hashtable->innerBatchFile[batchno]);
+								  &hashtable->innerBatchFile[batchno],
+								  hashtable->fileCxt);
+
 			pfree(hashTuple);
 			hashtable->spaceUsed -= tupleSize;
 			hashtable->spaceUsedSkew -= tupleSize;
diff --git a/src/backend/executor/nodeHashjoin.c b/src/backend/executor/nodeHashjoin.c
index 5454afbff2..5bc7f814c6 100644
--- a/src/backend/executor/nodeHashjoin.c
+++ b/src/backend/executor/nodeHashjoin.c
@@ -485,8 +485,10 @@ ExecHashJoinImpl(PlanState *pstate, bool parallel)
 					 */
 					Assert(parallel_state == NULL);
 					Assert(batchno > hashtable->curbatch);
+
 					ExecHashJoinSaveTuple(mintuple, hashvalue,
-										  &hashtable->outerBatchFile[batchno]);
+										  &hashtable->outerBatchFile[batchno],
+										  hashtable->fileCxt);
 
 					if (shouldFree)
 						heap_free_minimal_tuple(mintuple);
@@ -1297,21 +1299,27 @@ ExecParallelHashJoinNewBatch(HashJoinState *hjstate)
  * The data recorded in the file for each tuple is its hash value,
  * then the tuple in MinimalTuple format.
  *
- * Note: it is important always to call this in the regular executor
- * context, not in a shorter-lived context; else the temp file buffers
- * will get messed up.
+ * Note: it is important always to call this in the HashBatchFiles context,
+ * not in a shorter-lived context; else the temp file buffers will get messed
+ * up.
  */
 void
 ExecHashJoinSaveTuple(MinimalTuple tuple, uint32 hashvalue,
-					  BufFile **fileptr)
+					  BufFile **fileptr, MemoryContext filecxt)
 {
 	BufFile    *file = *fileptr;
 
 	if (file == NULL)
 	{
+		MemoryContext oldctx;
+
+		oldctx = MemoryContextSwitchTo(filecxt);
+
 		/* First write to this batch file, so open it. */
 		file = BufFileCreateTemp(false);
 		*fileptr = file;
+
+		MemoryContextSwitchTo(oldctx);
 	}
 
 	BufFileWrite(file, &hashvalue, sizeof(uint32));
diff --git a/src/include/executor/hashjoin.h b/src/include/executor/hashjoin.h
index 8ee59d2c71..74867c3e40 100644
--- a/src/include/executor/hashjoin.h
+++ b/src/include/executor/hashjoin.h
@@ -25,10 +25,14 @@
  *
  * Each active hashjoin has a HashJoinTable control block, which is
  * palloc'd in the executor's per-query context.  All other storage needed
- * for the hashjoin is kept in private memory contexts, two for each hashjoin.
+ * for the hashjoin is kept in private memory contexts, three for each
+ * hashjoin:
+ * - HashTableContext (hashCxt): the control block associated to the hash table
+ * - HashBatchContext (batchCxt): storages for batches
+ * - HashBatchFiles (fileCxt): storage for temp files buffers
+ *
  * This makes it easy and fast to release the storage when we don't need it
- * anymore.  (Exception: data associated with the temp files lives in the
- * per-query context too, since we always call buffile.c in that context.)
+ * anymore.
  *
  * The hashtable contexts are made children of the per-query context, ensuring
  * that they will be discarded at end of statement even if the join is
@@ -39,6 +43,10 @@
  * "hashCxt", while storage that is only wanted for the current batch is
  * allocated in the "batchCxt".  By resetting the batchCxt at the end of
  * each batch, we free all the per-batch storage reliably and without tedium.
+ * Note that data associated with the temp files lives in the "fileCxt" context
+ * which lives during the entire join as temp files might need to survives
+ * batches. These files are explicitly destroyed by calling BufFileClose()
+ * when the code is done with them.
  *
  * During first scan of inner relation, we get its tuples from executor.
  * If nbatch > 1 then tuples that don't belong in first batch get saved
@@ -350,6 +358,7 @@ typedef struct HashJoinTableData
 
 	MemoryContext hashCxt;		/* context for whole-hash-join storage */
 	MemoryContext batchCxt;		/* context for this-batch-only storage */
+	MemoryContext fileCxt;		/* context for the BufFile related storage */
 
 	/* used for dense allocation of tuples (into linked chunks) */
 	HashMemoryChunk chunks;		/* one list for the whole batch */
diff --git a/src/include/executor/nodeHashjoin.h b/src/include/executor/nodeHashjoin.h
index d367070883..a8f9ae1989 100644
--- a/src/include/executor/nodeHashjoin.h
+++ b/src/include/executor/nodeHashjoin.h
@@ -29,6 +29,6 @@ extern void ExecHashJoinInitializeWorker(HashJoinState *state,
 										 ParallelWorkerContext *pwcxt);
 
 extern void ExecHashJoinSaveTuple(MinimalTuple tuple, uint32 hashvalue,
-								  BufFile **fileptr);
+								  BufFile **fileptr, MemoryContext filecxt);
 
 #endif							/* NODEHASHJOIN_H */
-- 
2.39.2

  [text/x-patch] 0003-v4-Add-some-debug-and-metrics.patch (2.3K, ../20230408020119.32a0841b@karst/11-0003-v4-Add-some-debug-and-metrics.patch)
  download | inline diff:
From 606c3c8604acebe4291c21a8fd2a1bf1dddcd9d8 Mon Sep 17 00:00:00 2001
From: Jehan-Guillaume de Rorthais <jgdr@dalibo.com>
Date: Tue, 4 Apr 2023 16:24:40 +0200
Subject: [PATCH 3/4] Add some debug and metrics

---
 src/backend/executor/nodeHash.c | 35 +++++++++++++++++++++++++++++++++
 1 file changed, 35 insertions(+)

diff --git a/src/backend/executor/nodeHash.c b/src/backend/executor/nodeHash.c
index 4544296391..ec6b80121b 100644
--- a/src/backend/executor/nodeHash.c
+++ b/src/backend/executor/nodeHash.c
@@ -81,6 +81,35 @@ static bool ExecParallelHashTuplePrealloc(HashJoinTable hashtable,
 static void ExecParallelHashMergeCounters(HashJoinTable hashtable);
 static void ExecParallelHashCloseBatchAccessors(HashJoinTable hashtable);
 
+static void debugIncreaseBatches(HashJoinTable hashtable, char * label)
+{
+	int allocInnerBufFiles = 0;
+	int allocOuterBufFiles = 0;
+
+	for (int i=0; i<hashtable->nbatch; i++) {
+		if (hashtable->innerBatchFile != NULL &&
+			hashtable->innerBatchFile[i] != NULL)
+			allocInnerBufFiles++;
+		if (hashtable->outerBatchFile != NULL &&
+			hashtable->outerBatchFile[i] != NULL)
+			allocOuterBufFiles++;
+	}
+
+	elog(WARNING, "%25s: %d,\t%ld,\t%ld,\t%ld,\t%ld,\t%d,\t%d,\t%d,\t%d",
+		 label,
+		 hashtable->nbatch,
+		 hashtable->spaceAllowed,
+		 hashtable->spaceUsed,
+		 hashtable->spacePeak,
+		 hashtable->fileCxt->mem_allocated,
+		 allocInnerBufFiles,
+		 allocOuterBufFiles,
+		 hashtable->nbuckets,
+		 hashtable->growEnabled);
+
+	elog(LOG, "%s ======= context stats =======", label);
+	MemoryContextStats(TopMemoryContext);
+}
 
 /* ----------------------------------------------------------------
  *		ExecHash
@@ -889,6 +918,8 @@ ExecHashTableDestroy(HashJoinTable hashtable)
 {
 	int			i;
 
+	debugIncreaseBatches(hashtable, "destroying the hash table");
+
 	/*
 	 * Make sure all the temp files are closed.  We skip batch 0, since it
 	 * can't have any temp files (and the arrays might not even exist if
@@ -1685,7 +1716,11 @@ ExecHashTableInsert(HashJoinTable hashtable,
 		if (hashtable->spaceUsed +
 			hashtable->nbuckets_optimal * sizeof(HashJoinTuple)
 			> hashtable->spaceAllowed)
+		{
+			debugIncreaseBatches(hashtable, "trying to save memory");
 			ExecHashIncreaseNumBatches(hashtable);
+			debugIncreaseBatches(hashtable, "memory rescue done");
+		}
 	}
 	else
 	{
-- 
2.39.2

  [text/x-patch] 0004-v4-Limit-BufFile-memory-explosion-with-bad-HashJoin.patch (10.8K, ../20230408020119.32a0841b@karst/12-0004-v4-Limit-BufFile-memory-explosion-with-bad-HashJoin.patch)
  download | inline diff:
From 186cc5f60d8ffb11ca5150a3ece53e19368448d1 Mon Sep 17 00:00:00 2001
From: Jehan-Guillaume de Rorthais <jgdr@dalibo.com>
Date: Fri, 7 Apr 2023 19:24:20 +0200
Subject: [PATCH 4/4] Limit BufFile memory explosion with bad HashJoin

When hash join rely on a largely underestimated
statistics, splitting batches could result of a
huge memory consumption because BufFile buffers
associated to batches were not accounted.

This patch tries to keep the memory consummed
by batches and actual data balanced by rising
the memory limit to make more room for the data
before deciding to split again the batches.
---
 src/backend/executor/nodeHash.c     | 101 +++++++++++++++++++++++-----
 src/backend/executor/nodeHashjoin.c |  26 +++++--
 src/include/executor/nodeHashjoin.h |   2 +-
 3 files changed, 107 insertions(+), 22 deletions(-)

diff --git a/src/backend/executor/nodeHash.c b/src/backend/executor/nodeHash.c
index ec6b80121b..3a6e0bcfb7 100644
--- a/src/backend/executor/nodeHash.c
+++ b/src/backend/executor/nodeHash.c
@@ -81,6 +81,8 @@ static bool ExecParallelHashTuplePrealloc(HashJoinTable hashtable,
 static void ExecParallelHashMergeCounters(HashJoinTable hashtable);
 static void ExecParallelHashCloseBatchAccessors(HashJoinTable hashtable);
 
+static void ExecHashUpdateSpacePeak(HashJoinTable hashtable);
+
 static void debugIncreaseBatches(HashJoinTable hashtable, char * label)
 {
 	int allocInnerBufFiles = 0;
@@ -224,10 +226,8 @@ MultiExecPrivateHash(HashState *node)
 	if (hashtable->nbuckets != hashtable->nbuckets_optimal)
 		ExecHashIncreaseNumBuckets(hashtable);
 
-	/* Account for the buckets in spaceUsed (reported in EXPLAIN ANALYZE) */
-	hashtable->spaceUsed += hashtable->nbuckets * sizeof(HashJoinTuple);
-	if (hashtable->spaceUsed > hashtable->spacePeak)
-		hashtable->spacePeak = hashtable->spaceUsed;
+	/* refresh info about peak used memory */
+	ExecHashUpdateSpacePeak(hashtable);
 
 	hashtable->partialTuples = hashtable->totalTuples;
 }
@@ -969,6 +969,56 @@ ExecHashIncreaseNumBatches(HashJoinTable hashtable)
 	nbatch = oldnbatch * 2;
 	Assert(nbatch > 1);
 
+	/*
+	 * Each batch requires a non-trivial amount of memory, because BufFile
+	 * includes a PGAlignedBlock (typically 8kB buffer). So when doubling
+	 * the number of batches, we need to be careful and only allow that if
+	 * it actually has a chance of reducing memory usage.
+	 *
+	 * When doubling the number of batches, we expect to save roughly 1/2
+	 * of memory currently used for data (rows) at the price of doubling
+	 * the memory used for BufFile.
+	 * This becomes pointless when memory used for BufFile is greater than
+	 * the memory we expect to save:
+	 *
+	 *		(2 * nbatches * sizeof(BufFile)) > spaceUsed/2
+	 *
+	 * As far as this is confined inside work_mem, that's fine. But if
+	 * we significantly underestimate the number of batches, we may end
+	 * up in a situation where BufFile alone exceed work_mem.
+	 *
+	 * In such situation, move the threshold a bit, until the next point
+	 * where it'll make sense to consider adding batches again.
+	 * 
+	 * We can't stop adding batches entirely, because that would just mean
+	 * the batches would need more and more memory. So we need to increase
+	 * the number of batches, even if we can't enforce work_mem properly.
+	 * 
+	 * Note: This applies mostly to cases of significant underestimates,
+	 * resulting in an explosion of the number of batches. The properly
+	 * estimated cases should generally end up using merge join based on
+	 * high cost of the batched hash join.
+	 */
+	/*
+	 * fileCxt size is good enough estimation of BufFiles consumption.
+	 * Keep in mind spaceUsed includes real BufFile consumption as well
+	 */
+	if (2 * hashtable->fileCxt->mem_allocated > hashtable->spaceUsed / 2)
+	{
+		Size new_limit = nbatch * sizeof(PGAlignedBlock) * 3;
+		if (hashtable->spaceAllowed < new_limit)
+		{
+			ereport(WARNING, (
+				errmsg("Hash Join node must grow outside of work_mem"),
+				errdetail("Rising memory limit from %ld to %ld",
+						  hashtable->spaceAllowed, new_limit),
+				errhint("You might need to ANALYZE your table or tune its statistics collection.")));
+			hashtable->spaceAllowed = new_limit;
+		}
+
+		return;
+	}
+
 #ifdef HJDEBUG
 	printf("Hashjoin %p: increasing nbatch to %d because space = %zu\n",
 		   hashtable, nbatch, hashtable->spaceUsed);
@@ -1067,7 +1117,7 @@ ExecHashIncreaseNumBatches(HashJoinTable hashtable)
 				ExecHashJoinSaveTuple(HJTUPLE_MINTUPLE(hashTuple),
 									  hashTuple->hashvalue,
 									  &hashtable->innerBatchFile[batchno],
-									  hashtable->fileCxt);
+									  hashtable);
 
 				hashtable->spaceUsed -= hashTupleSize;
 				nfreed++;
@@ -1711,14 +1761,19 @@ ExecHashTableInsert(HashJoinTable hashtable,
 
 		/* Account for space used, and back off if we've used too much */
 		hashtable->spaceUsed += hashTupleSize;
-		if (hashtable->spaceUsed > hashtable->spacePeak)
-			hashtable->spacePeak = hashtable->spaceUsed;
+
+		/* refresh info about peak used memory */
+		ExecHashUpdateSpacePeak(hashtable);
+
+		/* Consider increasing number of batches. */
 		if (hashtable->spaceUsed +
 			hashtable->nbuckets_optimal * sizeof(HashJoinTuple)
 			> hashtable->spaceAllowed)
 		{
 			debugIncreaseBatches(hashtable, "trying to save memory");
+
 			ExecHashIncreaseNumBatches(hashtable);
+
 			debugIncreaseBatches(hashtable, "memory rescue done");
 		}
 	}
@@ -1732,7 +1787,7 @@ ExecHashTableInsert(HashJoinTable hashtable,
 		ExecHashJoinSaveTuple(tuple,
 							  hashvalue,
 							  &hashtable->innerBatchFile[batchno],
-							  hashtable->fileCxt);
+							  hashtable);
 	}
 
 	if (shouldFree)
@@ -1982,6 +2037,18 @@ ExecHashGetBucketAndBatch(HashJoinTable hashtable,
 	}
 }
 
+static void
+ExecHashUpdateSpacePeak(HashJoinTable hashtable)
+{
+	Size	spaceUsed = hashtable->spaceUsed;
+
+	/* Account for the buckets in spaceUsed (reported in EXPLAIN ANALYZE) */
+	spaceUsed += hashtable->nbuckets * sizeof(HashJoinTuple);
+
+	if (spaceUsed > hashtable->spacePeak)
+		hashtable->spacePeak = spaceUsed;
+}
+
 /*
  * ExecScanHashBucket
  *		scan a hash bucket for matches to the current outer tuple
@@ -2489,8 +2556,9 @@ ExecHashBuildSkewHash(HashJoinTable hashtable, Hash *node, int mcvsToUse)
 			+ mcvsToUse * sizeof(int);
 		hashtable->spaceUsedSkew += nbuckets * sizeof(HashSkewBucket *)
 			+ mcvsToUse * sizeof(int);
-		if (hashtable->spaceUsed > hashtable->spacePeak)
-			hashtable->spacePeak = hashtable->spaceUsed;
+
+		/* refresh info about peak used memory */
+		ExecHashUpdateSpacePeak(hashtable);
 
 		/*
 		 * Create a skew bucket for each MCV hash value.
@@ -2540,8 +2608,9 @@ ExecHashBuildSkewHash(HashJoinTable hashtable, Hash *node, int mcvsToUse)
 			hashtable->nSkewBuckets++;
 			hashtable->spaceUsed += SKEW_BUCKET_OVERHEAD;
 			hashtable->spaceUsedSkew += SKEW_BUCKET_OVERHEAD;
-			if (hashtable->spaceUsed > hashtable->spacePeak)
-				hashtable->spacePeak = hashtable->spaceUsed;
+
+			/* refresh info about peak used memory */
+			ExecHashUpdateSpacePeak(hashtable);
 		}
 
 		free_attstatsslot(&sslot);
@@ -2630,8 +2699,10 @@ ExecHashSkewTableInsert(HashJoinTable hashtable,
 	/* Account for space used, and back off if we've used too much */
 	hashtable->spaceUsed += hashTupleSize;
 	hashtable->spaceUsedSkew += hashTupleSize;
-	if (hashtable->spaceUsed > hashtable->spacePeak)
-		hashtable->spacePeak = hashtable->spaceUsed;
+
+	/* refresh info about peak used memory */
+	ExecHashUpdateSpacePeak(hashtable);
+
 	while (hashtable->spaceUsedSkew > hashtable->spaceAllowedSkew)
 		ExecHashRemoveNextSkewBucket(hashtable);
 
@@ -2715,7 +2786,7 @@ ExecHashRemoveNextSkewBucket(HashJoinTable hashtable)
 
 			ExecHashJoinSaveTuple(tuple, hashvalue,
 								  &hashtable->innerBatchFile[batchno],
-								  hashtable->fileCxt);
+								  hashtable);
 
 			pfree(hashTuple);
 			hashtable->spaceUsed -= tupleSize;
diff --git a/src/backend/executor/nodeHashjoin.c b/src/backend/executor/nodeHashjoin.c
index 5bc7f814c6..36fb024fab 100644
--- a/src/backend/executor/nodeHashjoin.c
+++ b/src/backend/executor/nodeHashjoin.c
@@ -488,7 +488,7 @@ ExecHashJoinImpl(PlanState *pstate, bool parallel)
 
 					ExecHashJoinSaveTuple(mintuple, hashvalue,
 										  &hashtable->outerBatchFile[batchno],
-										  hashtable->fileCxt);
+										  hashtable);
 
 					if (shouldFree)
 						heap_free_minimal_tuple(mintuple);
@@ -1041,8 +1041,11 @@ ExecHashJoinNewBatch(HashJoinState *hjstate)
 		 * away to free disk space.
 		 */
 		if (hashtable->outerBatchFile[curbatch])
+		{
 			BufFileClose(hashtable->outerBatchFile[curbatch]);
-		hashtable->outerBatchFile[curbatch] = NULL;
+			hashtable->outerBatchFile[curbatch] = NULL;
+			hashtable->spaceUsed -= sizeof(PGAlignedBlock);
+		}
 	}
 	else						/* we just finished the first batch */
 	{
@@ -1096,11 +1099,19 @@ ExecHashJoinNewBatch(HashJoinState *hjstate)
 		/* We can ignore this batch. */
 		/* Release associated temp files right away. */
 		if (hashtable->innerBatchFile[curbatch])
+		{
 			BufFileClose(hashtable->innerBatchFile[curbatch]);
-		hashtable->innerBatchFile[curbatch] = NULL;
+			hashtable->innerBatchFile[curbatch] = NULL;
+			hashtable->spaceUsed -= sizeof(PGAlignedBlock);
+		}
+
 		if (hashtable->outerBatchFile[curbatch])
+		{
 			BufFileClose(hashtable->outerBatchFile[curbatch]);
-		hashtable->outerBatchFile[curbatch] = NULL;
+			hashtable->outerBatchFile[curbatch] = NULL;
+			hashtable->spaceUsed -= sizeof(PGAlignedBlock);
+		}
+
 		curbatch++;
 	}
 
@@ -1141,6 +1152,7 @@ ExecHashJoinNewBatch(HashJoinState *hjstate)
 		 */
 		BufFileClose(innerFile);
 		hashtable->innerBatchFile[curbatch] = NULL;
+		hashtable->spaceUsed -= sizeof(PGAlignedBlock);
 	}
 
 	/*
@@ -1305,7 +1317,7 @@ ExecParallelHashJoinNewBatch(HashJoinState *hjstate)
  */
 void
 ExecHashJoinSaveTuple(MinimalTuple tuple, uint32 hashvalue,
-					  BufFile **fileptr, MemoryContext filecxt)
+					  BufFile **fileptr, HashJoinTable hashtable)
 {
 	BufFile    *file = *fileptr;
 
@@ -1313,13 +1325,15 @@ ExecHashJoinSaveTuple(MinimalTuple tuple, uint32 hashvalue,
 	{
 		MemoryContext oldctx;
 
-		oldctx = MemoryContextSwitchTo(filecxt);
+		oldctx = MemoryContextSwitchTo(hashtable->fileCxt);
 
 		/* First write to this batch file, so open it. */
 		file = BufFileCreateTemp(false);
 		*fileptr = file;
 
 		MemoryContextSwitchTo(oldctx);
+
+		hashtable->spaceUsed += sizeof(PGAlignedBlock);
 	}
 
 	BufFileWrite(file, &hashvalue, sizeof(uint32));
diff --git a/src/include/executor/nodeHashjoin.h b/src/include/executor/nodeHashjoin.h
index a8f9ae1989..ccb704ede1 100644
--- a/src/include/executor/nodeHashjoin.h
+++ b/src/include/executor/nodeHashjoin.h
@@ -29,6 +29,6 @@ extern void ExecHashJoinInitializeWorker(HashJoinState *state,
 										 ParallelWorkerContext *pwcxt);
 
 extern void ExecHashJoinSaveTuple(MinimalTuple tuple, uint32 hashvalue,
-								  BufFile **fileptr, MemoryContext filecxt);
+								  BufFile **fileptr, HashJoinTable hashtable);
 
 #endif							/* NODEHASHJOIN_H */
-- 
2.39.2

view thread (60+ messages)  latest in thread

Message-ID: <20230408020119.32a0841b@karst>
Permalink:  ../20230408020119.32a0841b@karst/
Also on:    postgresql.org/message-id/20230408020119.32a0841b@karst

reply

Reply instructions:

You may reply publicly to this message via plain-text email
using any one of the following methods:

* Reply to all the recipients using the --to and --cc options:
  reply via email

  To: pgsql-hackers@postgresql.org
  Cc: jgdr@dalibo.com, tomas.vondra@enterprisedb.com, melanieplageman@gmail.com, pgsql-hackers@lists.postgresql.org
  Subject: Re: Memory leak from ExecutorState context?
  In-Reply-To: <20230408020119.32a0841b@karst>

* Save the following mbox file, import it into your mail client,
  and reply-to-all from there: mbox

This inbox is served by agora; see mirroring instructions
for how to clone and mirror all data and code used for this inbox