This is an automated email from the ASF dual-hosted git repository.

HappenLee pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/master by this push:
     new da4ec5dca33 [fix](be) Prevent shared quantile state mutation with 
TDigest COW (#68498)
da4ec5dca33 is described below

commit da4ec5dca33e80ea70f3daf2a888fbdc759a432d
Author: linrrarity <[email protected]>
AuthorDate: Fri Oct 2 17:36:12 2026 +0800

    [fix](be) Prevent shared quantile state mutation with TDigest COW (#68498)
    
    ### What problem does this PR solve?
    
    Issue Number: close #xxx
    
    Related PR: #xxx
    
    Problem Summary:
    
    Copies of `QuantileState` share a `TDigest`. When CTEs, joins, or
    `explode` reuse a state, updating one copy can contaminate another
    group's quantile result. Concurrent percentile queries can also race
    while compressing the shared digest.
    
    Add copy-on-write in `QuantileState`, using the current reference count
    to detach before sample writes. Synchronize compression and shared
    reads, and preserve reserved vector capacity when copying `TDigest`.
    Serialization keeps its existing binary format; sizing finishes pending
    compression so subsequent queries cannot change the serialized length.
    This moves compression earlier and can affect approximate results.
    
    
    ```sql
    SET enable_cte_materialize = true;
    SET inline_cte_referenced_threshold = 0;
    
    WITH t AS (
        SELECT number % 2 AS g,
               quantile_union(to_quantile_state(100 * (number % 2), 2048)) AS q
        FROM numbers("number" = "8192")
        GROUP BY g
    )
    SELECT g, quantile_percent(q, 0.5) FROM t
    UNION ALL
    SELECT -1, quantile_percent(quantile_union(q), 0.5) FROM t;
    
    ```
    
    before:
    ```text
    std::vector<doris::Centroid>::operator[](size_type) const: Assertion '__n < 
this->size()' failed.
    *** Query id: e5437298fa6e4f14-9b6a546907f92425 ***
    *** tablet id: 0 ***
    *** Aborted at 1790250779 (unix time) try "date -d @1790250779" if you are 
using GNU date ***
    *** Current BE git commitID: 765e2dfefe ***
    *** SIGABRT unknown detail explain (@0x3fe0007a9c2) received by PID 502210 
(TID 547405 OR 0x113767afd640) from PID 502210; stack trace: ***
    F20260924 19:52:59.574612 547409 tdigest.h:458] Check failed: index >= 
_processed_weight - _weight(n - 1) / 2.0 (4095 vs. 8189)
    *** Check failure stack trace: ***
        @     0x560d751ce0b6  google::LogMessageFatal::~LogMessageFatal()
        @     0x560d4c8a1600  doris::TDigest::quantile_processed()
        @     0x560d4c884644  doris::TDigest::quantile()
        @     0x560d4c876beb  doris::QuantileState::get_value_by_percentile()
        @     0x560d555b14b5  
doris::FunctionQuantileStatePercent::execute_impl()
        @     0x560d555b1902  
doris::FunctionQuantileStatePercent::execute_impl()
        @     0x560d563a3d19  
doris::PreparedFunctionImpl::_execute_skipped_constant_deal()
        @     0x560d55df44cb  doris::PreparedFunctionImpl::default_execute()
        @     0x560d55d4aa7c  doris::PreparedFunctionImpl::execute()
        @     0x560d4adbbe7f  doris::IFunctionBase::execute()
        @     0x560d53f79b64  doris::VectorizedFnCall::_do_execute()
        @     0x560d53f147a3  doris::VectorizedFnCall::execute_column_impl()
        @     0x560d53f3d44c  doris::VExpr::execute_column()
        @     0x560d53f9d63b  doris::VExprContext::execute()
        @     0x560d535ee6cf  doris::OperatorXBase::do_projections()
        @     0x560d535f091a  doris::OperatorXBase::get_block_after_projects()
        @     0x560d4f121015  doris::PipelineTask::execute()
        @     0x560d51b28ca0  doris::TaskScheduler::_do_work()
        @     0x560d51d357c5  doris::TaskScheduler::start()::$_0::operator()()
        @     0x560d51d3570d  std::__invoke_impl<>()
        @     0x560d51d3562d  
_ZSt10__invoke_rIvRZN5doris13TaskScheduler5startEvE3$_0JEENSt9enable_ifIX16is_invocable_r_vIT_T0_DpT1_EES5_E4typeEOS6_DpOS7_
        @     0x560d51d35305  std::_Function_handler<>::_M_invoke()
        @     0x560d4a838b3e  std::function<>::operator()()
        @     0x560d710d3339  doris::FunctionRunnable::run()
        @     0x560d71013f64  doris::ThreadPool::dispatch_thread()
        @     0x560d710f38fd  std::__invoke_impl<>()
        @     0x560d710f36b5  std::__invoke<>()
        @     0x560d710f35e1  
_ZNSt5_BindIFMN5doris10ThreadPoolEFvvEPS1_EE6__callIvJEJLm0EEEET_OSt5tupleIJDpT0_EESt12_Index_tupleIJXspT1_EEE
        @     0x560d710f339c  std::_Bind<>::operator()<>()
        @     0x560d710f328d  std::__invoke_impl<>()
        @     0x560d710f318d  
_ZSt10__invoke_rIvRSt5_BindIFMN5doris10ThreadPoolEFvvEPS2_EEJEENSt9enable_ifIX16is_invocable_r_vIT_T0_DpT1_EESA_E4typeEOSB_DpOSC_
        @     0x560d710f2a65  std::_Function_handler<>::_M_invoke()
     0# doris::signal::(anonymous namespace)::FailureSignalHandler(int, 
siginfo_t*, void*) at ../src/common/signal_handler.h:418
     1# 0x00001543D423FC60 in /lib64/libc.so.6
     2# __pthread_kill_implementation in /lib64/libc.so.6
     3# gsignal in /lib64/libc.so.6
     4# abort in /lib64/libc.so.6
     5# 0x0000560D7AD406A7 in 
/mnt/disk9/linzhenqi/dv/doris/output/be/lib/doris_be
     6# std::vector<doris::Centroid, std::allocator<doris::Centroid> 
>::operator[](unsigned long) const at 
/mnt/disk9/linzhenqi/dv/version-toolchain/ldb_toolchain_v28/bin/../lib/gcc/x86_64-pc-linux-gnu/15/include/g++-v15/bits/stl_vector.h:1282
     7# doris::TDigest::_weight(long) const at ../src/util/tdigest.h:623
     8# doris::TDigest::_update_cumulative() at ../src/util/tdigest.h:701
     9# doris::TDigest::_process() at ../src/util/tdigest.h:749
    10# doris::TDigest::quantile(float) in 
/mnt/disk9/linzhenqi/dv/doris/output/be/lib/doris_be
    11# doris::QuantileState::get_value_by_percentile(float) const at 
./be/src/core/value/quantile_state.cpp:140
    12# 
doris::FunctionQuantileStatePercent::execute_impl(doris::FunctionContext*, 
doris::Block&, std::vector<unsigned int, std::allocator<unsigned int> > const&, 
unsigned int, unsigned long) const at 
./be/src/exprs/function/function_quantile_state.cpp:202
    13# non-virtual thunk to 
doris::FunctionQuantileStatePercent::execute_impl(doris::FunctionContext*, 
doris::Block&, std::vector<unsigned int, std::allocator<unsigned int> > const&, 
unsigned int, unsigned long) const in 
/mnt/disk9/linzhenqi/dv/doris/output/be/lib/doris_be
    14# 
doris::PreparedFunctionImpl::_execute_skipped_constant_deal(doris::FunctionContext*,
 doris::Block&, std::vector<unsigned int, std::allocator<unsigned int> > 
const&, unsigned int, unsigned long) const at 
./be/src/exprs/function/function.cpp:135
    15# doris::PreparedFunctionImpl::default_execute(doris::FunctionContext*, 
doris::Block&, std::vector<unsigned int, std::allocator<unsigned int> > const&, 
unsigned int, unsigned long) const at ./be/src/exprs/function/function.cpp:268
    16# doris::PreparedFunctionImpl::execute(doris::FunctionContext*, 
doris::Block&, std::vector<unsigned int, std::allocator<unsigned int> > const&, 
unsigned int, unsigned long) const at ./be/src/exprs/function/function.cpp:274
    17# doris::IFunctionBase::execute(doris::FunctionContext*, doris::Block&, 
std::vector<unsigned int, std::allocator<unsigned int> > const&, unsigned int, 
unsigned long) const at ../src/exprs/function/function.h:213
    18# doris::VectorizedFnCall::_do_execute(doris::VExprContext*, doris::Block 
const*, doris::PODArray<unsigned int, 4096ul, doris::Allocator<false, false, 
false, doris::DefaultMemoryAllocator, true>, 16ul, 15ul> const*, unsigned long, 
doris::COW<doris::IColumn>::immutable_ptr<doris::IColumn>&, 
doris::COW<doris::IColumn>::immutable_ptr<doris::IColumn>*) const at 
./be/src/exprs/vectorized_fn_call.cpp:441
    19# doris::VectorizedFnCall::execute_column_impl(doris::VExprContext*, 
doris::Block const*, doris::PODArray<unsigned int, 4096ul, 
doris::Allocator<false, false, false, doris::DefaultMemoryAllocator, true>, 
16ul, 15ul> const*, unsigned long, 
doris::COW<doris::IColumn>::immutable_ptr<doris::IColumn>&) const at 
./be/src/exprs/vectorized_fn_call.cpp:477
    20# doris::VExpr::execute_column(doris::VExprContext*, doris::Block const*, 
doris::PODArray<unsigned int, 4096ul, doris::Allocator<false, false, false, 
doris::DefaultMemoryAllocator, true>, 16ul, 15ul> const*, unsigned long, 
doris::COW<doris::IColumn>::immutable_ptr<doris::IColumn>&) const at 
./be/src/exprs/vexpr.cpp:1067
    21# doris::VExprContext::execute(doris::Block const*, 
doris::COW<doris::IColumn>::immutable_ptr<doris::IColumn>&) at 
./be/src/exprs/vexpr_context.cpp:91
    22# 
doris::VExprContext::get_output_block_after_execute_exprs(std::vector<std::shared_ptr<doris::VExprContext>,
 std::allocator<std::shared_ptr<doris::VExprContext> > > const&, doris::Block 
const&, doris::Block*, bool) at ./be/src/exprs/vexpr_context.cpp:464
    23# 
doris::MultiCastDataStreamerSourceOperatorX::get_block_impl(doris::RuntimeState*,
 doris::Block*, bool*) at 
./be/src/exec/operator/multi_cast_data_stream_source.cpp:112
    24# doris::OperatorXBase::get_block(doris::RuntimeState*, doris::Block*, 
bool*) at ../src/exec/operator/operator.h:898
    25# doris::OperatorXBase::get_block_after_projects(doris::RuntimeState*, 
doris::Block*, bool*) at ./be/build_ASAN/../src/exec/operator/operator.cpp:435
    26# doris::PipelineTask::execute(bool*) at 
./be/src/exec/pipeline/pipeline_task.cpp:655
    27# doris::TaskScheduler::_do_work(int) at 
./be/src/exec/pipeline/task_scheduler.cpp:155
    28# doris::TaskScheduler::start()::$_0::operator()() const at 
./be/src/exec/pipeline/task_scheduler.cpp:64
    29# void std::__invoke_impl<void, 
doris::TaskScheduler::start()::$_0&>(std::__invoke_other, 
doris::TaskScheduler::start()::$_0&) at 
/mnt/disk9/linzhenqi/dv/version-toolchain/ldb_toolchain_v28/bin/../lib/gcc/x86_64-pc-linux-gnu/15/include/g++-v15/bits/invoke.h:63
    30# std::enable_if<is_invocable_r_v<void, 
doris::TaskScheduler::start()::$_0&>, void>::type std::__invoke_r<void, 
doris::TaskScheduler::start()::$_0&>(doris::TaskScheduler::start()::$_0&) at 
/mnt/disk9/linzhenqi/dv/version-toolchain/ldb_toolchain_v28/bin/../lib/gcc/x86_64-pc-linux-gnu/15/include/g++-v15/bits/invoke.h:119
    31# std::_Function_handler<void (), 
doris::TaskScheduler::start()::$_0>::_M_invoke(std::_Any_data const&) at 
/mnt/disk9/linzhenqi/dv/version-toolchain/ldb_toolchain_v28/bin/../lib/gcc/x86_64-pc-linux-gnu/15/include/g++-v15/bits/std_function.h:292
    32# std::function<void ()>::operator()() const at 
/mnt/disk9/linzhenqi/dv/version-toolchain/ldb_toolchain_v28/bin/../lib/gcc/x86_64-pc-linux-gnu/15/include/g++-v15/bits/std_function.h:593
    33# doris::FunctionRunnable::run() at ./be/src/util/threadpool.cpp:60
    34# doris::ThreadPool::dispatch_thread() at ./be/src/util/threadpool.cpp:621
    35# void std::__invoke_impl<void, void (doris::ThreadPool::*&)(), 
doris::ThreadPool*&>(std::__invoke_memfun_deref, void 
(doris::ThreadPool::*&)(), doris::ThreadPool*&) at 
/mnt/disk9/linzhenqi/dv/version-toolchain/ldb_toolchain_v28/bin/../lib/gcc/x86_64-pc-linux-gnu/15/include/g++-v15/bits/invoke.h:76
    36# std::__invoke_result<void (doris::ThreadPool::*&)(), 
doris::ThreadPool*&>::type std::__invoke<void (doris::ThreadPool::*&)(), 
doris::ThreadPool*&>(void (doris::ThreadPool::*&)(), doris::ThreadPool*&) at 
/mnt/disk9/linzhenqi/dv/version-toolchain/ldb_toolchain_v28/bin/../lib/gcc/x86_64-pc-linux-gnu/15/include/g++-v15/bits/invoke.h:98
    37# void std::_Bind<void 
(doris::ThreadPool::*(doris::ThreadPool*))()>::__call<void, , 
0ul>(std::tuple<>&&, std::_Index_tuple<0ul>) at 
/mnt/disk9/linzhenqi/dv/version-toolchain/ldb_toolchain_v28/bin/../lib/gcc/x86_64-pc-linux-gnu/15/include/g++-v15/functional:515
    38# void std::_Bind<void 
(doris::ThreadPool::*(doris::ThreadPool*))()>::operator()<, void>() at 
/mnt/disk9/linzhenqi/dv/version-toolchain/ldb_toolchain_v28/bin/../lib/gcc/x86_64-pc-linux-gnu/15/include/g++-v15/functional:600
    39# void std::__invoke_impl<void, std::_Bind<void 
(doris::ThreadPool::*(doris::ThreadPool*))()>&>(std::__invoke_other, 
std::_Bind<void (doris::ThreadPool::*(doris::ThreadPool*))()>&) at 
/mnt/disk9/linzhenqi/dv/version-toolchain/ldb_toolchain_v28/bin/../lib/gcc/x86_64-pc-linux-gnu/15/include/g++-v15/bits/invoke.h:63
    40# std::enable_if<is_invocable_r_v<void, std::_Bind<void 
(doris::ThreadPool::*(doris::ThreadPool*))()>&>, void>::type 
std::__invoke_r<void, std::_Bind<void 
(doris::ThreadPool::*(doris::ThreadPool*))()>&>(std::_Bind<void 
(doris::ThreadPool::*(doris::ThreadPool*))()>&) at 
/mnt/disk9/linzhenqi/dv/version-toolchain/ldb_toolchain_v28/bin/../lib/gcc/x86_64-pc-linux-gnu/15/include/g++-v15/bits/invoke.h:119
    41# std::_Function_handler<void (), std::_Bind<void 
(doris::ThreadPool::*(doris::ThreadPool*))()> >::_M_invoke(std::_Any_data 
const&) at 
/mnt/disk9/linzhenqi/dv/version-toolchain/ldb_toolchain_v28/bin/../lib/gcc/x86_64-pc-linux-gnu/15/include/g++-v15/bits/std_function.h:292
    42# std::function<void ()>::operator()() const at 
/mnt/disk9/linzhenqi/dv/version-toolchain/ldb_toolchain_v28/bin/../lib/gcc/x86_64-pc-linux-gnu/15/include/g++-v15/bits/std_function.h:593
    43# doris::Thread::supervise_thread(void*) at ./be/src/util/thread.cpp:460
    44# asan_thread_start(void*) in 
/mnt/disk9/linzhenqi/dv/doris/output/be/lib/doris_be
    45# start_thread in /lib64/libc.so.6
    46# clone3 in /lib64/libc.so.6
    
    
    ```
    
    now:
    ```text
    +------+--------------------------+
    | g    | quantile_percent(q, 0.5) |
    +------+--------------------------+
    |   -1 |                        0 |
    |    0 |                        0 |
    |    1 |                      100 |
    +------+--------------------------+
    ```
    
    ### Release note
    
    Fix backend crashes and incorrect quantile results when aggregate states
    share inputs, including materialized CTE queries.
---
 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  |  22 +-
 be/src/util/tdigest.h                              |   2 +
 be/test/core/value/quantile_state_test.cpp         | 294 +++++++++++++++++++++
 be/test/exprs/aggregate/agg_percentile_test.cpp    | 126 +++++++++
 be/test/util/tdigest_test.cpp                      |  28 ++
 .../test_quantile_state_function.out               |  32 +++
 .../test_quantile_state_function.groovy            |  90 +++++++
 10 files changed, 708 insertions(+), 28 deletions(-)

diff --git a/be/src/core/value/quantile_state.cpp 
b/be/src/core/value/quantile_state.cpp
index a1821de29f9..c5008f59616 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,7 +30,68 @@
 #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();
+        // 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) {}
@@ -50,10 +113,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;
 }
 
@@ -137,7 +204,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.quantile_processed(percentile);
     }
     default:
         break;
