diff --git a/statshouse.hpp b/statshouse.hpp index 64065f7..145a3d7 100644 --- a/statshouse.hpp +++ b/statshouse.hpp @@ -1452,6 +1452,13 @@ class Registry { , unique{nullptr} , increment{false} { } + multivalue_view(size_t values_count, const double *values, double count) + : count{count} + , values_count{values_count} + , values{values} + , unique{nullptr} + , increment{false} { + } multivalue_view(const double (&values)[2]) : count{0} , values_count{2} @@ -1466,6 +1473,13 @@ class Registry { , unique{unique} , increment{false} { } + multivalue_view(size_t values_count, const uint64_t *unique, double count) + : count{count} + , values_count{values_count} + , values{nullptr} + , unique{unique} + , increment{false} { + } bool empty() const noexcept { return count == 0 && !values_count; } @@ -1492,6 +1506,7 @@ class Registry { ? src.values[0] + src.values[1] : src.values[src.values_count - 1]); size = 1; + weighted_count = 1; } else { values.back() = src.increment ? values.back() + src.values[1] @@ -1500,20 +1515,23 @@ class Registry { src.values += src.values_count; src.values_count = 0; } else { - write(values, src.values, src.values_count, r, b); + write(values, src.values, src.values_count, src.count, r, b); } + src.count = 0; } else if (src.unique) { - write(unique, src.unique, src.values_count, r, b); + write(unique, src.unique, src.values_count, src.count, r, b); + src.count = 0; } else { count += src.count; src.count = 0; } } template - void write(std::vector &dst, const T * &src, size_t &src_size, Registry *r, bucket *b) { + void write(std::vector &dst, const T * &src, size_t &src_size, double src_count, Registry *r, bucket *b) { if (!src_size) { return; } + const size_t source_size = src_size; size_t i = 0; size_t limit = r->max_bucket_size.load(std::memory_order_relaxed); dst.reserve(limit); // minimize the number of memory allocations @@ -1533,12 +1551,18 @@ class Registry { } src += i; src_size -= i; + if (src_count != 0) { + weighted_count += src_count; + } else { + weighted_count += static_cast(source_size); + } } bool empty() const noexcept { return count == 0 && values.empty() && unique.empty(); } double count{}; size_t size{}; + double weighted_count{}; std::vector values; std::vector unique; }; @@ -1651,12 +1675,12 @@ class Registry { } count *= m; registry->log_values(ptr->key, values, values_count, count, timestamp); - auto sampling = count != 0 && count != values_count; + auto sampling = count != 0 && count != values_count && m == 1.0; if (sampling || timestamp || values_count >= registry->max_bucket_size.load(std::memory_order_relaxed)) { std::lock_guard transport_lock{registry->transport_mu}; return ptr->key.write_values(values, values_count, count, timestamp); } - multivalue_view value{values_count, values}; + multivalue_view value{values_count, values, count}; return registry->update_multivalue_by_ref(ptr, value); } }; @@ -1676,12 +1700,12 @@ class Registry { } count *= m; registry->log_unique(ptr->key, values, values_count, count, timestamp); - auto sampling = count != 0 && count != values_count; + auto sampling = count != 0 && count != values_count && m == 1.0; if (sampling || timestamp || values_count >= registry->max_bucket_size.load(std::memory_order_relaxed)) { std::lock_guard transport_lock{registry->transport_mu}; return ptr->key.write_unique(values, values_count, count, timestamp); } - multivalue_view value{values_count, values}; + multivalue_view value{values_count, values, count}; return registry->update_multivalue_by_ref(ptr, value); } }; @@ -1767,12 +1791,12 @@ class Registry { } count *= m; registry->log_values(key, values, values_count, count, timestamp); - auto sampling = count != 0 && count != values_count; + auto sampling = count != 0 && count != values_count && m == 1.0; if (sampling || timestamp || values_count >= registry->max_bucket_size.load(std::memory_order_relaxed)) { std::lock_guard transport_lock{registry->transport_mu}; return key.write_values(values, values_count, count, timestamp); } - multivalue_view value{values_count, values}; + multivalue_view value{values_count, values, count}; return registry->update_multivalue_by_key(key, value); } bool write_unique(uint64_t value, uint32_t timestamp = 0) const { @@ -1785,12 +1809,12 @@ class Registry { } count *= m; registry->log_unique(key, values, values_count, count, timestamp); - auto sampling = count != 0 && count != values_count; + auto sampling = count != 0 && count != values_count && m == 1.0; if (sampling || timestamp || values_count >= registry->max_bucket_size.load(std::memory_order_relaxed)) { std::lock_guard transport_lock{registry->transport_mu}; return key.write_unique(values, values_count, count, timestamp); } - multivalue_view value{values_count, values}; + multivalue_view value{values_count, values, count}; return registry->update_multivalue_by_key(key, value); } private: @@ -1932,6 +1956,7 @@ class Registry { if (value.increment && !(*ptr->queue_ptr)->value.values.empty()) { ptr->value.values.push_back((*ptr->queue_ptr)->value.values.back()); ptr->value.size = 1; + ptr->value.weighted_count = 1; } ptr->value.write(value, this, ptr.get()); // assert(value.empty()) @@ -2014,9 +2039,9 @@ class Registry { std::lock_guard transport_lock{transport_mu}; auto &v = ptr->value; if (!v.values.empty()) { - ptr->key.write_values(v.values.data(), v.values.size(), v.size, ptr->timestamp); + ptr->key.write_values(v.values.data(), v.values.size(), v.weighted_count, ptr->timestamp); } else if (!v.unique.empty()) { - ptr->key.write_unique(v.unique.data(), v.unique.size(), v.size, ptr->timestamp); + ptr->key.write_unique(v.unique.data(), v.unique.size(), v.weighted_count, ptr->timestamp); } else if (v.count != 0) { ptr->key.write_count(v.count, ptr->timestamp); } @@ -2057,6 +2082,7 @@ class Registry { ptr->waterlevel = 0; ptr->value.count = 0; ptr->value.size = 0; + ptr->value.weighted_count = 0; ptr->value.values.clear(); ptr->value.unique.clear(); ptr->queue_ptr = nullptr;