github-actions[bot] commented on code in PR #68498:
URL: https://github.com/apache/doris/pull/68498#discussion_r4093607828


##########
be/src/core/value/quantile_state.h:
##########
@@ -67,12 +67,12 @@ class QuantileState {
         throw Exception(Status::NotSupported("QuantileState::to_string() not 
implemented"));
     }
 #endif
-    ~QuantileState() = default;
+    ~QuantileState();
 
 private:
     QuantileStateType _type = EMPTY;
-    std::shared_ptr<TDigest> _tdigest_ptr;
-    double _single_data;
+    std::unique_ptr<TDigest> _tdigest_ptr;

Review Comment:
   This member-size change also changes the raw bytes emitted by 
`DataTypeQuantileStateSerDe::write_column_to_pb()`: `get_data_at()` sends 
`sizeof(QuantileState)`, while `insert_data()` ignores the received length and 
reinterprets it as the local layout; scalar `FunctionRPC` uses this path for 
arguments/results and RPC UDAFs use it for inputs. The existing PB test shows 
EMPTY/SINGLE round-tripping today, but an old reader sees the new SINGLE bits 
as part of its two-pointer `shared_ptr`, and a new reader uses old member bytes 
at the shifted scalar/vector offsets. Please preserve the outer 
`shared_ptr<TDigest>` layout while allocating a distinct TDigest wrapper on 
copy, or version and migrate this PB path to `QuantileState::serialize()`, with 
an old/new SINGLE fixture.



##########
be/src/util/tdigest.h:
##########
@@ -422,222 +442,336 @@ class TDigest {
             return NAN;
         }
 
-        if (_processed.size() == 0) {
+        if (_data->_processed.size() == 0) {
             // no sorted means no data, no way to get a quantile
             return NAN;
-        } else if (_processed.size() == 1) {
+        } else if (_data->_processed.size() == 1) {
             // with one data point, all quantiles lead to Rome
 
             return _mean(0);
         }
 
         // we know that there are at least two sorted now
-        auto n = _processed.size();
+        auto n = _data->_processed.size();
 
         // if values were stored in a sorted array, index would be the offset 
we are Weighterested in
-        const auto index = q * _processed_weight;
+        const auto index = q * _data->_processed_weight;
 
-        // at the boundaries, we return _min or _max
+        // at the boundaries, we return _data->_min or _data->_max
         if (index <= _weight(0) / 2.0) {
             DCHECK_GT(_weight(0), 0);
-            return static_cast<Value>(_min + 2.0 * index / _weight(0) * 
(_mean(0) - _min));
+            return static_cast<Value>(_data->_min +
+                                      2.0 * index / _weight(0) * (_mean(0) - 
_data->_min));
         }
 
-        auto iter = std::lower_bound(_cumulative.cbegin(), _cumulative.cend(), 
index);
+        auto iter = std::lower_bound(_data->_cumulative.cbegin(), 
_data->_cumulative.cend(), index);
 
-        if (iter != _cumulative.cend() && iter != _cumulative.cbegin() &&
-            iter + 1 != _cumulative.cend()) {
-            auto i = std::distance(_cumulative.cbegin(), iter);
+        if (iter != _data->_cumulative.cend() && iter != 
_data->_cumulative.cbegin() &&
+            iter + 1 != _data->_cumulative.cend()) {
+            auto i = std::distance(_data->_cumulative.cbegin(), iter);
             auto z1 = index - *(iter - 1);
             auto z2 = *(iter)-index;
             // VLOG_CRITICAL << "z2 " << z2 << " index " << index << " z1 " << 
z1;
             return _weighted_average(_mean(i - 1), z2, _mean(i), z1);
         }
 
-        DCHECK_LE(index, _processed_weight);
-        DCHECK_GE(index, _processed_weight - _weight(n - 1) / 2.0);
+        DCHECK_LE(index, _data->_processed_weight);
+        DCHECK_GE(index, _data->_processed_weight - _weight(n - 1) / 2.0);
 
-        auto z1 = static_cast<Value>(index - _processed_weight - _weight(n - 
1) / 2.0);
+        auto z1 = static_cast<Value>(index - _data->_processed_weight - 
_weight(n - 1) / 2.0);
         auto z2 = static_cast<Value>(_weight(n - 1) / 2 - z1);
-        return _weighted_average(_mean(n - 1), z1, _max, z2);
+        return _weighted_average(_mean(n - 1), z1, _data->_max, z2);
     }
 
-    Value compression() const { return _compression; }
+    Value compression() const { return _data->_compression; }
 
     void add(Value x) { add(x, 1); }
 
-    void compress() { _process(); }
+    void compress() {
+        if (total_size() != 0) {
+            _prepare_for_write();
+            _process();
+        }
+    }
 
     // add a single centroid to the unprocessed vector, processing previously 
unprocessed sorted if our limit has
     // been reached.
     bool add(Value x, Weight w) {
         if (std::isnan(x)) {
             return false;
         }
-        _unprocessed.emplace_back(x, w);
-        _unprocessed_weight += w;
+        _prepare_for_write();
+        _data->_unprocessed.emplace_back(x, w);
+        _data->_unprocessed_weight += w;
         _process_if_necessary();
         return true;
     }
 
     void add(std::vector<Centroid>::const_iterator iter,
              std::vector<Centroid>::const_iterator end) {
+        const std::vector<Centroid> centroids(iter, end);
+        _prepare_for_write();
+        iter = centroids.cbegin();
+        end = centroids.cend();
         while (iter != end) {
             const size_t diff = std::distance(iter, end);
-            const size_t room = _max_unprocessed - _unprocessed.size();
+            const size_t room = _data->_max_unprocessed - 
_data->_unprocessed.size();
             auto mid = iter + std::min(diff, room);
             while (iter != mid) {
-                _unprocessed.push_back(*(iter++));
+                _data->_unprocessed_weight += iter->weight();
+                _data->_unprocessed.push_back(*(iter++));
             }
-            if (_unprocessed.size() >= _max_unprocessed) {
+            if (_data->_unprocessed.size() >= _data->_max_unprocessed) {
                 _process();
             }
         }
     }
 
-    uint32_t serialized_size() {
+    uint32_t serialized_size() const {
         return static_cast<uint32_t>(sizeof(uint32_t) + sizeof(Value) * 5 + 
sizeof(Index) * 2 +
-                                     sizeof(uint32_t) * 3 + _processed.size() 
* sizeof(Centroid) +
-                                     _unprocessed.size() * sizeof(Centroid) +
-                                     _cumulative.size() * sizeof(Weight));
+                                     sizeof(uint32_t) * 3 +
+                                     _data->_processed.size() * 
sizeof(Centroid) +
+                                     _data->_unprocessed.size() * 
sizeof(Centroid) +
+                                     _data->_cumulative.size() * 
sizeof(Weight));
     }
 
-    size_t serialize(uint8_t* writer) {
+    size_t serialize(uint8_t* writer) const {
         uint8_t* dst = writer;
         uint32_t total_size = serialized_size();
         memcpy(writer, &total_size, sizeof(uint32_t));
         writer += sizeof(uint32_t);
-        memcpy(writer, &_compression, sizeof(Value));
+        memcpy(writer, &_data->_compression, sizeof(Value));
         writer += sizeof(Value);
-        memcpy(writer, &_min, sizeof(Value));
+        memcpy(writer, &_data->_min, sizeof(Value));
         writer += sizeof(Value);
-        memcpy(writer, &_max, sizeof(Value));
+        memcpy(writer, &_data->_max, sizeof(Value));
         writer += sizeof(Value);
-        memcpy(writer, &_max_processed, sizeof(Index));
+        memcpy(writer, &_data->_max_processed, sizeof(Index));
         writer += sizeof(Index);
-        memcpy(writer, &_max_unprocessed, sizeof(Index));
+        memcpy(writer, &_data->_max_unprocessed, sizeof(Index));
         writer += sizeof(Index);
-        memcpy(writer, &_processed_weight, sizeof(Value));
+        memcpy(writer, &_data->_processed_weight, sizeof(Value));
         writer += sizeof(Value);
-        memcpy(writer, &_unprocessed_weight, sizeof(Value));
+        memcpy(writer, &_data->_unprocessed_weight, sizeof(Value));
         writer += sizeof(Value);
 
-        auto size = static_cast<uint32_t>(_processed.size());
+        auto size = static_cast<uint32_t>(_data->_processed.size());
         memcpy(writer, &size, sizeof(uint32_t));
         writer += sizeof(uint32_t);
         for (int i = 0; i < size; i++) {
-            memcpy(writer, &_processed[i], sizeof(Centroid));
+            memcpy(writer, &_data->_processed[i], sizeof(Centroid));
             writer += sizeof(Centroid);
         }
 
-        size = static_cast<uint32_t>(_unprocessed.size());
+        size = static_cast<uint32_t>(_data->_unprocessed.size());
         memcpy(writer, &size, sizeof(uint32_t));
         writer += sizeof(uint32_t);
         //TODO(weixiang): may be once memcpy is enough!
         for (int i = 0; i < size; i++) {
-            memcpy(writer, &_unprocessed[i], sizeof(Centroid));
+            memcpy(writer, &_data->_unprocessed[i], sizeof(Centroid));
             writer += sizeof(Centroid);
         }
 
-        size = static_cast<uint32_t>(_cumulative.size());
+        size = static_cast<uint32_t>(_data->_cumulative.size());
         memcpy(writer, &size, sizeof(uint32_t));
         writer += sizeof(uint32_t);
         for (int i = 0; i < size; i++) {
-            memcpy(writer, &_cumulative[i], sizeof(Weight));
+            memcpy(writer, &_data->_cumulative[i], sizeof(Weight));
             writer += sizeof(Weight);
         }
         return writer - dst;
     }
 
     void unserialize(const uint8_t* type_reader) {
+        if (_data->_is_shared.load(std::memory_order_acquire)) {
+            // Deserialization replaces every field, so do not copy the old 
payload.
+            _data = std::make_shared<Data>(0, 0, 0);
+        } else {
+            _data->_processed_snapshot.reset();
+        }
         uint32_t total_length = 0;
         memcpy(&total_length, type_reader, sizeof(uint32_t));
         type_reader += sizeof(uint32_t);
-        memcpy(&_compression, type_reader, sizeof(Value));
+        memcpy(&_data->_compression, type_reader, sizeof(Value));
         type_reader += sizeof(Value);
-        memcpy(&_min, type_reader, sizeof(Value));
+        memcpy(&_data->_min, type_reader, sizeof(Value));
         type_reader += sizeof(Value);
-        memcpy(&_max, type_reader, sizeof(Value));
+        memcpy(&_data->_max, type_reader, sizeof(Value));
         type_reader += sizeof(Value);
 
-        memcpy(&_max_processed, type_reader, sizeof(Index));
+        memcpy(&_data->_max_processed, type_reader, sizeof(Index));
         type_reader += sizeof(Index);
-        memcpy(&_max_unprocessed, type_reader, sizeof(Index));
+        memcpy(&_data->_max_unprocessed, type_reader, sizeof(Index));
         type_reader += sizeof(Index);
-        memcpy(&_processed_weight, type_reader, sizeof(Value));
+        memcpy(&_data->_processed_weight, type_reader, sizeof(Value));
         type_reader += sizeof(Value);
-        memcpy(&_unprocessed_weight, type_reader, sizeof(Value));
+        memcpy(&_data->_unprocessed_weight, type_reader, sizeof(Value));
         type_reader += sizeof(Value);
 
         uint32_t size;
         memcpy(&size, type_reader, sizeof(uint32_t));
         type_reader += sizeof(uint32_t);
-        _processed.resize(size);
+        _data->_processed.resize(size);
         for (int i = 0; i < size; i++) {
-            memcpy(&_processed[i], type_reader, sizeof(Centroid));
+            memcpy(&_data->_processed[i], type_reader, sizeof(Centroid));
             type_reader += sizeof(Centroid);
         }
         memcpy(&size, type_reader, sizeof(uint32_t));
         type_reader += sizeof(uint32_t);
-        _unprocessed.resize(size);
+        _data->_unprocessed.resize(size);
         for (int i = 0; i < size; i++) {
-            memcpy(&_unprocessed[i], type_reader, sizeof(Centroid));
+            memcpy(&_data->_unprocessed[i], type_reader, sizeof(Centroid));
             type_reader += sizeof(Centroid);
         }
         memcpy(&size, type_reader, sizeof(uint32_t));
         type_reader += sizeof(uint32_t);
-        _cumulative.resize(size);
+        _data->_cumulative.resize(size);
         for (int i = 0; i < size; i++) {
-            memcpy(&_cumulative[i], type_reader, sizeof(Weight));
+            memcpy(&_data->_cumulative[i], type_reader, sizeof(Weight));
             type_reader += sizeof(Weight);
         }
     }
 
 private:
-    Value _compression;
-
-    Value _min = std::numeric_limits<Value>::max();
-
-    // min() is the smallest positive value, so use lowest() for all-negative 
input,
-    // e.g. {-3, -2, -1} must set _max to -1.
-    Value _max = std::numeric_limits<Value>::lowest();
-
-    Index _max_processed;
-
-    Index _max_unprocessed;
+    struct Data {
+        enum class CopyMode { READ_SNAPSHOT, WRITE };
+
+        Data(Value compression, Index unmerged_size, Index merged_size)
+                : _compression(compression),
+                  _max_processed(processed_size(merged_size, compression)),
+                  _max_unprocessed(unprocessed_size(unmerged_size, 
compression)) {
+            _processed.reserve(_max_processed);
+            _unprocessed.reserve(_max_unprocessed + 1);
+        }
+
+        // A detached state starts without a read cache. The source cache may 
be
+        // initialized concurrently, so neither copy nor inspect it here.
+        Data(const Data& other, CopyMode mode = CopyMode::READ_SNAPSHOT)
+                : _compression(other._compression),
+                  _min(other._min),
+                  _max(other._max),
+                  _max_processed(other._max_processed),
+                  _max_unprocessed(other._max_unprocessed),
+                  _processed_weight(other._processed_weight),
+                  _unprocessed_weight(other._unprocessed_weight) {
+            // Reserve before copying to avoid allocating and moving the 
payload
+            // twice. Read snapshots do not need spare capacity for future 
adds.
+            const bool for_write = mode == CopyMode::WRITE;
+            _processed.reserve(for_write ? 
std::max(other._processed.capacity(), _max_processed)
+                                         : other._processed.size());
+            _unprocessed.reserve(
+                    for_write ? std::max(other._unprocessed.capacity(), 
_max_unprocessed + 1)
+                              : other._unprocessed.size());
+            _cumulative.reserve(for_write ? other._cumulative.capacity()
+                                          : other._cumulative.size());
+            _processed.assign(other._processed.begin(), 
other._processed.end());
+            _unprocessed.assign(other._unprocessed.begin(), 
other._unprocessed.end());
+            _cumulative.assign(other._cumulative.begin(), 
other._cumulative.end());
+        }
+
+        Value _compression;
+        Value _min = std::numeric_limits<Value>::max();
+        Value _max = std::numeric_limits<Value>::lowest();
+        Index _max_processed;
+        Index _max_unprocessed;
+        Value _processed_weight = 0.0;
+        Value _unprocessed_weight = 0.0;
+        std::vector<Centroid> _processed;
+        std::vector<Centroid> _unprocessed;
+        std::vector<Weight> _cumulative;
+        // Once published to another handle, the payload stays immutable even
+        // when its reference count returns to one. shared_ptr::use_count() 
does
+        // not synchronize with another thread finishing reads before release.
+        std::atomic<bool> _is_shared {false};
+        std::mutex _read_mutex;
+        std::condition_variable _snapshot_cv;
+        // Guarded by _read_mutex. A non-null snapshot represents the ready 
state.
+        bool _building_snapshot = false;
+        std::unique_ptr<const TDigest> _processed_snapshot;
+    };
 
-    Value _processed_weight = 0.0;
+    std::shared_ptr<Data> _data;
 
-    Value _unprocessed_weight = 0.0;
+    explicit TDigest(const Data& data) : _data(std::make_shared<Data>(data)) {}
 
-    std::vector<Centroid> _processed;
+    void _prepare_for_write() {
+        if (_data->_is_shared.load(std::memory_order_acquire)) {
+            _data = std::make_shared<Data>(*_data, Data::CopyMode::WRITE);
+        } else {
+            _data->_processed_snapshot.reset();
+        }
+    }
 
-    std::vector<Centroid> _unprocessed;
+    const TDigest& _processed_digest() const {
+        if (!have_unprocessed() && !is_dirty()) {
+            return *this;
+        }
+        std::unique_lock<std::mutex> lock(_data->_read_mutex);
+#ifdef BE_TEST
+        if (_data->_building_snapshot) {
+            TEST_SYNC_POINT("TDigest::_processed_digest:wait_snapshot");
+        }
+#endif
+        _data->_snapshot_cv.wait(lock, [this] { return 
!_data->_building_snapshot; });
+        if (_data->_processed_snapshot) {
+            return *_data->_processed_snapshot;
+        }
+        _data->_building_snapshot = true;
+        lock.unlock();
 
-    std::vector<Weight> _cumulative;
+        // Only this reader builds a snapshot. Other readers wait with the 
mutex
+        // released, while copies and writers on other handles can still 
proceed.
+        std::unique_ptr<TDigest> snapshot;
+        try {
+#ifdef BE_TEST
+            bool fail_allocation = false;
+            
TEST_SYNC_POINT_CALLBACK("TDigest::_processed_digest:build_snapshot", 
&fail_allocation);
+            if (fail_allocation) {
+                throw std::bad_alloc();
+            }
+#endif
+            snapshot = TDigest::create_unique(*_data);

Review Comment:
   For `percentile_approx(x, .5) OVER (... ROWS BETWEEN UNBOUNDED PRECEDING AND 
CURRENT ROW)`, the analytic executor adds one row and then reads one result. 
Because this leaves the owner dirty, each add invalidates the cache and the 
next read copies and sorts the full prefixes 1, 2, ... up to 80000 centroids at 
the default compression. Previously the read processed the owner in place, so 
each step sorted only the new delta and merged the bounded processed centroids 
instead of re-sorting the full raw prefix. Please preserve an incremental path 
for an exclusive owner (while keeping shared handles immutable), or update the 
cache incrementally, and add a large alternating add/read regression.



##########
be/src/util/tdigest.h:
##########
@@ -422,222 +442,336 @@ class TDigest {
             return NAN;
         }
 
-        if (_processed.size() == 0) {
+        if (_data->_processed.size() == 0) {
             // no sorted means no data, no way to get a quantile
             return NAN;
-        } else if (_processed.size() == 1) {
+        } else if (_data->_processed.size() == 1) {
             // with one data point, all quantiles lead to Rome
 
             return _mean(0);
         }
 
         // we know that there are at least two sorted now
-        auto n = _processed.size();
+        auto n = _data->_processed.size();
 
         // if values were stored in a sorted array, index would be the offset 
we are Weighterested in
-        const auto index = q * _processed_weight;
+        const auto index = q * _data->_processed_weight;
 
-        // at the boundaries, we return _min or _max
+        // at the boundaries, we return _data->_min or _data->_max
         if (index <= _weight(0) / 2.0) {
             DCHECK_GT(_weight(0), 0);
-            return static_cast<Value>(_min + 2.0 * index / _weight(0) * 
(_mean(0) - _min));
+            return static_cast<Value>(_data->_min +
+                                      2.0 * index / _weight(0) * (_mean(0) - 
_data->_min));
         }
 
-        auto iter = std::lower_bound(_cumulative.cbegin(), _cumulative.cend(), 
index);
+        auto iter = std::lower_bound(_data->_cumulative.cbegin(), 
_data->_cumulative.cend(), index);
 
-        if (iter != _cumulative.cend() && iter != _cumulative.cbegin() &&
-            iter + 1 != _cumulative.cend()) {
-            auto i = std::distance(_cumulative.cbegin(), iter);
+        if (iter != _data->_cumulative.cend() && iter != 
_data->_cumulative.cbegin() &&
+            iter + 1 != _data->_cumulative.cend()) {
+            auto i = std::distance(_data->_cumulative.cbegin(), iter);
             auto z1 = index - *(iter - 1);
             auto z2 = *(iter)-index;
             // VLOG_CRITICAL << "z2 " << z2 << " index " << index << " z1 " << 
z1;
             return _weighted_average(_mean(i - 1), z2, _mean(i), z1);
         }
 
-        DCHECK_LE(index, _processed_weight);
-        DCHECK_GE(index, _processed_weight - _weight(n - 1) / 2.0);
+        DCHECK_LE(index, _data->_processed_weight);
+        DCHECK_GE(index, _data->_processed_weight - _weight(n - 1) / 2.0);
 
-        auto z1 = static_cast<Value>(index - _processed_weight - _weight(n - 
1) / 2.0);
+        auto z1 = static_cast<Value>(index - _data->_processed_weight - 
_weight(n - 1) / 2.0);
         auto z2 = static_cast<Value>(_weight(n - 1) / 2 - z1);
-        return _weighted_average(_mean(n - 1), z1, _max, z2);
+        return _weighted_average(_mean(n - 1), z1, _data->_max, z2);
     }
 
-    Value compression() const { return _compression; }
+    Value compression() const { return _data->_compression; }
 
     void add(Value x) { add(x, 1); }
 
-    void compress() { _process(); }
+    void compress() {
+        if (total_size() != 0) {
+            _prepare_for_write();
+            _process();
+        }
+    }
 
     // add a single centroid to the unprocessed vector, processing previously 
unprocessed sorted if our limit has
     // been reached.
     bool add(Value x, Weight w) {
         if (std::isnan(x)) {
             return false;
         }
-        _unprocessed.emplace_back(x, w);
-        _unprocessed_weight += w;
+        _prepare_for_write();
+        _data->_unprocessed.emplace_back(x, w);
+        _data->_unprocessed_weight += w;
         _process_if_necessary();
         return true;
     }
 
     void add(std::vector<Centroid>::const_iterator iter,
              std::vector<Centroid>::const_iterator end) {
+        const std::vector<Centroid> centroids(iter, end);
+        _prepare_for_write();
+        iter = centroids.cbegin();
+        end = centroids.cend();
         while (iter != end) {
             const size_t diff = std::distance(iter, end);
-            const size_t room = _max_unprocessed - _unprocessed.size();
+            const size_t room = _data->_max_unprocessed - 
_data->_unprocessed.size();
             auto mid = iter + std::min(diff, room);
             while (iter != mid) {
-                _unprocessed.push_back(*(iter++));
+                _data->_unprocessed_weight += iter->weight();
+                _data->_unprocessed.push_back(*(iter++));
             }
-            if (_unprocessed.size() >= _max_unprocessed) {
+            if (_data->_unprocessed.size() >= _data->_max_unprocessed) {
                 _process();
             }
         }
     }
 
-    uint32_t serialized_size() {
+    uint32_t serialized_size() const {
         return static_cast<uint32_t>(sizeof(uint32_t) + sizeof(Value) * 5 + 
sizeof(Index) * 2 +
-                                     sizeof(uint32_t) * 3 + _processed.size() 
* sizeof(Centroid) +
-                                     _unprocessed.size() * sizeof(Centroid) +
-                                     _cumulative.size() * sizeof(Weight));
+                                     sizeof(uint32_t) * 3 +
+                                     _data->_processed.size() * 
sizeof(Centroid) +
+                                     _data->_unprocessed.size() * 
sizeof(Centroid) +
+                                     _data->_cumulative.size() * 
sizeof(Weight));
     }
 
-    size_t serialize(uint8_t* writer) {
+    size_t serialize(uint8_t* writer) const {
         uint8_t* dst = writer;
         uint32_t total_size = serialized_size();
         memcpy(writer, &total_size, sizeof(uint32_t));
         writer += sizeof(uint32_t);
-        memcpy(writer, &_compression, sizeof(Value));
+        memcpy(writer, &_data->_compression, sizeof(Value));
         writer += sizeof(Value);
-        memcpy(writer, &_min, sizeof(Value));
+        memcpy(writer, &_data->_min, sizeof(Value));
         writer += sizeof(Value);
-        memcpy(writer, &_max, sizeof(Value));
+        memcpy(writer, &_data->_max, sizeof(Value));
         writer += sizeof(Value);
-        memcpy(writer, &_max_processed, sizeof(Index));
+        memcpy(writer, &_data->_max_processed, sizeof(Index));
         writer += sizeof(Index);
-        memcpy(writer, &_max_unprocessed, sizeof(Index));
+        memcpy(writer, &_data->_max_unprocessed, sizeof(Index));
         writer += sizeof(Index);
-        memcpy(writer, &_processed_weight, sizeof(Value));
+        memcpy(writer, &_data->_processed_weight, sizeof(Value));
         writer += sizeof(Value);
-        memcpy(writer, &_unprocessed_weight, sizeof(Value));
+        memcpy(writer, &_data->_unprocessed_weight, sizeof(Value));
         writer += sizeof(Value);
 
-        auto size = static_cast<uint32_t>(_processed.size());
+        auto size = static_cast<uint32_t>(_data->_processed.size());
         memcpy(writer, &size, sizeof(uint32_t));
         writer += sizeof(uint32_t);
         for (int i = 0; i < size; i++) {
-            memcpy(writer, &_processed[i], sizeof(Centroid));
+            memcpy(writer, &_data->_processed[i], sizeof(Centroid));
             writer += sizeof(Centroid);
         }
 
-        size = static_cast<uint32_t>(_unprocessed.size());
+        size = static_cast<uint32_t>(_data->_unprocessed.size());
         memcpy(writer, &size, sizeof(uint32_t));
         writer += sizeof(uint32_t);
         //TODO(weixiang): may be once memcpy is enough!
         for (int i = 0; i < size; i++) {
-            memcpy(writer, &_unprocessed[i], sizeof(Centroid));
+            memcpy(writer, &_data->_unprocessed[i], sizeof(Centroid));
             writer += sizeof(Centroid);
         }
 
-        size = static_cast<uint32_t>(_cumulative.size());
+        size = static_cast<uint32_t>(_data->_cumulative.size());
         memcpy(writer, &size, sizeof(uint32_t));
         writer += sizeof(uint32_t);
         for (int i = 0; i < size; i++) {
-            memcpy(writer, &_cumulative[i], sizeof(Weight));
+            memcpy(writer, &_data->_cumulative[i], sizeof(Weight));
             writer += sizeof(Weight);
         }
         return writer - dst;
     }
 
     void unserialize(const uint8_t* type_reader) {
+        if (_data->_is_shared.load(std::memory_order_acquire)) {
+            // Deserialization replaces every field, so do not copy the old 
payload.
+            _data = std::make_shared<Data>(0, 0, 0);
+        } else {
+            _data->_processed_snapshot.reset();
+        }
         uint32_t total_length = 0;
         memcpy(&total_length, type_reader, sizeof(uint32_t));
         type_reader += sizeof(uint32_t);
-        memcpy(&_compression, type_reader, sizeof(Value));
+        memcpy(&_data->_compression, type_reader, sizeof(Value));
         type_reader += sizeof(Value);
-        memcpy(&_min, type_reader, sizeof(Value));
+        memcpy(&_data->_min, type_reader, sizeof(Value));
         type_reader += sizeof(Value);
-        memcpy(&_max, type_reader, sizeof(Value));
+        memcpy(&_data->_max, type_reader, sizeof(Value));
         type_reader += sizeof(Value);
 
-        memcpy(&_max_processed, type_reader, sizeof(Index));
+        memcpy(&_data->_max_processed, type_reader, sizeof(Index));
         type_reader += sizeof(Index);
-        memcpy(&_max_unprocessed, type_reader, sizeof(Index));
+        memcpy(&_data->_max_unprocessed, type_reader, sizeof(Index));
         type_reader += sizeof(Index);
-        memcpy(&_processed_weight, type_reader, sizeof(Value));
+        memcpy(&_data->_processed_weight, type_reader, sizeof(Value));
         type_reader += sizeof(Value);
-        memcpy(&_unprocessed_weight, type_reader, sizeof(Value));
+        memcpy(&_data->_unprocessed_weight, type_reader, sizeof(Value));
         type_reader += sizeof(Value);
 
         uint32_t size;
         memcpy(&size, type_reader, sizeof(uint32_t));
         type_reader += sizeof(uint32_t);
-        _processed.resize(size);
+        _data->_processed.resize(size);
         for (int i = 0; i < size; i++) {
-            memcpy(&_processed[i], type_reader, sizeof(Centroid));
+            memcpy(&_data->_processed[i], type_reader, sizeof(Centroid));
             type_reader += sizeof(Centroid);
         }
         memcpy(&size, type_reader, sizeof(uint32_t));
         type_reader += sizeof(uint32_t);
-        _unprocessed.resize(size);
+        _data->_unprocessed.resize(size);
         for (int i = 0; i < size; i++) {
-            memcpy(&_unprocessed[i], type_reader, sizeof(Centroid));
+            memcpy(&_data->_unprocessed[i], type_reader, sizeof(Centroid));
             type_reader += sizeof(Centroid);
         }
         memcpy(&size, type_reader, sizeof(uint32_t));
         type_reader += sizeof(uint32_t);
-        _cumulative.resize(size);
+        _data->_cumulative.resize(size);
         for (int i = 0; i < size; i++) {
-            memcpy(&_cumulative[i], type_reader, sizeof(Weight));
+            memcpy(&_data->_cumulative[i], type_reader, sizeof(Weight));
             type_reader += sizeof(Weight);
         }
     }
 
 private:
-    Value _compression;
-
-    Value _min = std::numeric_limits<Value>::max();
-
-    // min() is the smallest positive value, so use lowest() for all-negative 
input,
-    // e.g. {-3, -2, -1} must set _max to -1.
-    Value _max = std::numeric_limits<Value>::lowest();
-
-    Index _max_processed;
-
-    Index _max_unprocessed;
+    struct Data {
+        enum class CopyMode { READ_SNAPSHOT, WRITE };
+
+        Data(Value compression, Index unmerged_size, Index merged_size)
+                : _compression(compression),
+                  _max_processed(processed_size(merged_size, compression)),
+                  _max_unprocessed(unprocessed_size(unmerged_size, 
compression)) {
+            _processed.reserve(_max_processed);
+            _unprocessed.reserve(_max_unprocessed + 1);
+        }
+
+        // A detached state starts without a read cache. The source cache may 
be
+        // initialized concurrently, so neither copy nor inspect it here.
+        Data(const Data& other, CopyMode mode = CopyMode::READ_SNAPSHOT)
+                : _compression(other._compression),
+                  _min(other._min),
+                  _max(other._max),
+                  _max_processed(other._max_processed),
+                  _max_unprocessed(other._max_unprocessed),
+                  _processed_weight(other._processed_weight),
+                  _unprocessed_weight(other._unprocessed_weight) {
+            // Reserve before copying to avoid allocating and moving the 
payload
+            // twice. Read snapshots do not need spare capacity for future 
adds.
+            const bool for_write = mode == CopyMode::WRITE;
+            _processed.reserve(for_write ? 
std::max(other._processed.capacity(), _max_processed)
+                                         : other._processed.size());
+            _unprocessed.reserve(
+                    for_write ? std::max(other._unprocessed.capacity(), 
_max_unprocessed + 1)
+                              : other._unprocessed.size());
+            _cumulative.reserve(for_write ? other._cumulative.capacity()
+                                          : other._cumulative.size());
+            _processed.assign(other._processed.begin(), 
other._processed.end());
+            _unprocessed.assign(other._unprocessed.begin(), 
other._unprocessed.end());
+            _cumulative.assign(other._cumulative.begin(), 
other._cumulative.end());
+        }
+
+        Value _compression;
+        Value _min = std::numeric_limits<Value>::max();
+        Value _max = std::numeric_limits<Value>::lowest();
+        Index _max_processed;
+        Index _max_unprocessed;
+        Value _processed_weight = 0.0;
+        Value _unprocessed_weight = 0.0;
+        std::vector<Centroid> _processed;
+        std::vector<Centroid> _unprocessed;
+        std::vector<Weight> _cumulative;
+        // Once published to another handle, the payload stays immutable even
+        // when its reference count returns to one. shared_ptr::use_count() 
does
+        // not synchronize with another thread finishing reads before release.
+        std::atomic<bool> _is_shared {false};
+        std::mutex _read_mutex;
+        std::condition_variable _snapshot_cv;
+        // Guarded by _read_mutex. A non-null snapshot represents the ready 
state.
+        bool _building_snapshot = false;
+        std::unique_ptr<const TDigest> _processed_snapshot;
+    };
 
-    Value _processed_weight = 0.0;
+    std::shared_ptr<Data> _data;
 
-    Value _unprocessed_weight = 0.0;
+    explicit TDigest(const Data& data) : _data(std::make_shared<Data>(data)) {}
 
-    std::vector<Centroid> _processed;
+    void _prepare_for_write() {
+        if (_data->_is_shared.load(std::memory_order_acquire)) {
+            _data = std::make_shared<Data>(*_data, Data::CopyMode::WRITE);
+        } else {
+            _data->_processed_snapshot.reset();
+        }
+    }
 
-    std::vector<Centroid> _unprocessed;
+    const TDigest& _processed_digest() const {
+        if (!have_unprocessed() && !is_dirty()) {
+            return *this;
+        }
+        std::unique_lock<std::mutex> lock(_data->_read_mutex);
+#ifdef BE_TEST
+        if (_data->_building_snapshot) {
+            TEST_SYNC_POINT("TDigest::_processed_digest:wait_snapshot");
+        }
+#endif
+        _data->_snapshot_cv.wait(lock, [this] { return 
!_data->_building_snapshot; });
+        if (_data->_processed_snapshot) {
+            return *_data->_processed_snapshot;
+        }
+        _data->_building_snapshot = true;
+        lock.unlock();
 
-    std::vector<Weight> _cumulative;
+        // Only this reader builds a snapshot. Other readers wait with the 
mutex
+        // released, while copies and writers on other handles can still 
proceed.
+        std::unique_ptr<TDigest> snapshot;
+        try {
+#ifdef BE_TEST
+            bool fail_allocation = false;
+            
TEST_SYNC_POINT_CALLBACK("TDigest::_processed_digest:build_snapshot", 
&fail_allocation);
+            if (fail_allocation) {
+                throw std::bad_alloc();
+            }
+#endif
+            snapshot = TDigest::create_unique(*_data);
+            snapshot->compress();

Review Comment:
   After this compression, the snapshot's `_unprocessed` vector is empty but 
still retains the capacity copied in `Data(other)`. At compression 10000 that 
can be 80000 centroids (about 640 KB) per finalized state, so 1000 groups can 
retain about 610 MiB extra until teardown, and 
`ColumnQuantileState::allocated_bytes()` does not see it. Please release this 
empty buffer before publishing the cache and add a post-compression capacity 
test.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to