This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new 3d78b6b31f0 branch-4.1: [fix](be) Prevent shared quantile state
mutation with TDigest COW (#68498) (#68759)
3d78b6b31f0 is described below
commit 3d78b6b31f05793ee39ae67c2de630451cdc5848
Author: linrrarity <[email protected]>
AuthorDate: Fri Oct 9 09:31:06 2026 +0800
branch-4.1: [fix](be) Prevent shared quantile state mutation with TDigest
COW (#68498) (#68759)
pick: https://github.com/apache/doris/pull/68498
---
be/src/core/value/quantile_state.cpp | 122 +++++++--
be/src/core/value/quantile_state.h | 16 +-
.../aggregate_function_quantile_state.cpp | 4 +-
.../aggregate/aggregate_function_quantile_state.h | 24 +-
be/src/util/tdigest.h | 2 +
be/test/core/value/quantile_state_test.cpp | 294 +++++++++++++++++++++
be/test/exprs/aggregate/agg_percentile_test.cpp | 171 ++++++++++++
be/test/util/tdigest_test.cpp | 28 ++
.../test_quantile_state_function.out | 32 +++
.../test_quantile_state_function.groovy | 90 +++++++
10 files changed, 754 insertions(+), 29 deletions(-)
diff --git a/be/src/core/value/quantile_state.cpp
b/be/src/core/value/quantile_state.cpp
index 95598c2b38a..9ca0f0a0ca8 100644
--- a/be/src/core/value/quantile_state.cpp
+++ b/be/src/core/value/quantile_state.cpp
@@ -19,7 +19,9 @@
#include <string.h>
#include <cmath>
+#include <mutex>
#include <ostream>
+#include <shared_mutex>
#include <utility>
#include "common/logging.h"
@@ -28,8 +30,69 @@
#include "util/tdigest.h"
#include "util/unaligned.h"
+#ifdef BE_TEST
+#include "cpp/sync_point.h"
+#endif
+
namespace doris {
#include "common/compile_check_begin.h"
+
+// 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.haveUnprocessed()) {
+ read_lock.unlock();
+ {
+ std::unique_lock write_lock(mutex);
+ if (digest.haveUnprocessed()) {
+ 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();
+ // Vector copies retain the elements, not the accumulator's spare
write capacity.
+ result._tdigest_ptr = std::make_shared<TDigestHolder>(*_tdigest_ptr);
+ }
+ return result;
+}
+
+TDigest& QuantileState::_mutable_tdigest() {
+ if (_tdigest_ptr.use_count() == 1) {
+ return _tdigest_ptr->digest;
+ }
+ std::shared_ptr<TDigestHolder> detached;
+ {
+ std::shared_lock lock(_tdigest_ptr->mutex);
+#ifdef BE_TEST
+ TEST_SYNC_POINT("QuantileState::detach:source_locked");
+#endif
+ detached = std::make_shared<TDigestHolder>(*_tdigest_ptr);
+ }
+ _tdigest_ptr = std::move(detached);
+ return _tdigest_ptr->digest;
+}
+
QuantileState::QuantileState() : _type(EMPTY),
_compression(QUANTILE_STATE_COMPRESSION_MIN) {}
QuantileState::QuantileState(float compression) : _type(EMPTY),
_compression(compression) {}
@@ -51,10 +114,14 @@ size_t QuantileState::get_serialized_size() const {
case EXPLICIT:
size += sizeof(uint16_t) + sizeof(double) * _explicit_data.size();
break;
- case TDIGEST:
- size += _tdigest_ptr->serialized_size();
+ case TDIGEST: {
+ // Compress before sizing so concurrent queries cannot change the size
+ // before serialize(); writes through shared copies detach via COW.
+ auto lock = _tdigest_ptr->lock_processed_digest();
+ size += _tdigest_ptr->digest.serialized_size();
break;
}
+ }
return size;
}
@@ -138,7 +205,8 @@ double QuantileState::get_value_by_percentile(float
percentile) const {
return get_explicit_value_by_percentile(percentile);
}
case TDIGEST: {
- return _tdigest_ptr->quantile(percentile);
+ auto lock = _tdigest_ptr->lock_processed_digest();
+ return _tdigest_ptr->digest.quantileProcessed(percentile);
}
default:
break;
@@ -187,8 +255,8 @@ bool QuantileState::deserialize(const Slice& slice) {
}
case TDIGEST: {
// 4: Tdigest object value
- _tdigest_ptr = std::make_shared<TDigest>(0);
- _tdigest_ptr->unserialize(ptr);
+ _tdigest_ptr = std::make_shared<TDigestHolder>(0);
+ _tdigest_ptr->digest.unserialize(ptr);
break;
}
default:
@@ -225,7 +293,8 @@ size_t QuantileState::serialize(uint8_t* dst) const {
}
case TDIGEST: {
*ptr++ = TDIGEST;
- size_t tdigest_size = _tdigest_ptr->serialize(ptr);
+ auto lock = _tdigest_ptr->lock_processed_digest();
+ size_t tdigest_size = _tdigest_ptr->digest.serialize(ptr);
ptr += tdigest_size;
break;
}
@@ -236,6 +305,11 @@ size_t QuantileState::serialize(uint8_t* dst) const {
}
void QuantileState::merge(const QuantileState& other) {
+ if (this == &other) {
+ const QuantileState source(other);
+ merge(source);
+ return;
+ }
switch (other._type) {
case EMPTY:
break;
@@ -257,23 +331,25 @@ void QuantileState::merge(const QuantileState& other) {
case EXPLICIT:
if (_explicit_data.size() + other._explicit_data.size() >
QUANTILE_STATE_EXPLICIT_NUM) {
_type = TDIGEST;
- _tdigest_ptr = std::make_shared<TDigest>(_compression);
+ _tdigest_ptr = std::make_shared<TDigestHolder>(_compression);
for (int i = 0; i < _explicit_data.size(); i++) {
- _tdigest_ptr->add((float)_explicit_data[i]);
+ _tdigest_ptr->digest.add((float)_explicit_data[i]);
}
for (int i = 0; i < other._explicit_data.size(); i++) {
- _tdigest_ptr->add((float)other._explicit_data[i]);
+ _tdigest_ptr->digest.add((float)other._explicit_data[i]);
}
} else {
_explicit_data.insert(_explicit_data.end(),
other._explicit_data.begin(),
other._explicit_data.end());
}
break;
- case TDIGEST:
+ case TDIGEST: {
+ auto& digest = _mutable_tdigest();
for (int i = 0; i < other._explicit_data.size(); i++) {
- _tdigest_ptr->add((float)other._explicit_data[i]);
+ digest.add((float)other._explicit_data[i]);
}
break;
+ }
default:
break;
}
@@ -288,18 +364,26 @@ void QuantileState::merge(const QuantileState& other) {
case SINGLE:
_type = TDIGEST;
_tdigest_ptr = other._tdigest_ptr;
- _tdigest_ptr->add((float)_single_data);
+ _mutable_tdigest().add((float)_single_data);
break;
- case EXPLICIT:
+ case EXPLICIT: {
_type = TDIGEST;
_tdigest_ptr = other._tdigest_ptr;
+ auto& digest = _mutable_tdigest();
for (int i = 0; i < _explicit_data.size(); i++) {
- _tdigest_ptr->add((float)_explicit_data[i]);
+ digest.add((float)_explicit_data[i]);
}
break;
- case TDIGEST:
- _tdigest_ptr->merge(other._tdigest_ptr.get());
+ }
+ case TDIGEST: {
+ auto& digest = _mutable_tdigest();
+ std::shared_lock lock(other._tdigest_ptr->mutex);
+#ifdef BE_TEST
+ TEST_SYNC_POINT("QuantileState::merge:source_locked");
+#endif
+ digest.merge(&other._tdigest_ptr->digest);
break;
+ }
default:
break;
}
@@ -323,9 +407,9 @@ void QuantileState::add_value(const double& value) {
break;
case EXPLICIT:
if (_explicit_data.size() == QUANTILE_STATE_EXPLICIT_NUM) {
- _tdigest_ptr = std::make_shared<TDigest>(_compression);
+ _tdigest_ptr = std::make_shared<TDigestHolder>(_compression);
for (int i = 0; i < _explicit_data.size(); i++) {
- _tdigest_ptr->add((float)_explicit_data[i]);
+ _tdigest_ptr->digest.add((float)_explicit_data[i]);
}
_explicit_data.clear();
_explicit_data.shrink_to_fit();
@@ -336,7 +420,7 @@ void QuantileState::add_value(const double& value) {
}
break;
case TDIGEST:
- _tdigest_ptr->add((float)value);
+ _mutable_tdigest().add((float)value);
break;
}
}
diff --git a/be/src/core/value/quantile_state.h
b/be/src/core/value/quantile_state.h
index 50868290976..fc540848c26 100644
--- a/be/src/core/value/quantile_state.h
+++ b/be/src/core/value/quantile_state.h
@@ -47,11 +47,14 @@ public:
QuantileState();
explicit QuantileState(float compression);
explicit QuantileState(const Slice& slice);
- QuantileState& operator=(const QuantileState& other) noexcept = default;
- QuantileState(const QuantileState& other) noexcept = default;
+ QuantileState& operator=(const QuantileState& other) = default;
+ QuantileState(const QuantileState& other) = default;
QuantileState& operator=(QuantileState&& other) noexcept = default;
QuantileState(QuantileState&& other) noexcept = default;
+ // A compact, independent value for retained window results.
+ QuantileState copy_for_result() const;
+
void set_compression(float compression);
bool deserialize(const Slice& slice);
size_t serialize(uint8_t* dst) const;
@@ -70,9 +73,14 @@ public:
~QuantileState() = default;
private:
+ // Copies share a digest until the first sample write. Concurrent const
+ // operations are supported; mutating the same state requires exclusive
access.
+ struct TDigestHolder;
+ TDigest& _mutable_tdigest();
+
QuantileStateType _type = EMPTY;
- std::shared_ptr<TDigest> _tdigest_ptr;
- double _single_data;
+ std::shared_ptr<TDigestHolder> _tdigest_ptr;
+ double _single_data = 0;
std::vector<double> _explicit_data;
float _compression;
};
diff --git a/be/src/exprs/aggregate/aggregate_function_quantile_state.cpp
b/be/src/exprs/aggregate/aggregate_function_quantile_state.cpp
index 8e59828db3d..f31cb0d0cad 100644
--- a/be/src/exprs/aggregate/aggregate_function_quantile_state.cpp
+++ b/be/src/exprs/aggregate/aggregate_function_quantile_state.cpp
@@ -32,11 +32,11 @@ AggregateFunctionPtr
create_aggregate_function_quantile_state_union(
if (arg_is_nullable) {
return std::make_shared<
AggregateFunctionQuantileStateOp<true,
AggregateFunctionQuantileStateUnionOp>>(
- argument_types);
+ argument_types, attr.is_window_function);
} else {
return std::make_shared<
AggregateFunctionQuantileStateOp<false,
AggregateFunctionQuantileStateUnionOp>>(
- argument_types);
+ argument_types, attr.is_window_function);
}
}
diff --git a/be/src/exprs/aggregate/aggregate_function_quantile_state.h
b/be/src/exprs/aggregate/aggregate_function_quantile_state.h
index fb03d3eef8f..e85bd4ec40b 100644
--- a/be/src/exprs/aggregate/aggregate_function_quantile_state.h
+++ b/be/src/exprs/aggregate/aggregate_function_quantile_state.h
@@ -103,10 +103,11 @@ public:
String get_name() const override { return Op::name; }
- AggregateFunctionQuantileStateOp(const DataTypes& argument_types_)
+ AggregateFunctionQuantileStateOp(const DataTypes& argument_types_, bool
is_window_function)
:
IAggregateFunctionDataHelper<AggregateFunctionQuantileStateData<Op>,
AggregateFunctionQuantileStateOp<arg_is_nullable, Op>>(
- argument_types_) {}
+ argument_types_),
+ _is_window_function(is_window_function) {}
DataTypePtr get_return_type() const override {
return std::make_shared<DataTypeQuantileState>();
@@ -144,11 +145,26 @@ public:
}
void insert_result_into(ConstAggregateDataPtr __restrict place, IColumn&
to) const override {
- auto& column = assert_cast<ColVecResult&>(to);
- column.get_data().push_back(this->data(place).get());
+ auto& column = assert_cast<ColVecResult&,
TypeCheckOnRelease::DISABLE>(to);
+ const auto& value = this->data(place).get();
+ column.get_data().push_back(_is_window_function ?
value.copy_for_result() : value);
+ }
+
+ void insert_result_into_range(ConstAggregateDataPtr __restrict place,
IColumn& to,
+ const size_t start, const size_t end) const
override {
+ if (start == end) {
+ return;
+ }
+ insert_result_into(place, to);
+ auto& data = assert_cast<ColVecResult&,
TypeCheckOnRelease::DISABLE>(to).get_data();
+ // Rows with the same window result share one compact digest.
+ data.insert(data.end(), end - start - 1, data.back());
}
void reset(AggregateDataPtr __restrict place) const override {
this->data(place).reset(); }
+
+private:
+ const bool _is_window_function;
};
AggregateFunctionPtr create_aggregate_function_quantile_state_union(
diff --git a/be/src/util/tdigest.h b/be/src/util/tdigest.h
index 5dad6d5bac9..158dc748ce9 100644
--- a/be/src/util/tdigest.h
+++ b/be/src/util/tdigest.h
@@ -171,6 +171,8 @@ public:
return w;
}
+ TDigest(const TDigest&) = default;
+
TDigest& operator=(TDigest&& o) {
_compression = o._compression;
_max_processed = o._max_processed;
diff --git a/be/test/core/value/quantile_state_test.cpp
b/be/test/core/value/quantile_state_test.cpp
index ae3fc43b1ed..c025294caa3 100644
--- a/be/test/core/value/quantile_state_test.cpp
+++ b/be/test/core/value/quantile_state_test.cpp
@@ -20,10 +20,304 @@
#include <gtest/gtest-message.h>
#include <gtest/gtest-test-part.h>
+#include <atomic>
+#include <barrier>
+#include <chrono>
+#include <cstring>
+#include <future>
+#include <thread>
+
+#include "cpp/sync_point.h"
#include "gtest/gtest_pred_impl.h"
+#include "util/tdigest.h"
namespace doris {
+TEST(QuantileStateTest, SharedDigestWritesDoNotChangeSource) {
+ QuantileState source;
+ for (int i = 0; i < 4096; ++i) {
+ source.add_value(10);
+ }
+ QuantileState copied = source;
+ EXPECT_EQ(copied._tdigest_ptr, source._tdigest_ptr);
+ copied.add_value(110);
+ EXPECT_NE(copied._tdigest_ptr, source._tdigest_ptr);
+ EXPECT_EQ(110, copied.get_value_by_percentile(1));
+ EXPECT_EQ(10, source.get_value_by_percentile(1));
+
+ QuantileState merged;
+ merged.merge(source);
+ merged.add_value(-10);
+ EXPECT_EQ(-10, merged.get_value_by_percentile(0));
+ EXPECT_EQ(10, source.get_value_by_percentile(0));
+}
+
+static QuantileState constant_state(double value, int count = 4096) {
+ QuantileState state;
+ for (int i = 0; i < count; ++i) {
+ state.add_value(value);
+ }
+ return state;
+}
+
+TEST(QuantileStateTest, MergeDoesNotModifyDigestSourceForAnyTargetType) {
+ for (int count : {0, 1, 2, 4096}) {
+ auto source = constant_state(10);
+ auto target = constant_state(-10, count);
+ target.merge(source);
+ EXPECT_EQ(count == 0 ? 10 : -10, target.get_value_by_percentile(0));
+ EXPECT_EQ(10, target.get_value_by_percentile(1));
+ target.add_value(110);
+ EXPECT_EQ(110, target.get_value_by_percentile(1));
+ EXPECT_EQ(10, source.get_value_by_percentile(0));
+ EXPECT_EQ(10, source.get_value_by_percentile(1));
+ }
+}
+
+TEST(QuantileStateTest, AssignmentAndExplicitMergeDetachOnlyOnce) {
+ auto source = constant_state(10);
+ QuantileState target;
+ target = source;
+ auto extra = constant_state(-10, 2);
+ target.merge(extra);
+ auto* detached = target._tdigest_ptr.get();
+ EXPECT_NE(detached, source._tdigest_ptr.get());
+ for (int i = 0; i < 100; ++i) {
+ target.add_value(110);
+ EXPECT_EQ(detached, target._tdigest_ptr.get());
+ }
+ EXPECT_EQ(-10, target.get_value_by_percentile(0));
+ EXPECT_EQ(110, target.get_value_by_percentile(1));
+ EXPECT_EQ(10, source.get_value_by_percentile(0));
+ EXPECT_EQ(10, source.get_value_by_percentile(1));
+}
+
+TEST(QuantileStateTest, WindowReadsReuseDigestUntilResultIsCopied) {
+ auto state = constant_state(10);
+ auto* original = state._tdigest_ptr.get();
+ for (int i = 11; i < 100; ++i) {
+ state.add_value(i);
+ EXPECT_EQ(i, state.get_value_by_percentile(1));
+ EXPECT_EQ(original, state._tdigest_ptr.get());
+ }
+ auto saved_result = state;
+ state.add_value(110);
+ EXPECT_EQ(99, saved_result.get_value_by_percentile(1));
+ EXPECT_EQ(110, state.get_value_by_percentile(1));
+}
+
+TEST(QuantileStateTest, WritesReuseDigestAfterLastCopyIsDestroyed) {
+ auto state = constant_state(10);
+ auto* original = state._tdigest_ptr.get();
+ {
+ auto copy = state;
+ EXPECT_EQ(original, copy._tdigest_ptr.get());
+ }
+ state.add_value(110);
+ EXPECT_EQ(original, state._tdigest_ptr.get());
+ EXPECT_EQ(110, state.get_value_by_percentile(1));
+}
+
+TEST(QuantileStateTest, WritesReuseDigestAfterOtherCopyDetaches) {
+ auto state = constant_state(10);
+ auto* original = state._tdigest_ptr.get();
+ auto copy = state;
+ copy.add_value(110);
+ state.add_value(-10);
+ EXPECT_EQ(original, state._tdigest_ptr.get());
+ EXPECT_EQ(10, state.get_value_by_percentile(1));
+ EXPECT_EQ(110, copy.get_value_by_percentile(1));
+ EXPECT_EQ(10, copy.get_value_by_percentile(0));
+}
+
+TEST(QuantileStateTest, SerializedSizeRemainsStableAcrossSharedQueries) {
+ auto source = constant_state(10);
+ auto copy = source;
+ const auto size = source.get_serialized_size();
+ EXPECT_EQ(10, copy.get_value_by_percentile(0.5));
+ std::vector<uint8_t> bytes(size);
+ ASSERT_EQ(size, source.serialize(bytes.data()));
+ QuantileState restored(Slice(reinterpret_cast<char*>(bytes.data()),
bytes.size()));
+ EXPECT_EQ(10, restored.get_value_by_percentile(0.5));
+}
+
+static long serialized_digest_weight(const QuantileState& state) {
+ std::vector<uint8_t> bytes(state.get_serialized_size());
+ state.serialize(bytes.data());
+ TDigest digest(0);
+ digest.unserialize(bytes.data() + sizeof(float) + sizeof(uint8_t));
+ return digest.totalWeight();
+}
+
+TEST(QuantileStateTest, SelfMergeAndSharedSourceMerge) {
+ for (int count : {2, 4096}) {
+ auto state = constant_state(10, count);
+ auto unchanged = state;
+ state.merge(state);
+ state.merge(unchanged);
+ if (count == 2) {
+ EXPECT_EQ(6, state._explicit_data.size());
+ } else {
+ EXPECT_EQ(3 * serialized_digest_weight(unchanged),
serialized_digest_weight(state));
+ }
+ state.add_value(110);
+ EXPECT_EQ(10, state.get_value_by_percentile(0));
+ EXPECT_EQ(110, state.get_value_by_percentile(1));
+ EXPECT_EQ(10, unchanged.get_value_by_percentile(1));
+ }
+}
+
+static void expect_serialized_maximum(const QuantileState& state, double
expected) {
+ std::vector<uint8_t> bytes(state.get_serialized_size());
+ ASSERT_EQ(bytes.size(), state.serialize(bytes.data()));
+ QuantileState restored(Slice(reinterpret_cast<char*>(bytes.data()),
bytes.size()));
+ EXPECT_EQ(expected, restored.get_value_by_percentile(1));
+}
+
+TEST(QuantileStateTest, ConcurrentReadsSerializationAndIndependentWrites) {
+ auto source = constant_state(10);
+ std::vector<QuantileState> writers(4, source);
+ std::barrier start(8);
+ std::vector<std::thread> threads;
+ for (int i = 0; i < 2; ++i) {
+ threads.emplace_back([&] {
+ start.arrive_and_wait();
+ for (int j = 0; j < 100; ++j) {
+ EXPECT_EQ(10, source.get_value_by_percentile(0.5));
+ }
+ });
+ }
+ threads.emplace_back([&] {
+ start.arrive_and_wait();
+ for (int j = 0; j < 100; ++j) {
+ expect_serialized_maximum(source, 10);
+ }
+ });
+ for (int i = 0; i < 4; ++i) {
+ threads.emplace_back([&, i] {
+ start.arrive_and_wait();
+ for (int j = 0; j < 100; ++j) {
+ writers[i].add_value(110 + i);
+ writers[i].merge(source);
+ EXPECT_EQ(110 + i, writers[i].get_value_by_percentile(1));
+ }
+ });
+ }
+ start.arrive_and_wait();
+ for (auto& thread : threads) {
+ thread.join();
+ }
+ EXPECT_EQ(10, source.get_value_by_percentile(1));
+}
+
+TEST(QuantileStateTest, SharedSourceDetachesAndReadsCanOverlap) {
+ auto source = constant_state(10);
+ // A processed query only needs a shared lock; compression still needs an
exclusive lock.
+ ASSERT_EQ(10, source.get_value_by_percentile(0.5));
+ auto first = source;
+ auto second = source;
+ std::promise<void> first_entered;
+ std::promise<void> release_first;
+ auto released = release_first.get_future();
+ std::atomic<int> arrivals {0};
+ auto* sync = SyncPoint::get_instance();
+ SyncPoint::CallbackGuard guard;
+ sync->set_call_back(
+ "QuantileState::detach:source_locked",
+ [&](auto&&) {
+ if (arrivals.fetch_add(1) == 0) {
+ first_entered.set_value();
+ released.wait();
+ }
+ },
+ &guard);
+ sync->enable_processing();
+ std::thread first_thread([&] { first.add_value(-10); });
+ auto first_ready =
first_entered.get_future().wait_for(std::chrono::seconds(10));
+ std::promise<void> second_finished;
+ std::thread second_thread([&] {
+ second.add_value(110);
+ second_finished.set_value();
+ });
+ std::promise<double> read_finished;
+ auto read_result = read_finished.get_future();
+ std::thread reader([&] {
read_finished.set_value(source.get_value_by_percentile(0.5)); });
+ auto second_ready =
second_finished.get_future().wait_for(std::chrono::seconds(10));
+ auto read_ready = read_result.wait_for(std::chrono::seconds(10));
+ release_first.set_value();
+ first_thread.join();
+ second_thread.join();
+ reader.join();
+ sync->disable_processing();
+ EXPECT_EQ(std::future_status::ready, first_ready);
+ EXPECT_EQ(std::future_status::ready, second_ready);
+ EXPECT_EQ(std::future_status::ready, read_ready);
+ EXPECT_EQ(10, read_result.get());
+ EXPECT_EQ(-10, first.get_value_by_percentile(0));
+ EXPECT_EQ(10, first.get_value_by_percentile(1));
+ EXPECT_EQ(10, second.get_value_by_percentile(0));
+ EXPECT_EQ(110, second.get_value_by_percentile(1));
+ EXPECT_EQ(10, source.get_value_by_percentile(0.5));
+}
+
+TEST(QuantileStateTest, SharedSourceMergesCanOverlap) {
+ auto source = constant_state(10);
+ auto first = constant_state(-10);
+ auto second = constant_state(20);
+ std::promise<void> first_entered;
+ std::promise<void> second_entered;
+ std::promise<void> release_first;
+ auto released = release_first.get_future();
+ std::atomic<int> arrivals {0};
+ auto* sync = SyncPoint::get_instance();
+ SyncPoint::CallbackGuard guard;
+ sync->set_call_back(
+ "QuantileState::merge:source_locked",
+ [&](auto&&) {
+ if (arrivals.fetch_add(1) == 0) {
+ first_entered.set_value();
+ released.wait();
+ } else {
+ second_entered.set_value();
+ }
+ },
+ &guard);
+ sync->enable_processing();
+ std::thread first_thread([&] { first.merge(source); });
+ auto first_ready =
first_entered.get_future().wait_for(std::chrono::seconds(10));
+ std::thread second_thread([&] { second.merge(source); });
+ auto second_ready =
second_entered.get_future().wait_for(std::chrono::seconds(10));
+ release_first.set_value();
+ first_thread.join();
+ second_thread.join();
+ sync->disable_processing();
+ EXPECT_EQ(std::future_status::ready, first_ready);
+ EXPECT_EQ(std::future_status::ready, second_ready);
+ EXPECT_EQ(-10, first.get_value_by_percentile(0));
+ EXPECT_EQ(20, second.get_value_by_percentile(1));
+ EXPECT_EQ(10, source.get_value_by_percentile(0.5));
+}
+
+TEST(QuantileStateTest, ReadsLegacyUnprocessedDigestAndKeepsWireLayout) {
+ TDigest legacy(2048);
+ for (int i = 0; i < 4096; ++i) {
+ legacy.add(10);
+ }
+ constexpr size_t header_size = sizeof(float) + sizeof(uint8_t);
+ std::vector<uint8_t> bytes(header_size + legacy.serialized_size());
+ const float compression = 2048;
+ memcpy(bytes.data(), &compression, sizeof(compression));
+ bytes[sizeof(float)] = TDIGEST;
+ legacy.serialize(bytes.data() + header_size);
+ QuantileState state(Slice(reinterpret_cast<char*>(bytes.data()),
bytes.size()));
+ EXPECT_EQ(10, state.get_value_by_percentile(0.5));
+ bytes.resize(state.get_serialized_size());
+ ASSERT_EQ(bytes.size(), state.serialize(bytes.data()));
+ legacy.unserialize(bytes.data() + header_size);
+ EXPECT_EQ(4096, legacy.totalWeight());
+ EXPECT_EQ(10, legacy.quantile(0.5));
+}
+
TEST(QuantileStateTest, merge) {
QuantileState empty;
EXPECT_EQ(EMPTY, empty._type);
diff --git a/be/test/exprs/aggregate/agg_percentile_test.cpp
b/be/test/exprs/aggregate/agg_percentile_test.cpp
new file mode 100644
index 00000000000..97f74de21fc
--- /dev/null
+++ b/be/test/exprs/aggregate/agg_percentile_test.cpp
@@ -0,0 +1,171 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+#include <gtest/gtest.h>
+
+#include <memory>
+#include <vector>
+
+#include "core/column/column_complex.h"
+#include "core/data_type/data_type_quantilestate.h"
+#include "exprs/aggregate/aggregate_function_quantile_state.h"
+#include "util/defer_op.h"
+#include "util/tdigest.h"
+
+namespace doris {
+namespace {
+
+std::vector<double> expected_quantiles(const std::vector<double>& values,
+ const std::vector<double>& quantiles,
float compression) {
+ TDigest digest(compression);
+ for (double value : values) {
+ digest.add(value);
+ }
+ std::vector<double> result;
+ result.reserve(quantiles.size());
+ for (double quantile : quantiles) {
+ result.push_back(digest.quantile(quantile));
+ }
+ return result;
+}
+
+} // namespace
+
+TEST(AggregateFunctionQuantileStateTest, GrowingWindowKeepsResultsCompact) {
+ auto type = std::make_shared<DataTypeQuantileState>();
+ auto function = create_aggregate_function_quantile_state_union(
+ "quantile_union", {type}, type, false,
+ {.is_window_function = true, .column_names = {}});
+ std::unique_ptr<char[]> memory(new char[function->size_of_data()]);
+ auto* place = memory.get();
+ function->create(place);
+ Defer destroy([&] { function->destroy(place); });
+ Arena arena;
+ auto input = ColumnQuantileState::create();
+ QuantileState seed(10000);
+ constexpr size_t seed_count = 4096;
+ for (size_t i = 0; i < seed_count; ++i) {
+ seed.add_value(10);
+ }
+ input->insert_value(std::move(seed));
+ const IColumn* columns[] = {input.get()};
+ function->add(place, columns, 0, arena);
+ input->clear();
+ using Data =
AggregateFunctionQuantileStateData<AggregateFunctionQuantileStateUnionOp>;
+ auto& accumulator = reinterpret_cast<Data*>(place)->value;
+ auto* original = accumulator._tdigest_ptr.get();
+ const size_t initial_unprocessed =
accumulator._mutable_tdigest().unprocessed().size();
+ auto results = ColumnQuantileState::create();
+ constexpr size_t result_count = 12000;
+ results->reserve(result_count);
+ bool reused_accumulator = true;
+ size_t unprocessed_centroids = 0;
+ for (size_t row = 0; row < result_count; ++row) {
+ QuantileState value;
+ value.add_value(20 + row);
+ input->clear();
+ input->insert_value(std::move(value));
+ function->add(place, columns, 0, arena);
+ unprocessed_centroids +=
accumulator._mutable_tdigest().unprocessed().size();
+ function->insert_result_into(place, *results);
+ reused_accumulator &= accumulator._tdigest_ptr.get() == original;
+ }
+ EXPECT_TRUE(reused_accumulator);
+ // Each sample should enter result-insertion sorting only once across the
window.
+ EXPECT_EQ(initial_unprocessed + result_count, unprocessed_centroids);
+ RecordProperty("unprocessed_centroids",
std::to_string(unprocessed_centroids));
+ EXPECT_NE(accumulator._tdigest_ptr, results->get_element(result_count -
1)._tdigest_ptr);
+ // Inspect actual capacities independently of the column's approximate
accounting.
+ size_t digest_bytes = 0;
+ for (auto& result : results->get_data()) {
+ auto& digest = result._mutable_tdigest();
+ ASSERT_EQ(0, digest._unprocessed.capacity());
+ ASSERT_EQ(digest._processed.size(), digest._processed.capacity());
+ ASSERT_EQ(digest._cumulative.size(), digest._cumulative.capacity());
+ digest_bytes +=
+ (digest._processed.capacity() +
digest._unprocessed.capacity()) * sizeof(Centroid) +
+ digest._cumulative.capacity() * sizeof(Weight);
+ }
+ RecordProperty("retained_digest_bytes", std::to_string(digest_bytes));
+ EXPECT_LT(digest_bytes, result_count * 160 * 1024);
+ for (size_t row : {size_t(0), result_count / 2, result_count - 1}) {
+ auto& result = results->get_element(row);
+ EXPECT_EQ(0, result._mutable_tdigest().unprocessed().capacity());
+ EXPECT_EQ(10, result.get_value_by_percentile(0));
+ EXPECT_EQ(20 + row, result.get_value_by_percentile(1));
+ std::vector<double> values(seed_count, 10);
+ for (size_t i = 0; i <= row; ++i) {
+ values.push_back(20 + i);
+ }
+ const std::vector<double> quantiles {0.5, 0.9, 0.99};
+ const auto expected = expected_quantiles(values, quantiles, 10000);
+ for (size_t i = 0; i < quantiles.size(); ++i) {
+ EXPECT_NEAR(expected[i],
result.get_value_by_percentile(quantiles[i]), 1.0);
+ }
+ }
+ EXPECT_GE(accumulator._mutable_tdigest().unprocessed().capacity(), 80001);
+ EXPECT_FALSE(accumulator._mutable_tdigest().haveUnprocessed());
+}
+
+class AggregateFunctionQuantileStateRangeTest : public
testing::TestWithParam<bool> {};
+
+TEST_P(AggregateFunctionQuantileStateRangeTest,
RangeResultsShareOneSavedDigest) {
+ const bool is_window = GetParam();
+ auto type = std::make_shared<DataTypeQuantileState>();
+ auto function = create_aggregate_function_quantile_state_union(
+ "quantile_union", {type}, type, false,
+ {.is_window_function = is_window, .column_names = {}});
+ std::unique_ptr<char[]> memory(new char[function->size_of_data()]);
+ auto* place = memory.get();
+ function->create(place);
+ Defer destroy([&] { function->destroy(place); });
+ Arena arena;
+ auto input = ColumnQuantileState::create();
+ QuantileState state(10000);
+ for (int i = 0; i < 4096; ++i) {
+ state.add_value(10);
+ }
+ input->insert_value(state);
+ const IColumn* columns[] = {input.get()};
+ function->add(place, columns, 0, arena);
+ auto results = ColumnQuantileState::create();
+ results->insert_many_defaults(3);
+ function->insert_result_into_range(place, *results, 3, 12003);
+ function->insert_result_into_range(place, *results, 12003, 12003);
+ ASSERT_EQ(12003, results->size());
+ const size_t column_bytes = results->allocated_bytes();
+ const auto& first = results->get_element(3);
+ for (size_t row = 3; row < results->size(); ++row) {
+ EXPECT_EQ(first._tdigest_ptr, results->get_element(row)._tdigest_ptr);
+ }
+ if (is_window) {
+ EXPECT_NE(state._tdigest_ptr, first._tdigest_ptr);
+ } else {
+ EXPECT_EQ(state._tdigest_ptr, first._tdigest_ptr);
+ }
+ auto modified = first;
+ modified.add_value(110);
+ EXPECT_EQ(110, modified.get_value_by_percentile(1));
+ EXPECT_EQ(10, first.get_value_by_percentile(1));
+ EXPECT_EQ(10, state.get_value_by_percentile(1));
+ EXPECT_EQ(column_bytes, results->allocated_bytes());
+}
+
+INSTANTIATE_TEST_SUITE_P(AggregateAndWindow,
AggregateFunctionQuantileStateRangeTest,
+ testing::Bool());
+
+} // namespace doris
diff --git a/be/test/util/tdigest_test.cpp b/be/test/util/tdigest_test.cpp
index efb40774324..0804fa3ad7f 100644
--- a/be/test/util/tdigest_test.cpp
+++ b/be/test/util/tdigest_test.cpp
@@ -76,6 +76,34 @@ static double quantile(const double q, const
std::vector<double>& values) {
return q1;
}
+TEST_F(TDigestTest, CopyPreservesValuesWithoutSpareCapacity) {
+ const auto allocated_bytes = [](const TDigest& digest) {
+ return (digest._processed.capacity() + digest._unprocessed.capacity())
* sizeof(Centroid) +
+ digest._cumulative.capacity() * sizeof(Weight);
+ };
+ TDigest source(10000);
+ for (int i = 0; i < 300; ++i) {
+ source.add(i);
+ }
+ source.compress();
+ TDigest processed_copy(source);
+ EXPECT_EQ(0, processed_copy.unprocessed().capacity());
+ EXPECT_LT(allocated_bytes(processed_copy), 16 * 1024);
+ source.add(1000);
+ TDigest copy(source);
+ EXPECT_EQ(301, copy.totalWeight());
+ EXPECT_LT(allocated_bytes(copy), 16 * 1024);
+ // Continued writes may grow the copy's buffers, without changing the
source.
+ for (int i = 0; i < 1000; ++i) {
+ copy.add(2000);
+ }
+ EXPECT_EQ(1301, copy.totalWeight());
+ EXPECT_EQ(0, copy.quantile(0));
+ EXPECT_EQ(2000, copy.quantile(1));
+ EXPECT_EQ(1000, source.quantile(1));
+ EXPECT_EQ(299, processed_copy.quantile(1));
+}
+
TEST_F(TDigestTest, CrashAfterMerge) {
TDigest digest(1000);
std::uniform_real_distribution<> reals(0.0, 1.0);
diff --git
a/regression-test/data/query_p0/sql_functions/quantile_state_functions/test_quantile_state_function.out
b/regression-test/data/query_p0/sql_functions/quantile_state_functions/test_quantile_state_function.out
index 1e9b9630e05..42cd7965e97 100644
---
a/regression-test/data/query_p0/sql_functions/quantile_state_functions/test_quantile_state_function.out
+++
b/regression-test/data/query_p0/sql_functions/quantile_state_functions/test_quantile_state_function.out
@@ -35,3 +35,35 @@ AAAARQEAAAAAAAAkQA==
-- !sql_quantile_state_base64_12 --
true
+-- !cow_independent --
+0 10 10 10 10 110
+1 110 110 110 10 110
+
+-- !cow_shared --
+0 10 10 10 10 110
+1 110 110 110 10 110
+
+-- !cow_window --
+0 0 10 10
+0 1 10 10
+0 2 10 10
+0 3 10 10
+1 0 10 110
+1 1 10 110
+1 2 10 110
+1 3 10 110
+
+-- !cow_join_expand --
+0 10
+1 110
+
+-- !cow_explode_expand --
+0 10
+1 110
+
+-- !cow_long_high_compression_window --
+0 10 10
+1 10 20
+6000 10 6019
+12000 10 12019
+
diff --git
a/regression-test/suites/query_p0/sql_functions/quantile_state_functions/test_quantile_state_function.groovy
b/regression-test/suites/query_p0/sql_functions/quantile_state_functions/test_quantile_state_function.groovy
index b1480f0971b..d01983f5c69 100644
---
a/regression-test/suites/query_p0/sql_functions/quantile_state_functions/test_quantile_state_function.groovy
+++
b/regression-test/suites/query_p0/sql_functions/quantile_state_functions/test_quantile_state_function.groovy
@@ -79,4 +79,94 @@ suite("test_quantile_state_function") {
)
) = quantile_state_to_base64(to_quantile_state(10.0, 2048))
"""
+
+ sql "DROP TABLE IF EXISTS test_quantile_state_cow"
+ sql """
+ CREATE TABLE test_quantile_state_cow (
+ topic INT NOT NULL,
+ chunk INT NOT NULL,
+ q QUANTILE_STATE QUANTILE_UNION NOT NULL
+ ) AGGREGATE KEY(topic, chunk)
+ DISTRIBUTED BY HASH(topic, chunk) BUCKETS 4
+ PROPERTIES("replication_num" = "1")
+ """
+ // Each stored state has 4096 inputs, exceeding the EXPLICIT-state limit.
+ sql """
+ INSERT INTO test_quantile_state_cow
+ SELECT number % 2, (number DIV 2) % 4,
+ quantile_union(to_quantile_state(10 + 100 * (number % 2), 2048))
+ FROM numbers("number" = "32768")
+ GROUP BY 1, 2
+ """
+
+ def query = """
+ WITH shared AS (SELECT * FROM test_quantile_state_cow),
+ topics AS (SELECT topic, quantile_union(q) AS q FROM shared GROUP BY
topic),
+ total AS (SELECT quantile_union(q) AS q FROM shared)
+ SELECT topic, quantile_percent(topics.q, 0),
quantile_percent(topics.q, 0.5),
+ quantile_percent(topics.q, 1),
+ quantile_percent(total.q, 0), quantile_percent(total.q, 1)
+ FROM topics CROSS JOIN total ORDER BY topic
+ """
+ sql "SET enable_cte_materialize = false"
+ qt_cow_independent query
+ sql "SET enable_cte_materialize = true"
+ sql "SET inline_cte_referenced_threshold = 0"
+ explain {
+ sql(query)
+ contains "MultiCastDataSinks"
+ }
+ // Both plans are checked against fixed results: topic values are 10 or
110.
+ qt_cow_shared query
+
+ qt_cow_window """
+ SELECT topic, chunk, quantile_percent(q, 0), quantile_percent(q, 1)
+ FROM (
+ SELECT topic, chunk,
+ quantile_union(q) OVER (ORDER BY topic, chunk
+ ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS q
+ FROM test_quantile_state_cow
+ ) t ORDER BY topic, chunk
+ """
+
+ // Each state has 4096 samples. Process the shared 10 state first so both
+ // groups start from it before group 1 merges 110; group 0 must stay at 10.
+ qt_cow_join_expand """
+ SELECT g.k, quantile_percent(quantile_union(t.q), 1)
+ FROM (
+ SELECT number % 2 AS id,
+ quantile_union(to_quantile_state(10 + 100 * (number % 2),
2048)) AS q
+ FROM numbers("number" = "8192") GROUP BY 1 ORDER BY 1 LIMIT 2
+ ) t
+ JOIN (SELECT 0 AS k UNION ALL SELECT 1) g ON t.id = 0 OR g.k = 1
+ GROUP BY g.k ORDER BY g.k
+ """
+ qt_cow_explode_expand """
+ SELECT k, quantile_percent(quantile_union(q), 1)
+ FROM (
+ SELECT number % 2 AS id,
+ quantile_union(to_quantile_state(10 + 100 * (number % 2),
2048)) AS q
+ FROM numbers("number" = "8192") GROUP BY 1 ORDER BY 1 LIMIT 2
+ ) t
+ LATERAL VIEW explode(IF(id = 0, [0, 1], [1])) e AS k
+ GROUP BY k ORDER BY k
+ """
+
+ qt_cow_long_high_compression_window """
+ SELECT id, quantile_percent(q, 0), quantile_percent(q, 1)
+ FROM (
+ SELECT id, quantile_union(q) OVER (
+ ORDER BY id ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW)
AS q
+ FROM (
+ SELECT 0 AS id, quantile_union(to_quantile_state(10, 10000))
AS q
+ FROM numbers("number" = "4096")
+ UNION ALL
+ SELECT number + 1 AS id, to_quantile_state(20 + number, 10000)
AS q
+ FROM numbers("number" = "12000")
+ ) input
+ ) window_results
+ WHERE id IN (0, 1, 6000, 12000)
+ ORDER BY id
+ """
+
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]