Skip to content

Commit ce7aa63

Browse files
authored
[VL]Add hash table memory usage metric (#12528)
1 parent 352e05a commit ce7aa63

8 files changed

Lines changed: 47 additions & 3 deletions

File tree

backends-velox/src/main/java/org/apache/gluten/vectorized/HashJoinBuilder.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -48,6 +48,8 @@ public static native long deserializeHashTableDirect(
4848

4949
public static native long getHashTableBloomFilterBlocksByteSize(long hashTableHandle);
5050

51+
public static native long getHashTableMemoryUsage(long hashTableHandle);
52+
5153
public static native long serializedHashTableSizeDirect(long hashTableHandle);
5254

5355
public static native void serializeHashTableDirect(long hashTableHandle, long address, long size);

backends-velox/src/main/scala/org/apache/gluten/backendsapi/velox/VeloxMetricsApi.scala

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -727,7 +727,10 @@ class VeloxMetricsApi extends MetricsApi with Logging {
727727
"time to build hash table"),
728728
"deserializeHashTableTime" -> SQLMetrics.createTimingMetric(
729729
sparkContext,
730-
"time to deserialize hash table")
730+
"time to deserialize hash table"),
731+
"hashTableMemorySize" -> SQLMetrics.createSizeMetric(
732+
sparkContext,
733+
"hash table memory size")
731734
)
732735

733736
override def genHashJoinTransformerMetricsUpdater(

backends-velox/src/main/scala/org/apache/gluten/execution/HashJoinExecTransformer.scala

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -181,7 +181,8 @@ case class BroadcastHashJoinExecTransformer(
181181
metrics.get("buildHashTableTime"),
182182
metrics.get("serializeHashTableTime"),
183183
metrics.get("deserializeHashTableTime"),
184-
metrics.get("serializedHashTableSize")
184+
metrics.get("serializedHashTableSize"),
185+
metrics.get("hashTableMemorySize")
185186
)
186187

187188
// Check the type of broadcast relation to determine the approach
@@ -265,7 +266,8 @@ case class BroadcastHashJoinContext(
265266
buildHashTableTimeMetric: Option[SQLMetric] = None,
266267
serializeHashTableTimeMetric: Option[SQLMetric] = None,
267268
deserializeHashTableTimeMetric: Option[SQLMetric] = None,
268-
serializedHashTableSizeMetric: Option[SQLMetric] = None) {
269+
serializedHashTableSizeMetric: Option[SQLMetric] = None,
270+
hashTableMemorySizeMetric: Option[SQLMetric] = None) {
269271
def droppedDuplicates: Boolean = {
270272
!hasMixedFiltCondition && (
271273
substraitJoinType == JoinRel.JoinType.JOIN_TYPE_LEFT_SEMI ||

backends-velox/src/main/scala/org/apache/gluten/execution/VeloxBroadcastBuildSideCache.scala

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -98,6 +98,9 @@ object VeloxBroadcastBuildSideCache
9898
unsafe.buildHashTable(broadcastContext)
9999
}
100100

101+
broadcastContext.hashTableMemorySizeMetric.foreach(
102+
_ += HashJoinBuilder.getHashTableMemoryUsage(pointer))
103+
101104
BroadcastHashTable(pointer, relation, droppedDuplicates)
102105
}
103106
)

cpp/velox/jni/JniHashTable.cc

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -180,6 +180,11 @@ size_t serializedHashTableSize(std::shared_ptr<HashTableBuilder> builder) {
180180
return HashTableSerializer::serializedSize<true>(hashTableTrue);
181181
}
182182

183+
int64_t hashTableMemoryUsage(std::shared_ptr<HashTableBuilder> builder) {
184+
VELOX_CHECK_NOT_NULL(builder, "Hash table builder cannot be null");
185+
return builder->hashTableMemoryUsage();
186+
}
187+
183188
void serializeHashTableTo(std::shared_ptr<HashTableBuilder> builder, uint8_t* data, size_t size) {
184189
VELOX_CHECK_NOT_NULL(builder, "Hash table builder cannot be null");
185190
VELOX_CHECK_NOT_NULL(data, "Serialized buffer cannot be null");

cpp/velox/jni/JniHashTable.h

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -94,6 +94,9 @@ long getJoin(const std::string& hashTableId);
9494
// Return the exact serialized hash table size for direct buffer allocation.
9595
size_t serializedHashTableSize(std::shared_ptr<HashTableBuilder> builder);
9696

97+
// Return retained bytes for the hash table, including parallel-build subtables.
98+
int64_t hashTableMemoryUsage(std::shared_ptr<HashTableBuilder> builder);
99+
97100
// Serialize hash table directly to a caller-provided buffer.
98101
void serializeHashTableTo(std::shared_ptr<HashTableBuilder> builder, uint8_t* data, size_t size);
99102

cpp/velox/jni/VeloxJniWrapper.cc

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1054,6 +1054,7 @@ JNIEXPORT jlong JNICALL Java_org_apache_gluten_vectorized_HashJoinBuilder_native
10541054
builder->dropDuplicates(),
10551055
nullptr);
10561056
builder->setHashTable(std::move(mainTable));
1057+
builder->setHashTableMemoryUsage(builder->hashTable()->allocatedBytes());
10571058

10581059
auto* cache = facebook::velox::exec::HashTableCache::instance();
10591060

@@ -1122,6 +1123,10 @@ JNIEXPORT jlong JNICALL Java_org_apache_gluten_vectorized_HashJoinBuilder_native
11221123
for (int i = 1; i < numThreads; ++i) {
11231124
tables.push_back(std::move(otherTables[i]));
11241125
}
1126+
int64_t otherTablesMemoryUsage = 0;
1127+
for (const auto& table : tables) {
1128+
otherTablesMemoryUsage += table->allocatedBytes();
1129+
}
11251130

11261131
// TODO: Get accurate signal if parallel join build is going to be applied
11271132
// from hash table. Currently there is still a chance inside hash table that
@@ -1143,6 +1148,8 @@ JNIEXPORT jlong JNICALL Java_org_apache_gluten_vectorized_HashJoinBuilder_native
11431148
}
11441149

