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]

Reply via email to