Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
52 changes: 39 additions & 13 deletions statshouse.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -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}
Expand All @@ -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;
}
Expand All @@ -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]
Expand All @@ -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 <typename T>
void write(std::vector<T> &dst, const T * &src, size_t &src_size, Registry *r, bucket *b) {
void write(std::vector<T> &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
Expand All @@ -1533,12 +1551,18 @@ class Registry {
}
src += i;
src_size -= i;
if (src_count != 0) {
weighted_count += src_count;
} else {
weighted_count += static_cast<double>(source_size);
}
}
bool empty() const noexcept {
return count == 0 && values.empty() && unique.empty();
}
double count{};
size_t size{};
double weighted_count{};
std::vector<double> values;
std::vector<uint64_t> unique;
};
Expand Down Expand Up @@ -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<std::mutex> 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);
}
};
Expand All @@ -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<std::mutex> 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);
}
};
Expand Down Expand Up @@ -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<std::mutex> 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 {
Expand All @@ -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<std::mutex> 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:
Expand Down Expand Up @@ -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())
Expand Down Expand Up @@ -2014,9 +2039,9 @@ class Registry {
std::lock_guard<std::mutex> 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);
}
Expand Down Expand Up @@ -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;
Expand Down