github-actions[bot] commented on code in PR #68498:
URL: https://github.com/apache/doris/pull/68498#discussion_r4137227788
##########
be/src/core/value/quantile_state.cpp:
##########
@@ -28,7 +31,81 @@
#include "util/tdigest.h"
#include "util/unaligned.h"
+#ifdef BE_TEST
+#include "cpp/sync_point.h"
+#endif
+
namespace doris {
+
+// Shares a digest across QuantileState copies and detaches it before sample
+// writes. Readers use shared locks; compression uses an exclusive lock.
+struct QuantileState::TDigestHolder {
+ explicit TDigestHolder(float compression) : digest(compression) {}
+ TDigestHolder(const TDigestHolder& other) : digest(other.digest) {}
+
+ std::shared_lock<std::shared_mutex> lock_processed_digest() {
+ std::shared_lock read_lock(mutex);
+ if (digest.have_unprocessed()) {
+ read_lock.unlock();
+ {
+ std::unique_lock write_lock(mutex);
+ if (digest.have_unprocessed()) {
+ digest.compress();
+ }
+ }
+ read_lock.lock();
+ }
+ return read_lock;
+ }
+
+ TDigest digest;
+ std::shared_mutex mutex;
+};
+
+QuantileState QuantileState::copy_for_result() const {
+ if (_type != TDIGEST) {
+ return *this;
+ }
+ QuantileState result(_compression);
+ result._type = TDIGEST;
+ {
+ // Reuse processed centroids on the next row instead of sorting the
prefix again.
Review Comment:
[P3] Keep the analytic Blocks charge stable when processing a shared digest.
A grouped `quantile_union` result can still have unprocessed centroids when it
feeds `quantile_union(q) OVER (...)`. The window accumulator shares its holder
with the retained input `q` column, and this call grows processed/cumulative
vector capacity after the sink has charged the input Block's
`allocated_bytes()`. `_output_current_block()` later subtracts the Block's
now-larger allocation, so its `Blocks` high-water counter can become negative
and understate subsequent peaks. Store the charged byte count for removal, or
adjust the counter when the holder's capacity changes.
##########
be/src/core/column/column_complex.h:
##########
@@ -109,7 +109,12 @@ class ColumnComplexType final : public COWHelper<IColumn,
ColumnComplexType<T>>
// calculate the memory requested by value_type
size_t byte_size() const override { return data.size() * sizeof(data[0]); }
- size_t allocated_bytes() const override { return byte_size(); }
+ size_t allocated_bytes() const override {
+ if constexpr (T == TYPE_QUANTILE_STATE) {
+ return value_type::allocated_bytes(data);
Review Comment:
[P2] Keep the join build reservation estimate on comparable byte bases.
`HashJoinBuildSinkLocalState::get_reserve_mem_size()` still estimates the next
batch from `byte_size()` (only the outer `QuantileState` rows), but this new
`allocated_bytes()` includes digest buffers in the denominator of its 85%
growth test. For 4096 rows sharing a high-compression digest, roughly 256 KiB
of outer rows plus at least 625 KiB of digest buffer makes the next 4096-row
batch score below 85%, so the sink skips reservation even when appending
reallocates the row vector. Under workload-group memory pressure this bypasses
the normal reserve/pause path and can fail the query during that allocation.
Use comparable outer-column capacity for the growth test, or separately
estimate incoming digest memory.
##########
be/src/core/value/quantile_state.cpp:
##########
@@ -28,7 +31,81 @@
#include "util/tdigest.h"
#include "util/unaligned.h"
+#ifdef BE_TEST
+#include "cpp/sync_point.h"
+#endif
+
namespace doris {
+
+// Shares a digest across QuantileState copies and detaches it before sample
+// writes. Readers use shared locks; compression uses an exclusive lock.
+struct QuantileState::TDigestHolder {
+ explicit TDigestHolder(float compression) : digest(compression) {}
+ TDigestHolder(const TDigestHolder& other) : digest(other.digest) {}
+
+ std::shared_lock<std::shared_mutex> lock_processed_digest() {
+ std::shared_lock read_lock(mutex);
+ if (digest.have_unprocessed()) {
+ read_lock.unlock();
+ {
+ std::unique_lock write_lock(mutex);
+ if (digest.have_unprocessed()) {
+ digest.compress();
+ }
+ }
+ read_lock.lock();
+ }
+ return read_lock;
+ }
+
+ TDigest digest;
+ std::shared_mutex mutex;
+};
+
+QuantileState QuantileState::copy_for_result() const {
+ if (_type != TDIGEST) {
+ return *this;
+ }
+ QuantileState result(_compression);
+ result._type = TDIGEST;
+ {
+ // Reuse processed centroids on the next row instead of sorting the
prefix again.
+ auto lock = _tdigest_ptr->lock_processed_digest();
+ result._tdigest_ptr = std::make_shared<TDigestHolder>(*_tdigest_ptr);
+ }
+ return result;
+}
+
+size_t QuantileState::allocated_bytes(const std::vector<QuantileState>&
states) {
+ size_t bytes = states.capacity() * sizeof(QuantileState);
+ std::unordered_set<const TDigestHolder*> counted;
Review Comment:
[P2] Avoid scanning the entire quantile column on every allocation check.
`QuantileState::allocated_bytes(data)` now rebuilds a hash set and locks each
distinct digest while walking every row. The hash-join build sink retains one
growing mutable block and calls `allocated_bytes()` after each incoming batch
(and during reservation). With one million quantile rows arriving in 4096-row
batches, those post-merge checks alone revisit about 123 million states; the
previous implementation was constant time per check. Please maintain an
incrementally updated byte count or use a bounded-cost accounting path so
growing joins do not spend quadratic time measuring memory.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]