Skip to content

Commit 2278e75

Browse files
authored
fix builtin sample transport opt (#37)
1 parent 908fcd7 commit 2278e75

1 file changed

Lines changed: 39 additions & 13 deletions

File tree

statshouse.hpp

Lines changed: 39 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -1452,6 +1452,13 @@ class Registry {
14521452
, unique{nullptr}
14531453
, increment{false} {
14541454
}
1455+
multivalue_view(size_t values_count, const double *values, double count)
1456+
: count{count}
1457+
, values_count{values_count}
1458+
, values{values}
1459+
, unique{nullptr}
1460+
, increment{false} {
1461+
}
14551462
multivalue_view(const double (&values)[2])
14561463
: count{0}
14571464
, values_count{2}
@@ -1466,6 +1473,13 @@ class Registry {
14661473
, unique{unique}
14671474
, increment{false} {
14681475
}
1476+
multivalue_view(size_t values_count, const uint64_t *unique, double count)
1477+
: count{count}
1478+
, values_count{values_count}
1479+
, values{nullptr}
1480+
, unique{unique}
1481+
, increment{false} {
1482+
}
14691483
bool empty() const noexcept {
14701484
return count == 0 && !values_count;
14711485
}
@@ -1492,6 +1506,7 @@ class Registry {
14921506
? src.values[0] + src.values[1]
14931507
: src.values[src.values_count - 1]);
14941508
size = 1;
1509+
weighted_count = 1;
14951510
} else {
14961511
values.back() = src.increment
14971512
? values.back() + src.values[1]
@@ -1500,20 +1515,23 @@ class Registry {
15001515
src.values += src.values_count;
15011516
src.values_count = 0;
15021517
} else {
1503-
write(values, src.values, src.values_count, r, b);
1518+
write(values, src.values, src.values_count, src.count, r, b);
15041519
}
1520+
src.count = 0;
15051521
} else if (src.unique) {
1506-
write(unique, src.unique, src.values_count, r, b);
1522+
write(unique, src.unique, src.values_count, src.count, r, b);
1523+
src.count = 0;
15071524
} else {
15081525
count += src.count;
15091526
src.count = 0;
15101527
}
15111528
}
15121529
template <typename T>
1513-
void write(std::vector<T> &dst, const T * &src, size_t &src_size, Registry *r, bucket *b) {
1530+
void write(std::vector<T> &dst, const T * &src, size_t &src_size, double src_count, Registry *r, bucket *b) {
15141531
if (!src_size) {
15151532
return;
15161533
}
1534+
const size_t source_size = src_size;
15171535
size_t i = 0;
15181536
size_t limit = r->max_bucket_size.load(std::memory_order_relaxed);
15191537
dst.reserve(limit); // minimize the number of memory allocations
@@ -1533,12 +1551,18 @@ class Registry {
15331551
}
15341552
src += i;
15351553
src_size -= i;
1554+
if (src_count != 0) {
1555+
weighted_count += src_count;
1556+
} else {
1557+
weighted_count += static_cast<double>(source_size);
1558+
}
15361559
}
15371560
bool empty() const noexcept {
15381561
return count == 0 && values.empty() && unique.empty();
15391562
}
15401563
double count{};
15411564
size_t size{};
1565+
double weighted_count{};
15421566
std::vector<double> values;
15431567
std::vector<uint64_t> unique;
15441568
};
@@ -1651,12 +1675,12 @@ class Registry {
16511675
}
16521676
count *= m;
16531677
registry->log_values(ptr->key, values, values_count, count, timestamp);
1654-
auto sampling = count != 0 && count != values_count;
1678+
auto sampling = count != 0 && count != values_count && m == 1.0;
16551679
if (sampling || timestamp || values_count >= registry->max_bucket_size.load(std::memory_order_relaxed)) {
16561680
std::lock_guard<std::mutex> transport_lock{registry->transport_mu};
16571681
return ptr->key.write_values(values, values_count, count, timestamp);
16581682
}
1659-
multivalue_view value{values_count, values};
1683+
multivalue_view value{values_count, values, count};
16601684
return registry->update_multivalue_by_ref(ptr, value);
16611685
}
16621686
};
@@ -1676,12 +1700,12 @@ class Registry {
16761700
}
16771701
count *= m;
16781702
registry->log_unique(ptr->key, values, values_count, count, timestamp);
1679-
auto sampling = count != 0 && count != values_count;
1703+
auto sampling = count != 0 && count != values_count && m == 1.0;
16801704
if (sampling || timestamp || values_count >= registry->max_bucket_size.load(std::memory_order_relaxed)) {
16811705
std::lock_guard<std::mutex> transport_lock{registry->transport_mu};
16821706
return ptr->key.write_unique(values, values_count, count, timestamp);
16831707
}
1684-
multivalue_view value{values_count, values};
1708+
multivalue_view value{values_count, values, count};
16851709
return registry->update_multivalue_by_ref(ptr, value);
16861710
}
16871711
};
@@ -1767,12 +1791,12 @@ class Registry {
17671791
}
17681792
count *= m;
17691793
registry->log_values(key, values, values_count, count, timestamp);
1770-
auto sampling = count != 0 && count != values_count;
1794+
auto sampling = count != 0 && count != values_count && m == 1.0;
17711795
if (sampling || timestamp || values_count >= registry->max_bucket_size.load(std::memory_order_relaxed)) {
17721796
std::lock_guard<std::mutex> transport_lock{registry->transport_mu};
17731797
return key.write_values(values, values_count, count, timestamp);
17741798
}
1775-
multivalue_view value{values_count, values};
1799+
multivalue_view value{values_count, values, count};
17761800
return registry->update_multivalue_by_key(key, value);
17771801
}
17781802
bool write_unique(uint64_t value, uint32_t timestamp = 0) const {
@@ -1785,12 +1809,12 @@ class Registry {
17851809
}
17861810
count *= m;
17871811
registry->log_unique(key, values, values_count, count, timestamp);
1788-
auto sampling = count != 0 && count != values_count;
1812+
auto sampling = count != 0 && count != values_count && m == 1.0;
17891813
if (sampling || timestamp || values_count >= registry->max_bucket_size.load(std::memory_order_relaxed)) {
17901814
std::lock_guard<std::mutex> transport_lock{registry->transport_mu};
17911815
return key.write_unique(values, values_count, count, timestamp);
17921816
}
1793-
multivalue_view value{values_count, values};
1817+
multivalue_view value{values_count, values, count};
17941818
return registry->update_multivalue_by_key(key, value);
17951819
}
17961820
private:
@@ -1932,6 +1956,7 @@ class Registry {
19321956
if (value.increment && !(*ptr->queue_ptr)->value.values.empty()) {
19331957
ptr->value.values.push_back((*ptr->queue_ptr)->value.values.back());
19341958
ptr->value.size = 1;
1959+
ptr->value.weighted_count = 1;
19351960
}
19361961
ptr->value.write(value, this, ptr.get());
19371962
// assert(value.empty())
@@ -2014,9 +2039,9 @@ class Registry {
20142039
std::lock_guard<std::mutex> transport_lock{transport_mu};
20152040
auto &v = ptr->value;
20162041
if (!v.values.empty()) {
2017-
ptr->key.write_values(v.values.data(), v.values.size(), v.size, ptr->timestamp);
2042+
ptr->key.write_values(v.values.data(), v.values.size(), v.weighted_count, ptr->timestamp);
20182043
} else if (!v.unique.empty()) {
2019-
ptr->key.write_unique(v.unique.data(), v.unique.size(), v.size, ptr->timestamp);
2044+
ptr->key.write_unique(v.unique.data(), v.unique.size(), v.weighted_count, ptr->timestamp);
20202045
} else if (v.count != 0) {
20212046
ptr->key.write_count(v.count, ptr->timestamp);
20222047
}
@@ -2057,6 +2082,7 @@ class Registry {
20572082
ptr->waterlevel = 0;
20582083
ptr->value.count = 0;
20592084
ptr->value.size = 0;
2085+
ptr->value.weighted_count = 0;
20602086
ptr->value.values.clear();
20612087
ptr->value.unique.clear();
20622088
ptr->queue_ptr = nullptr;

0 commit comments

Comments
 (0)