diff --git a/db/art/clf_model.cc b/db/art/clf_model.cc index ba667f50c..2c8224979 100644 --- a/db/art/clf_model.cc +++ b/db/art/clf_model.cc @@ -35,7 +35,7 @@ void ClfModel::write_debug_dataset() { std::vector 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"); @@ -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 @@ -94,7 +94,7 @@ void ClfModel::write_real_dataset(std::vector>& 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 diff --git a/db/art/clf_model.h b/db/art/clf_model.h index 3c168ee62..7ea317cba 100644 --- a/db/art/clf_model.h +++ b/db/art/clf_model.h @@ -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 { @@ -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& 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; diff --git a/db/art/filter_cache.cc b/db/art/filter_cache.cc index 66ed8fd1d..b5744e9fb 100644 --- a/db/art/filter_cache.cc +++ b/db/art/filter_cache.cc @@ -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) { @@ -176,8 +177,14 @@ void FilterCacheManager::hit_heat_buckets(const std::string& key) { } } +bool FilterCacheManager::make_clf_model_ready(std::vector& 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); } @@ -285,7 +292,7 @@ void FilterCacheManager::estimate_counts_for_all(std::map& a void FilterCacheManager::try_retrain_model(std::map& level_recorder, - std::map>& segment_ranges_recorder, + std::map>& segment_ranges_recorder, std::map& unit_size_recorder) { // we should guarantee these 3 external recorder share the same keys set // we need to do this job outside FilterCacheManager @@ -347,28 +354,29 @@ void FilterCacheManager::try_retrain_model(std::map& level_r std::vector 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 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 @@ -378,7 +386,7 @@ void FilterCacheManager::try_retrain_model(std::map& level_r } level_it ++; - ranges_it ++; + range_it ++; label_it ++; } } @@ -389,7 +397,7 @@ void FilterCacheManager::try_retrain_model(std::map& level_r } void FilterCacheManager::update_cache_and_heap(std::map& level_recorder, - std::map>& segment_ranges_recorder) { + std::map>& segment_ranges_recorder) { assert(level_recorder.size() == segment_ranges_recorder.size()); std::vector segment_ids; std::vector> datas; @@ -410,14 +418,14 @@ void FilterCacheManager::update_cache_and_heap(std::map& 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 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); } @@ -485,7 +493,7 @@ bool FilterCacheManager::adjust_cache_and_heap() { void FilterCacheManager::insert_segments(std::vector& merged_segment_ids, std::vector& new_segment_ids, std::map>& inherit_infos_recorder, std::map& level_recorder, const uint32_t& level_0_base_count, - std::map>& segment_ranges_recorder) { + std::map>& segment_ranges_recorder) { std::unordered_map segment_units_num_recorder; std::map approximate_counts_recorder; std::set failed_segment_ids; @@ -580,10 +588,10 @@ void FilterCacheManager::insert_segments(std::vector& merged_segment_i std::vector 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); } diff --git a/db/art/filter_cache.h b/db/art/filter_cache.h index 5bf96e31e..4ba109503 100644 --- a/db/art/filter_cache.h +++ b/db/art/filter_cache.h @@ -110,17 +110,22 @@ class FilterCacheManager { // TODO: one background thread monitor this func, if return true, call try_retrain_model at once, wait for training end, and call update_cache_and_heap bool need_retrain() { return train_signal_; } - // TODO: one background thread monitor this func, if return true, call try_retrain_model at once and wait for training end. then call update_cache_and_heap - // TODO: if training end, stop this thread, because if is_ready_ is true, is_ready_ will never change to false + // TODO: one background thread monitor this func, if return true, call make_clf_model_ready first, then call try_retrain_model at once and wait for training end. + // TODO: lastly call update_cache_and_heap. if all end, stop this thread, because if is_ready_ is true, is_ready_ will never change to false bool ready_work() { return is_ready_; } // input segment id and target key, check whether target key exist in this segment // return true when target key may exist (may cause false positive fault) // if there is no cache item for this segment, always return true // TODO: normal bloom filter units query, can we put hit_count_recorder outside this func? this will make get opt faster + // TODO: will be called by a get operation, this will block get operation + // TODO: remember to call hit_count_recorder in a background thread bool check_key(const uint32_t& segment_id, const std::string& key); // add 1 to get cnt of specified segment in current long period + // TODO: will be called when calling check_key + // TODO: remember to move this func to a single background thread aside check_key + // TODO: because this func shouldn't block get operations void hit_count_recorder(const uint32_t& segment_id); // copy counts to last_count_recorder and reset counts of current_count_recorder @@ -139,14 +144,26 @@ class FilterCacheManager { // it is like { segment 1: [min_key_1, max_key_1], segment 2: [min_key_2, max_key_2], ... } // return true when heat_buckets is ready, so no need to call this func again // TODO: remember to be called when receiving put opt. Normally, we can make heat_buckets_ ready before YCSB load end, so we can use it in YCSB testing phase + // TODO: every put operation will call a background thread to call make_heat_buckets_ready + // TODO: after this func return true, no need to call this func in put operation + // TODO: remember to move this func to a single background thread aside put operations bool make_heat_buckets_ready(const std::string& key, std::unordered_map>& segment_info_recorder); + // clf_model_ need to determine feature nums before training + // actually, YCSB will load data before testing + // features_nums: [feature_num_1, feature_num_2, ...], it includes all feature_num of all alive segments + // feature_num_k is 2 * (number of key ranges intersecting with segment k) + 1 + // return true when clf_model_ set to ready successfully + // TODO: we need to call make_clf_model_ready before we first call try_retrain_model + // TODO: simply, if ready_work return true, we call make_clf_model_ready at once + bool make_clf_model_ready(std::vector& features_nums); + // add 1 to get cnt of target key range for every get operation // update short periods if get cnt exceeds PERIOD_COUNT // every get opt will make add 1 to only one heat bucket counter // also need to update count records if one long period end // also need to re-calcuate estimated count of current segments and update FilterCacheHeap if one short period end - // TODO: we should call one background thread to exec this func, reducing tail latency + // TODO: we should use one background thread to call this func in every get operation void hit_heat_buckets(const std::string& key); // if one long period end, we need to check effectiveness of model. @@ -165,7 +182,7 @@ class FilterCacheManager { // TODO: because of the time cost of writing csv file, we need to do this func with a background thread // TODO: need real benchmark data to debug this func void try_retrain_model(std::map& level_recorder, - std::map>& segment_ranges_recorder, + std::map>& segment_ranges_recorder, std::map& unit_size_recorder); // after one long period end, we may retrain model. when try_retrain_model func end, we need to predict units num for every segments @@ -177,15 +194,17 @@ class FilterCacheManager { // and update their filter units in filter cache and nodes in heap // so level_recorder keys set and segment_ranges_recorder keys set can be different // TODO: only be called after try_retrain_model (should be guaranteed) - // TODO: we can guarantee this by putting try_retrain_model and update_cache_and_heap into one background thread + // TODO: we can guarantee this by putting try_retrain_model and update_cache_and_heap into only one background thread void update_cache_and_heap(std::map& level_recorder, - std::map>& segment_ranges_recorder); + std::map>& segment_ranges_recorder); // remove merged segments' filter units in the filter cache // also remove related items in FilterCacheHeap // segment_ids: [level_1_segment_1, level_0_segment_1, ...] // level_0_segment_ids: [level_0_segment_1, ...] // TODO: should be called by one background thread + // TODO: this func will be called by insert_segments + // TODO: you can also call this func alone after segments are merged (not suggested) void remove_segments(std::vector& segment_ids, std::set& level_0_segment_ids); // insert new segments into cache @@ -208,10 +227,11 @@ class FilterCacheManager { // level_recorder keys set and segment_ranges_recorder keys set can be different // but should ensure all new segments are in both level_recorder and segment_ranges_recorder // TODO: should be called by one background thread! + // TODO: when old segments are merged into some new segments, call this func in one background thread void insert_segments(std::vector& merged_segment_ids, std::vector& new_segment_ids, std::map>& inherit_infos_recorder, std::map& level_recorder, const uint32_t& level_0_base_count, - std::map>& segment_ranges_recorder); + std::map>& segment_ranges_recorder); // make filter unit adjustment based on two heaps (benefit of enabling one unit & cost of disabling one unit) // simply, we disable one unit of one segment and enable one unit of another segment and guarantee cost < benefit diff --git a/db/art/macros.h b/db/art/macros.h index 566634834..c227accf7 100644 --- a/db/art/macros.h +++ b/db/art/macros.h @@ -162,8 +162,12 @@ namespace ROCKSDB_NAMESPACE { // the path to save model txt file and train dataset csv file #define MODEL_PATH "/pg_wal/ycc/" // we cannot send hotness value (double) to model side, -// so we try multiple hotness value by SIGNIFICANT_DIGITS_FACTOR, then send its integer part to model -#define SIGNIFICANT_DIGITS_FACTOR 1e6 +// so we try multiple hotness value by HOTNESS_SIGNIFICANT_DIGITS_FACTOR, then send its integer part to model +// also we need to multiple key range rate by RATE_SIGNIFICANT_DIGITS_FACTOR +#define HOTNESS_SIGNIFICANT_DIGITS_FACTOR 1e6 +#define RATE_SIGNIFICANT_DIGITS_FACTOR 1e3 +// model feature num max limit : 3 * 30 + 1 +#define MAX_FEATURES_NUM 91 // config micro connecting to LightGBM server @@ -171,7 +175,7 @@ namespace ROCKSDB_NAMESPACE { #define HOST "127.0.0.1" #define PORT "9090" // max size of socket receive buffer size -#define BUFFER_SIZE 8 +#define BUFFER_SIZE 1024 // socket message prefix #define TRAIN_PREFIX "t " #define PREDICT_PREFIX "p "