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]