@@ -186,8 +254,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:
@@ -224,7 +292,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;
     }
@@ -235,6 +304,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;
@@ -256,23 +330,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;
         }
@@ -287,18 +363,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;
         }
@@ -322,9 +406,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();
@@ -335,7 +419,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 f980c4fe1c5..18ea406de97 100644
--- a/be/src/exprs/aggregate/aggregate_function_quantile_state.cpp
+++ b/be/src/exprs/aggregate/aggregate_function_quantile_state.cpp
@@ -31,11 +31,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 8f90d91cdfe..7f5f09710bd 100644
--- a/be/src/exprs/aggregate/aggregate_function_quantile_state.h
+++ b/be/src/exprs/aggregate/aggregate_function_quantile_state.h
@@ -102,10 +102,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,10 +145,25 @@ public:
 
     void insert_result_into(ConstAggregateDataPtr __restrict place, IColumn& 
to) const override {
         auto& column = assert_cast<ColVecResult&, 
TypeCheckOnRelease::DISABLE>(to);
-        column.get_data().push_back(this->data(place).get());
+        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 cd5570edeb6..1d5202a2d9b 100644
--- a/be/src/util/tdigest.h
+++ b/be/src/util/tdigest.h
@@ -170,6 +170,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..80ac1aa1ccf 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.total_weight();
+}
+
+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.total_weight());
+    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
index bfb911bdb4c..28c31f0737f 100644
--- a/be/test/exprs/aggregate/agg_percentile_test.cpp
+++ b/be/test/exprs/aggregate/agg_percentile_test.cpp
@@ -23,14 +23,17 @@
 #include <vector>
 
 #include "core/column/column_array.h"