11451150
hashTableBuilders[0]->setHashTable(std::move(mainTable));
1151+
hashTableBuilders[0]->setHashTableMemoryUsage(
1152+
hashTableBuilders[0]->hashTable()->allocatedBytes() + otherTablesMemoryUsage);
11461153

11471154
auto* cache = facebook::velox::exec::HashTableCache::instance();
11481155
if (!cache->exist(hashTableId)) {
@@ -1260,6 +1267,16 @@ Java_org_apache_gluten_vectorized_HashJoinBuilder_getHashTableBloomFilterBlocksB
12601267
JNI_METHOD_END(0L)
12611268
}
12621269

1270+
JNIEXPORT jlong JNICALL Java_org_apache_gluten_vectorized_HashJoinBuilder_getHashTableMemoryUsage( // NOLINT
1271+
JNIEnv* env,
1272+
jclass,
1273+
jlong hashTableHandle) {
1274+
JNI_METHOD_START
1275+
auto builder = ObjectStore::retrieve<gluten::HashTableBuilder>(hashTableHandle);
1276+
return static_cast<jlong>(gluten::hashTableMemoryUsage(builder));
1277+
JNI_METHOD_END(0L)
1278+
}
1279+
12631280
JNIEXPORT jlong JNICALL Java_org_apache_gluten_vectorized_HashJoinBuilder_serializedHashTableSizeDirect( // NOLINT
12641281
JNIEnv* env,
12651282
jclass,

cpp/velox/operators/hashjoin/HashTableBuilder.h

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -75,6 +75,14 @@ class HashTableBuilder {
7575
return noMoreInput_;
7676
}
7777

78+
void setHashTableMemoryUsage(int64_t hashTableMemoryUsage) {
79+
hashTableMemoryUsage_ = hashTableMemoryUsage;
80+
}
81+
82+
int64_t hashTableMemoryUsage() const {
83+
return hashTableMemoryUsage_;
84+
}
85+
7886
uint32_t joinBuildVectorHasherMaxNumDistinct() const {
7987
return joinBuildVectorHasherMaxNumDistinct_;
8088
}
@@ -150,6 +158,7 @@ class HashTableBuilder {
150158
uint32_t abandonHashBuildDedupMinRows_{100'000};
151159
uint32_t abandonHashBuildDedupMinPct_{0};
152160
bool filterPropagatesNulls_{false};
161+
int64_t hashTableMemoryUsage_{0};
153162
};
154163

155164
} // namespace gluten

0 commit comments

Comments
 (0)