Skip to content
Merged
Show file tree
Hide file tree
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
8 changes: 4 additions & 4 deletions db/art/clf_model.cc
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@ void ClfModel::write_debug_dataset() {
std::vector<std::string> header;
header.emplace_back("Level");
for (int i = 0; i < 20; i ++) {
header.emplace_back("Range_" + std::to_string(i));
header.emplace_back("Rate_" + std::to_string(i));
header.emplace_back("Hotness_" + std::to_string(i));
}
header.emplace_back("Target");
Expand Down Expand Up @@ -65,8 +65,8 @@ void ClfModel::write_debug_dataset() {
std::shuffle(ids.begin(), ids.end(), std::default_random_engine(seed));
values.emplace_back(std::to_string(level));
for (int j = 0; j < 20; j ++) {
values.emplace_back(std::to_string(ids[j]));
values.emplace_back(std::to_string(uint32_t(SIGNIFICANT_DIGITS_FACTOR * hotness_map[ids[j]])));
values.emplace_back(std::to_string(uint32_t(ids[j] * 0.005 * RATE_SIGNIFICANT_DIGITS_FACTOR)));
values.emplace_back(std::to_string(uint32_t(HOTNESS_SIGNIFICANT_DIGITS_FACTOR * hotness_map[ids[j]])));
}
values.emplace_back(std::to_string(target)); // Target column
values.emplace_back(std::to_string(count)); // Count column
Expand Down Expand Up @@ -94,7 +94,7 @@ void ClfModel::write_real_dataset(std::vector<std::vector<uint32_t>>& datas, std
uint16_t ranges_num = (feature_num_ - 1) / 2;
header.emplace_back("Level");
for (int i = 0; i < ranges_num; i ++) {
header.emplace_back("Range_" + std::to_string(i));
header.emplace_back("Rate_" + std::to_string(i));
header.emplace_back("Hotness_" + std::to_string(i));
}
// remind that targeted class is in csv Target column
Expand Down
44 changes: 28 additions & 16 deletions db/art/clf_model.h
Original file line number Diff line number Diff line change
Expand Up @@ -13,34 +13,41 @@
// supposed that considering r key range
// every key range have id and hotness ( see heat_buckets )
// so data point features format :
// LSM-Tree level, Key Range 1 id, Key Range 1 hotness, Key Range 2 id, Key Range 2 hotness, ..., Key Range r id, Key Range r hotness
// we also need to append best units num and visit count to every row
// LSM-Tree level, Key Range 1 rate, Key Range 1 hotness, Key Range 2 rate, Key Range 2 hotness, ..., Key Range r rate, Key Range r hotness
// we also need to append best units num (from solving programming problem) and visit count to every row
// so in data csv, one row would be like:
// LSM-Tree level, key range 1 id, key range 1 hotness, ..., best units num (for the segment), visit count to this segment (in last long peroid)
// LSM-Tree level, key range 1 rate, key range 1 hotness, ..., best units num (for the segment), visit count to this segment (in last long peroid)
// for example, assume that segment 1 can be devide into key range 3 (50% keys), key range 4 (30% keys), key range 6 (20% keys)
// sort these key ranges by their keys rate (e.g. key range 3 rate = 50%) and set feature num to 7
// data format can be like:
// 5 (segment 1 level in LSM-Tree), 50% (key range 3 rate), 5234 (key range 3 hotness), 30% (key range 4 rate), 2222 (key range 4 hotness), 20% (key range 6 rate), 11111 (key range 6 hotness)
// remind that heat_buckets recorded hotness is double type,
// we use uint32_t(uint32_t(SIGNIFICANT_DIGITS * hotness)) to closely estimate its hotness value
// because data feature only accept uint32_t type,
// we use uint32_t(HOTNESS_SIGNIFICANT_DIGITS_FACTOR * hotness) to closely estimate its hotness value
// also use uint32_t(RATE_SIGNIFICANT_DIGITS_FACTOR * rate) to closely estimate its rate in this segment

namespace ROCKSDB_NAMESPACE {

struct RangeHeatPair;
struct RangeRatePair;
class ClfModel;
bool RangeHeatPairLesserComparor(const RangeHeatPair& pair_1, const RangeHeatPair& pair_2);
bool RangeHeatPairGreaterComparor(const RangeHeatPair& pair_1, const RangeHeatPair& pair_2);

struct RangeHeatPair {
bool RangeRatePairLessorComparor(const RangeRatePair& pair_1, const RangeRatePair& pair_2);
bool RangeRatePairGreatorComparor(const RangeRatePair& pair_1, const RangeRatePair& pair_2);

struct RangeRatePair {
uint32_t range_id;
double hotness_val;
RangeHeatPair(const uint32_t& id, const double& hotness) {
range_id = id; hotness_val = hotness;
double rate_in_segment;
RangeRatePair(const uint32_t& id, const double& rate) {
range_id = id; rate_in_segment = rate;
}
};

bool RangeHeatPairLesserComparor(const RangeHeatPair& pair_1, const RangeHeatPair& pair_2) {
return pair_1.hotness_val < pair_2.hotness_val;
bool RangeRatePairLessorComparor(const RangeRatePair& pair_1, const RangeRatePair& pair_2) {
return pair_1.rate_in_segment < pair_2.rate_in_segment;
}

bool RangeHeatPairGreaterComparor(const RangeHeatPair& pair_1, const RangeHeatPair& pair_2) {
return pair_1.hotness_val > pair_2.hotness_val;
bool RangeRatePairGreatorComparor(const RangeRatePair& pair_1, const RangeRatePair& pair_2) {
return pair_1.rate_in_segment > pair_2.rate_in_segment;
}

class ClfModel {
Expand Down Expand Up @@ -68,13 +75,18 @@ class ClfModel {
// make ready for training, only need init feature_nums_ now
// when first call ClfModel, we need to use current segments information to init features_num_
// we can calcuate feature nums for every segment,
// feature num = level feature num (1) + 2 * num of key ranges segment covers
// feature num = level feature num (1) + 2 * num of key ranges
// we set features_num_ to largest feature num
void make_ready(std::vector<uint16_t>& features_nums) {
if (features_nums.empty()) {
feature_num_ = 41; // debug feature num, see ../lgb_server files
} else {
// we may limit feature_num_ because of the socket transmit size limit is 1024 bytes
// so feature_num_ may be limit to at most about 3 * 30 + 1 = 91
feature_num_ = *max_element(features_nums.begin(), features_nums.end());
if (feature_num_ > MAX_FEATURES_NUM) {
feature_num_ = MAX_FEATURES_NUM;
}
}

// std::cout << "[DEBUG] ClfModel ready, feature_num_: " << feature_num_ << std::endl;
Expand Down
54 changes: 31 additions & 23 deletions db/art/filter_cache.cc
Original file line number Diff line number Diff line change
Expand Up @@ -138,6 +138,7 @@ bool FilterCacheManager::make_heat_buckets_ready(const std::string& key,
}
heat_buckets_.sample(key, segments_infos);
}
return heat_buckets_.is_ready();
}

void FilterCacheManager::hit_heat_buckets(const std::string& key) {
Expand Down Expand Up @@ -176,8 +177,14 @@ void FilterCacheManager::hit_heat_buckets(const std::string& key) {
}
}

bool FilterCacheManager::make_clf_model_ready(std::vector<uint16_t>& features_nums) {
clf_model_.make_ready(features_nums);
return clf_model_.is_ready();
}

bool FilterCacheManager::check_key(const uint32_t& segment_id, const std::string& key) {
hit_count_recorder(segment_id); // one get opt will cause query to many segments.
// move hit_count_recorder to a background thread
// hit_count_recorder(segment_id); // one get opt will cause query to many segments.
// so one get opt only call one hit_heat_buckets, but call many hit_count_recorder
return filter_cache_.check_key(segment_id, key);
}
Expand Down Expand Up @@ -285,7 +292,7 @@ void FilterCacheManager::estimate_counts_for_all(std::map<uint32_t, uint32_t>& a


void FilterCacheManager::try_retrain_model(std::map<uint32_t, uint16_t>& level_recorder,
std::map<uint32_t, std::vector<uint32_t>>& segment_ranges_recorder,
std::map<uint32_t, std::vector<RangeRatePair>>& segment_ranges_recorder,
std::map<uint32_t, uint32_t>& unit_size_recorder) {
// we should guarantee these 3 external recorder share the same keys set
// we need to do this job outside FilterCacheManager
Expand Down Expand Up @@ -347,28 +354,29 @@ void FilterCacheManager::try_retrain_model(std::map<uint32_t, uint16_t>& level_r
std::vector<uint32_t> get_cnts;

auto level_it = level_recorder.begin(); // key range id start with 0
auto ranges_it = segment_ranges_recorder.begin();
auto range_it = segment_ranges_recorder.begin();
auto count_it = last_count_recorder_.begin();
auto label_it = label_recorder.begin();
while (level_it != level_recorder.end() && ranges_it != segment_ranges_recorder.end() &&
while (level_it != level_recorder.end() && range_it != segment_ranges_recorder.end() &&
count_it != last_count_recorder_.end() && label_it != label_recorder.end()) {
assert(level_it->first == ranges_it->first);
assert(level_it->first == range_it->first);
assert(level_it->first == label_it->first);
if (count_it->first < level_it->first) {
count_it ++;
} else if (count_it->first > level_it->first) {
level_it ++;
ranges_it ++;
range_it ++;
label_it ++;
} else {
if (level_it->second > 0) {
// add data row
std::vector<uint32_t> data;
std::sort((range_it->second).begin(), (range_it->second).end(), RangeRatePairGreatorComparor);
data.emplace_back(level_it->second);
for (uint32_t& range_id : ranges_it->second) {
assert(range_id >= 0 && range_id < buckets.size());
data.emplace_back(range_id);
data.emplace_back(uint32_t(SIGNIFICANT_DIGITS_FACTOR * buckets[range_id].hotness_));
for (RangeRatePair& pair : range_it->second) {
assert(pair.range_id >= 0 && pair.range_id < buckets.size());
data.emplace_back(uint32_t(RATE_SIGNIFICANT_DIGITS_FACTOR * pair.rate_in_segment));
data.emplace_back(uint32_t(HOTNESS_SIGNIFICANT_DIGITS_FACTOR * buckets[pair.range_id].hotness_));
}
datas.emplace_back(data);
// add label row
Expand All @@ -378,7 +386,7 @@ void FilterCacheManager::try_retrain_model(std::map<uint32_t, uint16_t>& level_r
}

level_it ++;
ranges_it ++;
range_it ++;
label_it ++;
}
}
Expand All @@ -389,7 +397,7 @@ void FilterCacheManager::try_retrain_model(std::map<uint32_t, uint16_t>& level_r
}

void FilterCacheManager::update_cache_and_heap(std::map<uint32_t, uint16_t>& level_recorder,
std::map<uint32_t, std::vector<uint32_t>>& segment_ranges_recorder) {
std::map<uint32_t, std::vector<RangeRatePair>>& segment_ranges_recorder) {
assert(level_recorder.size() == segment_ranges_recorder.size());
std::vector<uint32_t> segment_ids;
std::vector<std::vector<uint32_t>> datas;
Expand All @@ -410,14 +418,14 @@ void FilterCacheManager::update_cache_and_heap(std::map<uint32_t, uint16_t>& lev
assert(level_it->first == range_it->first);

if (level_it->second > 0) {
segment_ids.emplace_back(level_it->first);

// add data row
std::vector<uint32_t> data;
std::sort((range_it->second).begin(), (range_it->second).end(), RangeRatePairGreatorComparor);
data.emplace_back(level_it->second);
for (uint32_t& range_id : range_it->second) {
assert(range_id >= 0 && range_id < buckets.size());
data.emplace_back(range_id);
data.emplace_back(uint32_t(SIGNIFICANT_DIGITS_FACTOR * buckets[range_id].hotness_));
for (RangeRatePair& pair : range_it->second) {
assert(pair.range_id >= 0 && pair.range_id < buckets.size());
data.emplace_back(uint32_t(RATE_SIGNIFICANT_DIGITS_FACTOR * pair.rate_in_segment));
data.emplace_back(uint32_t(HOTNESS_SIGNIFICANT_DIGITS_FACTOR * buckets[pair.range_id].hotness_));
}
datas.emplace_back(data);
}
Expand Down Expand Up @@ -485,7 +493,7 @@ bool FilterCacheManager::adjust_cache_and_heap() {
void FilterCacheManager::insert_segments(std::vector<uint32_t>& merged_segment_ids, std::vector<uint32_t>& new_segment_ids,
std::map<uint32_t, std::unordered_map<uint32_t, double>>& inherit_infos_recorder,
std::map<uint32_t, uint16_t>& level_recorder, const uint32_t& level_0_base_count,
std::map<uint32_t, std::vector<uint32_t>>& segment_ranges_recorder) {
std::map<uint32_t, std::vector<RangeRatePair>>& segment_ranges_recorder) {
std::unordered_map<uint32_t, uint16_t> segment_units_num_recorder;
std::map<uint32_t, uint32_t> approximate_counts_recorder;
std::set<uint32_t> failed_segment_ids;
Expand Down Expand Up @@ -580,10 +588,10 @@ void FilterCacheManager::insert_segments(std::vector<uint32_t>& merged_segment_i

std::vector<uint32_t> pred_data;
pred_data.emplace_back(level_recorder[new_segment_id]);
for (uint32_t& range_id : segment_ranges_recorder[new_segment_id]) {
assert(range_id >= 0 && range_id < buckets.size());
pred_data.emplace_back(range_id);
pred_data.emplace_back(uint32_t(SIGNIFICANT_DIGITS_FACTOR * buckets[range_id].hotness_));
for (RangeRatePair& pair : segment_ranges_recorder[new_segment_id]) {
assert(pair.range_id >= 0 && pair.range_id < buckets.size());
pred_data.emplace_back(uint32_t(RATE_SIGNIFICANT_DIGITS_FACTOR * pair.rate_in_segment));
pred_data.emplace_back(uint32_t(HOTNESS_SIGNIFICANT_DIGITS_FACTOR * buckets[pair.range_id].hotness_));
}
pred_datas.emplace_back(pred_data);
}
Expand Down
Loading