+#include "core/column/column_complex.h"
 #include "core/column/column_nullable.h"
 #include "core/column/column_string.h"
 #include "core/column/column_vector.h"
 #include "core/data_type/data_type_array.h"
 #include "core/data_type/data_type_nullable.h"
 #include "core/data_type/data_type_number.h"
+#include "core/data_type/data_type_quantilestate.h"
 #include "core/string_buffer.hpp"
 #include "exprs/aggregate/aggregate_function_percentile.h"
+#include "exprs/aggregate/aggregate_function_quantile_state.h"
 #include "exprs/aggregate/aggregate_function_simple_factory.h"
 #include "util/tdigest.h"
 
@@ -110,6 +113,129 @@ void expect_results_equal(const std::vector<double>& 
actual, const std::vector<d
 
 } // 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().have_unprocessed());
+}
+
+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());
+
 TEST(AggregateFunctionPercentileApproxArrayTest, AddAndBatchPaths) {
     const std::vector<double> values {1, 2, 3, 4, 5, 100, 
std::numeric_limits<double>::quiet_NaN()};
     const std::vector<double> quantiles {0.9, 0.0, 0.5, 0.5, 1.0};
diff --git a/be/test/util/tdigest_test.cpp b/be/test/util/tdigest_test.cpp
index b5fec84ba48..392724e1862 100644
--- a/be/test/util/tdigest_test.cpp
+++ b/be/test/util/tdigest_test.cpp
@@ -77,6 +77,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.total_weight());
+    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.total_weight());
+    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