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]