diff --git a/examples/base.py b/examples/base.py deleted file mode 100644 index 4dfd4fb..0000000 --- a/examples/base.py +++ /dev/null @@ -1,199 +0,0 @@ -import os -import joblib -import zipfile -from shutil import copyfile, rmtree - -import json -from learnware.client import LearnwareClient -from learnware.logger import get_module_logger -from learnware.market import instantiate_learnware_market -from multiprocessing import Pool - -from benchmarks import DataLoader -from config import * -from methods import * -from utils import process_single_aug - -logger = get_module_logger("TableWorkflow", level="INFO") - - -class TableWorkflow: - def __init__(self, learnware_market): - self.learnware_market = learnware_market - - self.root_path = os.path.abspath(os.path.join(__file__, "..")) - self.learnware_pool_path = os.path.join(self.root_path, "data/learnware_pool") - self.learnware_zip_pool_path = os.path.join(self.root_path, "data/zips") - self.example_learnware_path = os.path.join(self.root_path, "data/example_files") - self.model_save_path = os.path.join(self.root_path, "data/uploader_models") - self.result_path = os.path.join(self.root_path, "results") - - os.makedirs(self.learnware_pool_path, exist_ok=True) - os.makedirs(self.learnware_zip_pool_path, exist_ok=True) - os.makedirs(self.model_save_path, exist_ok=True) - os.makedirs(self.result_path, exist_ok=True) - - def _init_dataset(self): - self._prepare_data() - self._prepare_model() - - @staticmethod - def _limited_data(method, test_info, loss_func): - all_scores = [] - for subset in test_info["train_subsets"]: - subset_scores = [] - for sample in subset: - x_train, y_train = sample["x_train"], sample["y_train"] - model = method(x_train, y_train, test_info) - subset_scores.append(loss_func(model.predict(test_info["test_x"]), test_info["test_y"])) - all_scores.append(np.mean(subset_scores)) - return all_scores - - # @staticmethod - # def _limited_data_single_learnware(method, test_info, learnware): - # test_info['single_learnware'] = learnware - # return TableWorkflow._limited_data(method, test_info) - - def test_method(self, test_info, recorders, loss_func=loss_func_rmse): - method_name_full = test_info["method_name"] - method_name = method_name_full if method_name_full == "user_model" else "_".join(method_name_full.split("_")[1:]) - user, idx = test_info["user"], test_info["idx"] - recorder = recorders[method_name_full] - - save_root_path = os.path.join(self.curves_result_path, f"{user}/{user}_{idx}") - os.makedirs(save_root_path, exist_ok=True) - save_path = os.path.join(save_root_path, f"{method_name}.json") - - if method_name == "single_aug": - if test_info["force"] or recorder.should_test_method(user, idx, save_path): - # with Pool() as pool: - # learnware_results = pool.starmap( - # self._limited_data_single_learnware, - # [(test_methods[method_name], test_info, learnware) for learnware in test_info['learnwares']] - # ) - # for scores in learnware_results: - # recorders[method_name].record(user, idx, scores) - - for learnware in test_info['learnwares']: - test_info['single_learnware'] = learnware - scores = self._limited_data(test_methods[method_name_full], test_info, loss_func) - recorder.record(user, idx, scores) - - process_single_aug(user, idx, scores, recorders, save_root_path) - recorder.save(save_path) - logger.info(f"Method {method_name} on {user}_{idx} finished") - else: - process_single_aug(user, idx, recorder.data[user][str(idx)], recorders, save_root_path) - logger.info(f"Method {method_name} on {user}_{idx} already exists") - else: - if test_info["force"] or recorder.should_test_method(user, idx, save_path): - scores = self._limited_data(test_methods[method_name_full], test_info, loss_func) - recorder.record(user, idx, scores) - recorder.save(save_path) - logger.info(f"Method {method_name} on {user}_{idx} finished") - else: - logger.info(f"Method {method_name} on {user}_{idx} already exists") - - def prepare_market(self, name, market_id, regenerate_flag=False): - if regenerate_flag: - self._init_dataset() - market = instantiate_learnware_market(name=name, market_id=market_id, rebuild=True) - client = LearnwareClient() - - full_descriptions_dir = os.path.join("./data/full_descriptions.json") - with open(full_descriptions_dir, "rb") as f: - full_descriptions = json.load(f) - - for uploader in self.learnware_market: - data_loader = DataLoader(uploader) - idx_list = data_loader.get_shop_ids() - for i, idx in enumerate(idx_list): - feature_descriptions = data_loader.get_raw_data(idx)[-1] - feature_dim = len(feature_descriptions) - feature_descriptions_dict = {str(i): feature_descriptions[i] for i in range(feature_dim)} - input_description = {"Dimension": feature_dim, "Description": feature_descriptions_dict} - - name_and_description = full_descriptions[uploader][i] - semantic_spec = client.create_semantic_specification( - name=name_and_description["name"], - description=name_and_description["description"], - data_type="Table", - task_type="Regression", - library_type="Others", - license=["MIT"], - scenarios=["Business"], - input_description=input_description, - output_description=output_description, - ) - - learnware_zip_path = self._prepare_learnware(data_loader, idx) - market.add_learnware(learnware_zip_path, semantic_spec) - - # if use pretrained market mapping - if name == "hetero": - learnware_ids = market.get_learnware_ids() - market.learnware_organizer._update_learware_hetero_spec(learnware_ids) - - logger.info("Total Item: %d" % (len(market))) - - def _prepare_data(self): - for uploader in self.learnware_market: - data_loader = DataLoader(uploader) - data_loader.regenerate_raw_data() - - def _prepare_model(self, use_exist=True): - self.learnware_num = 0 - for uploader in self.learnware_market: - data_loader = DataLoader(uploader) - idx_list = data_loader.get_shop_ids() - self.learnware_num += len(idx_list) - for idx in idx_list: - logger.info(f"Train on uploader: {uploader}_{idx}") - idx_model_save_path = os.path.join(self.model_save_path, f"{uploader}_{idx}.out") - if not use_exist: - x_train, y_train, x_val, y_val, _ = data_loader.get_raw_data(idx) - data_loader.train_a_model(x_train, y_train, x_val, y_val, save_dir=idx_model_save_path) - else: - uploader_dataset = uploader.split("_")[0] - model = data_loader.get_model(idx) - if uploader_dataset == "corporacion": - model.save_model(idx_model_save_path) - elif uploader_dataset == "pfs": - joblib.dump(model, idx_model_save_path) - else: - logger.error(f"Not supported dataset type {uploader_dataset}") - - logger.info(f"Model saved to {idx_model_save_path}") - - def _prepare_learnware(self, data_loader, idx): - zip_path = os.path.join(self.learnware_zip_pool_path, f"{data_loader.dataset}_{idx}") - dir_path = os.path.join(self.learnware_pool_path, f"{data_loader.dataset}_{idx}") - model_path = os.path.join(self.model_save_path, f"{data_loader.dataset}_{idx}.out") - os.makedirs(dir_path, exist_ok=True) - - stat_spec, _ = data_loader.get_rkme(idx) - init_file = os.path.join(dir_path, "__init__.py") - yaml_file = os.path.join(dir_path, "learnware.yaml") - env_file = os.path.join(dir_path, "environment.yaml") - model_file = os.path.join(dir_path, "model.out") - - stat_spec.save(os.path.join(dir_path, "rkme.json")) - copyfile(os.path.join(self.example_learnware_path, f"{data_loader.dataset}/__init__.py"), init_file) - copyfile(os.path.join(self.example_learnware_path, f"{data_loader.dataset}/learnware.yaml"), yaml_file) - copyfile(os.path.join(self.example_learnware_path, "environment.yaml"), env_file) - copyfile(model_path, model_file) - - zip_file = zip_path + ".zip" - with zipfile.ZipFile(zip_file, "w") as zip_obj: - for foldername, _, filenames in os.walk(dir_path): - for filename in filenames: - file_path = os.path.join(foldername, filename) - zip_info = zipfile.ZipInfo(filename) - zip_info.compress_type = zipfile.ZIP_STORED - with open(file_path, "rb") as file: - zip_obj.writestr(zip_info, file.read()) - - rmtree(dir_path) # rm -r dir_path - return zip_file - - \ No newline at end of file diff --git a/examples/dataset_table_workflow/base.py b/examples/dataset_table_workflow/base.py new file mode 100644 index 0000000..fe31bd6 --- /dev/null +++ b/examples/dataset_table_workflow/base.py @@ -0,0 +1,105 @@ +import os +import time +import pandas +import random +import tempfile +import numpy as np +from learnware.client import LearnwareClient +from learnware.logger import get_module_logger +from learnware.market import instantiate_learnware_market +from learnware.tests.benchmarks import LearnwareBenchmark + +from config import * +from methods import * +from utils import process_single_aug + +logger = get_module_logger("base_table", level="INFO") + + +class TableWorkflow: + def __init__(self, benchmark_config, name="easy", rebuild=False): + self.root_path = os.path.abspath(os.path.join(__file__, "..")) + self.result_path = os.path.join(self.root_path, "results") + self.curves_result_path = os.path.join(self.root_path, "curves") + os.makedirs(self.result_path, exist_ok=True) + os.makedirs(self.curves_result_path, exist_ok=True) + self._prepare_market(benchmark_config, name, rebuild) + + @staticmethod + def _limited_data(method, test_info, loss_func): + all_scores = [] + for subset in test_info["train_subsets"]: + subset_scores = [] + for sample in subset: + x_train, y_train = sample["x_train"], sample["y_train"] + model = method(x_train, y_train, test_info) + subset_scores.append(loss_func(model.predict(test_info["test_x"]), test_info["test_y"])) + all_scores.append(np.mean(subset_scores)) + return all_scores + + @staticmethod + def get_train_subsets(train_x, train_y): + np.random.seed(1) + random.seed(1) + train_subsets = [] + for n_label, repeated in zip(n_labeled_list, n_repeat_list): + train_subsets.append([]) + if n_label > len(train_x): + n_label = len(train_x) + for _ in range(repeated): + x_train, y_train = zip(*random.sample(list(zip(train_x, train_y)), k=n_label)) + train_subsets[-1].append({"x_train": np.array(x_train), "y_train": np.array(list(y_train))}) + return train_subsets + + def _prepare_market(self, benchmark_config, name, rebuild): + client = LearnwareClient() + self.benchmark = LearnwareBenchmark().get_benchmark(benchmark_config) + self.market = instantiate_learnware_market(market_id=self.benchmark.name, name=name, rebuild=rebuild) + self.user_semantic = client.get_semantic_specification(self.benchmark.learnware_ids[0]) + self.user_semantic["Name"]["Values"] = "" + + if len(self.market) == 0 or rebuild == True: + for learnware_id in self.benchmark.learnware_ids: + with tempfile.TemporaryDirectory(prefix="table_benchmark_") as tempdir: + zip_path = os.path.join(tempdir, f"{learnware_id}.zip") + for i in range(20): + try: + semantic_spec = client.get_semantic_specification(learnware_id) + client.download_learnware(learnware_id, zip_path) + self.market.add_learnware(zip_path, semantic_spec) + break + except: + time.sleep(1) + continue + + def test_method(self, test_info, recorders, loss_func=loss_func_rmse): + method_name_full = test_info["method_name"] + method_name = method_name_full if method_name_full == "user_model" else "_".join(method_name_full.split("_")[1:]) + user, idx = test_info["user"], test_info["idx"] + recorder = recorders[method_name_full] + + save_root_path = os.path.join(self.curves_result_path, user, f"{user}_{idx}") + os.makedirs(save_root_path, exist_ok=True) + save_path = os.path.join(save_root_path, f"{method_name}.json") + + if method_name == "single_aug": + if test_info["force"] or recorder.should_test_method(user, idx, save_path): + for learnware in test_info["learnwares"]: + test_info["single_learnware"] = [learnware] + scores = self._limited_data(test_methods[method_name_full], test_info, loss_func) + recorder.record(user, idx, scores) + + process_single_aug(user, idx, scores, recorders, save_root_path) + recorder.save(save_path) + logger.info(f"Method {method_name} on {user}_{idx} finished") + else: + process_single_aug(user, idx, recorder.data[user][str(idx)], recorders, save_root_path) + logger.info(f"Method {method_name} on {user}_{idx} already exists") + else: + if test_info["force"] or recorder.should_test_method(user, idx, save_path): + scores = self._limited_data(test_methods[method_name_full], test_info, loss_func) + recorder.record(user, idx, scores) + recorder.save(save_path) + logger.info(f"Method {method_name} on {user}_{idx} finished") + else: + logger.info(f"Method {method_name} on {user}_{idx} already exists") \ No newline at end of file diff --git a/examples/dataset_table_workflow/config.py b/examples/dataset_table_workflow/config.py index c1e0ab8..14b60ed 100644 --- a/examples/dataset_table_workflow/config.py +++ b/examples/dataset_table_workflow/config.py @@ -1,3 +1,6 @@ +from learnware.tests.benchmarks import BenchmarkConfig + + n_labeled_list = [100, 200, 500, 1000, 2000, 4000, 6000, 8000, 10000] n_repeat_list = [10, 10, 10, 3, 3, 3, 3, 3, 3] @@ -16,30 +19,11 @@ labels = { 'user_model': "User Model", 'single_aug': "Single Learnware Reuse (Select)", "select_score": "Single Learnware Reuse (Select)", - # "Single Learnware Reuse (Avg)", - # "Single Learnware Reuse (Oracle)", 'multiple_aug': "Multiple Learnware Reuse (FeatAug)", 'ensemble_pruning': "Multiple Learnware Reuse (EnsemblePrune)", 'multiple_avg': "Multiple Learnware Reuse (Averaging)" } -output_description = { - "Dimension": 1, - "Description": { - "0": "Product sales on the date.", - }, -} - -user_semantic = { - "Data": {"Values": ["Table"], "Type": "Class"}, - "Task": {"Values": ["Regression"], "Type": "Class"}, - "Library": {"Values": ["Others"], "Type": "Class"}, - "Scenario": {"Values": ["Business"], "Type": "Tag"}, - "Description": {"Values": "", "Type": "String"}, - "Name": {"Values": "", "Type": "String"}, - "Output": output_description, -} - align_model_params = { "network_type": "ArbitraryMapping", # ["ArbitraryMapping", "BaseMapping", "BaseMapping_BN", "BaseMapping_Dropout"] "num_epoch": 50, @@ -63,3 +47,87 @@ market_mapping_params = { "ffn_dim": 512, # [128, 256, 512, 768, 1024], # the dimension of feed-forward layer in the transformer layer "activation": "leakyrelu", } + +user_model_params = { + "Corporacion": { + "lgb": { + "params": { + "num_leaves": 31, + "objective": "regression", + "learning_rate": 0.1, + "feature_fraction": 0.8, + "bagging_fraction": 0.8, + "bagging_freq": 2, + "metric": "l2", + "num_threads": 4, + "verbose": -1, + }, + "MAX_ROUNDS": 500, + "early_stopping_rounds": 50, + } + } +} + +homo_table_benchmark_config = BenchmarkConfig( + name="Corporacion", + user_num=54, + learnware_ids=[ + "00000912", + "00000911", + "00000910", + "00000909", + "00000908", + "00000907", + "00000906", + "00000905", + "00000904", + "00000903", + "00000902", + "00000901", + "00000900", + "00000899", + "00000898", + "00000897", + "00000896", + "00000895", + "00000894", + "00000893", + "00000892", + "00000891", + "00000890", + "00000889", + "00000888", + "00000887", + "00000886", + "00000885", + "00000884", + "00000883", + "00000882", + "00000881", + "00000880", + "00000879", + "00000878", + "00000877", + "00000876", + "00000875", + "00000874", + "00000873", + "00000872", + "00000871", + "00000870", + "00000869", + "00000868", + "00000867", + "00000866", + "00000865", + "00000864", + "00000863", + "00000862", + "00000861", + "00000860", + "00000859" + ], + test_data_path="Corporacion/test_data.zip", + train_data_path="Corporacion/train_data.zip", + extra_info_path="Corporacion/extra_info.zip", +) \ No newline at end of file diff --git a/examples/dataset_table_workflow/homo.py b/examples/dataset_table_workflow/homo.py index efa6deb..d6bb8df 100644 --- a/examples/dataset_table_workflow/homo.py +++ b/examples/dataset_table_workflow/homo.py @@ -1,204 +1,175 @@ import os import warnings -from collections import defaultdict - +import numpy as np warnings.filterwarnings("ignore") -import numpy as np -from matplotlib import pyplot as plt -from functools import partial -import learnware.specification as specification +from learnware.market import BaseUserInfo from learnware.logger import get_module_logger -from learnware.market import instantiate_learnware_market, BaseUserInfo +from learnware.specification import generate_stat_spec from learnware.reuse import AveragingReuser, JobSelectorReuser, EnsemblePruningReuser -from benchmarks import DataLoader -from base import TableWorkflow, user_semantic from methods import * -from config import n_labeled_list, n_repeat_list +from base import TableWorkflow +from config import n_labeled_list from utils import Recorder, plot_performance_curves -logger = get_module_logger("corporacion_test", level="INFO") -learnware_market = ["corporacion_bojan", "corporacion_lee", "corporacion_lingzhi"] -users = ["corporacion_lingzhi"] +logger = get_module_logger("homo_table", level="INFO") -class CorporacionDatasetWorkflow(TableWorkflow): - def __init__(self, reload_market=False, regenerate_flag=False): - super(CorporacionDatasetWorkflow, self).__init__(learnware_market) - self.curves_result_path = os.path.join(self.result_path, "curves") - self.figs_result_path = os.path.join(self.result_path, "figs") - - os.makedirs(self.curves_result_path, exist_ok=True) - os.makedirs(self.figs_result_path, exist_ok=True) - - if reload_market: - self.prepare_market(name="easy", market_id="corporacion", regenerate_flag=regenerate_flag) - - def test_homo_unlabeled(self): - corporacion_market = instantiate_learnware_market(market_id="corporacion") - logger.info("Total Item: %d" % len(corporacion_market)) - - learnware_rmse_list = defaultdict(list) - job_selector_score_list = defaultdict(list) - ensemble_score_list = defaultdict(list) - pruning_score_list = defaultdict(list) - - for user in learnware_market: - corporacion = DataLoader(user) - idx_list = corporacion.get_shop_ids() - for idx in idx_list: - _, _, test_x, test_y, _ = corporacion.get_raw_data(idx) - user_stat_spec = specification.RKMETableSpecification() - user_stat_spec.generate_stat_spec_from_data(X=test_x) - user_info = BaseUserInfo( - semantic_spec=user_semantic, stat_info={"RKMETableSpecification": user_stat_spec} - ) - logger.info(f"Searching Market for user: {user}_{idx}") - - search_result = corporacion_market.search_learnware(user_info, max_search_num=10) - single_result = search_result.get_single_results() - multiple_result = search_result.get_multiple_results() - - logger.info(f"search result of user {user}_{idx}:") - logger.info( - f"single model num: {len(single_result)}, max_score: {single_result[0].score}, min_score: {single_result[-1].score}" - ) - - l = len(single_result) - rmse_list = [] - for idx in range(l): - learnware = single_result[idx].learnware - pred_y = learnware.predict(test_x) - rmse_list.append(loss_func_mse(pred_y, test_y)) - logger.info( - f"Top1-score: {single_result[0].score}, learnware_id: {single_result[0].learnware.id}, rmse: {rmse_list[0]}" - ) - - if len(multiple_result) > 0: - mixture_id = " ".join([learnware.id for learnware in multiple_result[0].learnwares]) - logger.info(f"mixture_score: {multiple_result[0].score}, mixture_learnware: {mixture_id}") - mixture_learnware_list = multiple_result[0].learnwares - else: - mixture_learnware_list = [single_result[0].learnware] - - # test reuse (job selector) - reuse_baseline = JobSelectorReuser(learnware_list=mixture_learnware_list, herding_num=100) - reuse_predict = reuse_baseline.predict(user_data=test_x) - reuse_score = loss_func_mse(reuse_predict, test_y) - job_selector_score_list[user].append(reuse_score) - logger.info(f"mixture reuse rmse (job selector): {reuse_score}") - - # test reuse (ensemble) - reuse_ensemble = AveragingReuser(learnware_list=mixture_learnware_list, mode="mean") - ensemble_predict_y = reuse_ensemble.predict(user_data=test_x) - ensemble_score = loss_func_mse(ensemble_predict_y, test_y) - ensemble_score_list[user].append(ensemble_score) - logger.info(f"mixture reuse rmse (ensemble): {ensemble_score}") - - # test reuse (ensemblePruning) - reuse_pruning = EnsemblePruningReuser(learnware_list=mixture_learnware_list, mode="regression") - pruning_predict_y = reuse_pruning.predict(user_data=test_x) - pruning_score = loss_func_mse(pruning_predict_y, test_y) - pruning_score_list[user].append(pruning_score) - logger.info(f"mixture reuse rmse (ensemble Pruning): {pruning_score}\n") - - learnware_rmse_list[user].append(rmse_list) - - for user in learnware_market: - logger.info(f"User Dataset: {user}") - - single_list = np.array(learnware_rmse_list[user]) - select_score_list = [lst[0] for lst in single_list] - avg_score_list = [np.mean(lst, axis=0) for lst in single_list] - oracle_score_list = [np.min(lst, axis=0) for lst in single_list] - logger.info( - "RMSE of selected learnware: %.3f +/- %.3f, Average performance: %.3f +/- %.3f, Oracle performace: %.3f +/- %.3f" - % ( - np.mean(select_score_list), - np.std(select_score_list), - np.mean(avg_score_list), - np.std(avg_score_list), - np.mean(oracle_score_list), - np.std(oracle_score_list), - ) - ) - logger.info( - "Average Job Selector Reuse Performance: %.3f +/- %.3f" - % (np.mean(job_selector_score_list[user]), np.std(job_selector_score_list[user])) +class CorporacionDatasetWorkflow(TableWorkflow): + def unlabeled_homo_table_example(self): + logger.info("Total Item: %d" % (len(self.market))) + learnware_rmse_list = [] + single_score_list = [] + job_selector_score_list = [] + ensemble_score_list = [] + pruning_score_list = [] + all_learnwares = self.market.get_learnwares() + + user = self.benchmark.name + for idx in range(self.benchmark.user_num): + test_x, test_y = self.benchmark.get_test_data(user_ids=idx) + test_x, test_y = test_x.values, test_y.values + user_stat_spec = generate_stat_spec(type="table", X=test_x) + user_info = BaseUserInfo( + semantic_spec=self.user_semantic, stat_info={user_stat_spec.type: user_stat_spec} ) + logger.info(f"Searching Market for user: {user}_{idx}") + + search_result = self.market.search_learnware(user_info) + single_result = search_result.get_single_results() + multiple_result = search_result.get_multiple_results() + + logger.info(f"search result of user {user}_{idx}:") logger.info( - "Averaging Ensemble Reuse Performance: %.3f +/- %.3f" - % (np.mean(ensemble_score_list[user]), np.std(ensemble_score_list[user])) + f"single model num: {len(single_result)}, max_score: {single_result[0].score}, min_score: {single_result[-1].score}" ) + + pred_y = single_result[0].learnware.predict(test_x) + single_score_list.append(loss_func_rmse(pred_y, test_y)) + + rmse_list = [] + for learnware in all_learnwares: + pred_y = learnware.predict(test_x) + rmse_list.append(loss_func_rmse(pred_y, test_y)) logger.info( - "Selective Ensemble Reuse Performance: %.3f +/- %.3f" - % (np.mean(pruning_score_list[user]), np.std(pruning_score_list[user])) + f"Top1-score: {single_result[0].score}, learnware_id: {single_result[0].learnware.id}, rmse: {single_score_list[-1]}" ) - def test_homo_labeled(self): - corporacion_market = instantiate_learnware_market(market_id="corporacion") - logger.info("Total Item: %d" % len(corporacion_market)) - + if len(multiple_result) > 0: + mixture_id = " ".join([learnware.id for learnware in multiple_result[0].learnwares]) + logger.info(f"mixture_score: {multiple_result[0].score}, mixture_learnware: {mixture_id}") + mixture_learnware_list = multiple_result[0].learnwares + else: + mixture_learnware_list = [single_result[0].learnware] + + # test reuse (job selector) + reuse_baseline = JobSelectorReuser(learnware_list=mixture_learnware_list, herding_num=100) + reuse_predict = reuse_baseline.predict(user_data=test_x) + reuse_score = loss_func_rmse(reuse_predict, test_y) + job_selector_score_list.append(reuse_score) + logger.info(f"mixture reuse rmse (job selector): {reuse_score}") + + # test reuse (ensemble) + reuse_ensemble = AveragingReuser(learnware_list=mixture_learnware_list, mode="mean") + ensemble_predict_y = reuse_ensemble.predict(user_data=test_x) + ensemble_score = loss_func_rmse(ensemble_predict_y, test_y) + ensemble_score_list.append(ensemble_score) + logger.info(f"mixture reuse rmse (ensemble): {ensemble_score}") + + # test reuse (ensemblePruning) + reuse_pruning = EnsemblePruningReuser(learnware_list=mixture_learnware_list, mode="regression") + pruning_predict_y = reuse_pruning.predict(user_data=test_x) + pruning_score = loss_func_rmse(pruning_predict_y, test_y) + pruning_score_list.append(pruning_score) + logger.info(f"mixture reuse rmse (ensemble Pruning): {pruning_score}\n") + + learnware_rmse_list.append(rmse_list) + + single_list = np.array(learnware_rmse_list) + avg_score_list = [np.mean(lst, axis=0) for lst in single_list] + oracle_score_list = [np.min(lst, axis=0) for lst in single_list] + + logger.info( + "RMSE of selected learnware: %.3f +/- %.3f, Average performance: %.3f +/- %.3f, Oracle performace: %.3f +/- %.3f" + % ( + np.mean(single_score_list), + np.std(single_score_list), + np.mean(avg_score_list), + np.std(avg_score_list), + np.mean(oracle_score_list), + np.std(oracle_score_list), + ) + ) + logger.info( + "Average Job Selector Reuse Performance: %.3f +/- %.3f" + % (np.mean(job_selector_score_list), np.std(job_selector_score_list)) + ) + logger.info( + "Averaging Ensemble Reuse Performance: %.3f +/- %.3f" + % (np.mean(ensemble_score_list), np.std(ensemble_score_list)) + ) + logger.info( + "Selective Ensemble Reuse Performance: %.3f +/- %.3f" + % (np.mean(pruning_score_list), np.std(pruning_score_list)) + ) + + def labeled_homo_table_example(self): + logger.info("Total Item: %d" % (len(self.market))) methods = ["user_model", "homo_single_aug", "homo_multiple_aug", "homo_multiple_avg", "homo_ensemble_pruning"] recorders = {method: Recorder() for method in methods} - methods_to_retest = [] - for user in users: - data_loader = DataLoader(user) - idx_list = data_loader.get_shop_ids() - for idx in idx_list: - _, _, test_x, test_y, _ = data_loader.get_raw_data(idx) - train_subsets = data_loader.get_labeled_training_data( - idx, - size_list=n_labeled_list, - n_repeat_list=n_repeat_list - ) - - user_stat_spec = specification.RKMETableSpecification() - user_stat_spec.generate_stat_spec_from_data(X=test_x) - user_info = BaseUserInfo( - semantic_spec=user_semantic, stat_info={"RKMETableSpecification": user_stat_spec} - ) - logger.info(f"Searching Market for user: {user}_{idx}") - - search_result = corporacion_market.search_learnware(user_info, max_search_num=10) - single_result = search_result.get_single_results() - multiple_result = search_result.get_multiple_results() - - logger.info(f"search result of user {user}_{idx}:") - logger.info( - f"single model num: {len(single_result)}, max_score: {single_result[0].score}, min_score: {single_result[-1].score}" - ) - - if len(multiple_result) > 0: - mixture_id = " ".join([learnware.id for learnware in multiple_result[0].learnwares]) - logger.info(f"mixture_score: {multiple_result[0].score}, mixture_learnware: {mixture_id}") - mixture_learnware_list = multiple_result[0].learnwares - else: - mixture_learnware_list = [single_result[0].learnware] - - test_info = {"user": user, "idx": idx, "train_subsets": train_subsets, "test_x": test_x, "test_y": test_y} - common_config = {"multiple_learnwares": mixture_learnware_list} - method_configs = { - "user_model": {"data_loader": data_loader}, - "homo_single_aug": {"single_learnware": [single_result[0].learnware]}, - "homo_multiple_aug": common_config, - "homo_multiple_avg": common_config, - "homo_ensemble_pruning": common_config - } - - for method_name in methods: - # self.test_method(method_name, HeteroScoringMethods.__dict__[f"{method_name}_score"], recorders, test_info) - logger.info(f"Testing method {method_name}") - test_info["method_name"] = method_name - test_info["force"] = method_name in methods_to_retest - test_info.update(method_configs[method_name]) - self.test_method(test_info, recorders, loss_func=loss_func_mse) + user = self.benchmark.name + for idx in range(self.benchmark.user_num): + test_x, test_y = self.benchmark.get_test_data(user_ids=idx) + test_x, test_y = test_x.values, test_y.values - for method, recorder in recorders.items(): - recorder.save(os.path.join(self.curves_result_path, f"{user}_{method}_performance.json")) - - methods_to_plot = ["user_model", "homo_single_aug", "homo_ensemble_pruning"] - plot_performance_curves(user, {method: recorders[method] for method in methods_to_plot}, task="Homo", n_labeled_list=n_labeled_list) \ No newline at end of file + train_x, train_y = self.benchmark.get_train_data(user_ids=idx) + train_x, train_y = train_x.values, train_y.values + train_subsets = self.get_train_subsets(train_x, train_y) + + user_stat_spec = generate_stat_spec(type="table", X=test_x) + user_info = BaseUserInfo( + semantic_spec=self.user_semantic, stat_info={"RKMETableSpecification": user_stat_spec} + ) + logger.info(f"Searching Market for user: {user}_{idx}") + + search_result = self.market.search_learnware(user_info) + single_result = search_result.get_single_results() + multiple_result = search_result.get_multiple_results() + + logger.info(f"search result of user {user}_{idx}:") + logger.info( + f"single model num: {len(single_result)}, max_score: {single_result[0].score}, min_score: {single_result[-1].score}" + ) + + if len(multiple_result) > 0: + mixture_id = " ".join([learnware.id for learnware in multiple_result[0].learnwares]) + logger.info(f"mixture_score: {multiple_result[0].score}, mixture_learnware: {mixture_id}") + mixture_learnware_list = multiple_result[0].learnwares + else: + mixture_learnware_list = [single_result[0].learnware] + + test_info = {"user": user, "idx": idx, "train_subsets": train_subsets, "test_x": test_x, "test_y": test_y} + common_config = {"learnwares": mixture_learnware_list} + method_configs = { + "user_model": {"dataset": self.benchmark.name, "model_type": "lgb"}, + "homo_single_aug": {"learnwares": [single_result[0].learnware]}, + "homo_multiple_aug": common_config, + "homo_multiple_avg": common_config, + "homo_ensemble_pruning": common_config + } + + for method_name in methods: + logger.info(f"Testing method {method_name}") + test_info["method_name"] = method_name + test_info["force"] = method_name in methods_to_retest + test_info.update(method_configs[method_name]) + self.test_method(test_info, recorders, loss_func=loss_func_rmse) + + for method, recorder in recorders.items(): + recorder.save(os.path.join(self.curves_result_path, f"{user}_{method}_performance.json")) + + methods_to_plot = ["user_model", "homo_single_aug", "homo_ensemble_pruning"] + plot_performance_curves(user, {method: recorders[method] for method in methods_to_plot}, task="Homo", n_labeled_list=n_labeled_list) \ No newline at end of file diff --git a/examples/dataset_table_workflow/main.py b/examples/dataset_table_workflow/main.py deleted file mode 100644 index 19911b8..0000000 --- a/examples/dataset_table_workflow/main.py +++ /dev/null @@ -1,34 +0,0 @@ -import fire -from pyinstrument import Profiler - -from dataset_corporacion_workflow import CorporacionDatasetWorkflow -from dataset_heterogeneous_workflow import HeterogeneousWorkflow - -workflow_mapping = { - "test_homo_unlabeled": CorporacionDatasetWorkflow, - "test_homo_labeled": CorporacionDatasetWorkflow, - "test_hetero_unlabeled": HeterogeneousWorkflow, - "test_hetero_labeled": HeterogeneousWorkflow, -} - - -def main(): - def dispatch(command): - if command in workflow_mapping: - workflow_class = workflow_mapping[command] - workflow_instance = workflow_class() - getattr(workflow_instance, command)() - else: - print(f"No workflow found for command: {command}") - - fire.Fire(dispatch) - - -if __name__ == "__main__": - profiler = Profiler() - profiler.start() - - main() - - profiler.stop() - profiler.print() diff --git a/examples/dataset_table_workflow/methods.py b/examples/dataset_table_workflow/methods.py index 11580fb..fe40eef 100644 --- a/examples/dataset_table_workflow/methods.py +++ b/examples/dataset_table_workflow/methods.py @@ -1,23 +1,22 @@ import numpy as np from sklearn.metrics import mean_squared_error -from sklearn.model_selection import train_test_split # Add missing import -from loguru import logger +from sklearn.model_selection import train_test_split from learnware.reuse import AveragingReuser, EnsemblePruningReuser, FeatureAugmentReuser, HeteroMapAlignLearnware -from examples.dataset_table_workflow.config import align_model_params +from config import align_model_params +from train import train_model + def loss_func_rmse(y_true, y_pred): return np.sqrt(mean_squared_error(y_true, y_pred)) -def loss_func_mse(y_true, y_pred): - return mean_squared_error(y_true, y_pred) def user_model_score(x_train, y_train, test_info): - data_loader = test_info["data_loader"] x_train, x_val, y_train, y_val = train_test_split(x_train, y_train, test_size=0.2, random_state=42) - user_model = data_loader.train_a_model(x_train, y_train, x_val, y_val) + user_model = train_model(x_train, y_train, x_val, y_val, test_info) return user_model + class HomoScoringMethods: @staticmethod def single_aug_score(x_train, y_train, test_info): @@ -28,20 +27,20 @@ class HomoScoringMethods: @staticmethod def multiple_aug_score(x_train, y_train, test_info): - multiple_learnwares = test_info["multiple_learnwares"] + multiple_learnwares = test_info["learnwares"] reuse_multiple_augment = FeatureAugmentReuser(multiple_learnwares, mode="regression") reuse_multiple_augment.fit(x_train=x_train, y_train=y_train) return reuse_multiple_augment @staticmethod def multiple_avg_score(x_train, y_train, test_info): - multiple_learnwares = test_info["multiple_learnwares"] + multiple_learnwares = test_info["learnwares"] reuse_multiple_avg = AveragingReuser(multiple_learnwares, mode="mean") return reuse_multiple_avg @staticmethod def multiple_ensemble_pruning_score(x_train, y_train, test_info): - multiple_learnwares = test_info["multiple_learnwares"] + multiple_learnwares = test_info["learnwares"] if len(multiple_learnwares) == 1: return multiple_learnwares[0] reuse_pruning = EnsemblePruningReuser(multiple_learnwares, mode="regression") @@ -51,7 +50,7 @@ class HomoScoringMethods: class HeteroMethods: @staticmethod - def create_hetero_learnware_list(learnware_list, user_rkme, x_train, y_train): # Fix typo in method name + def create_hetero_learnware_list(learnware_list, user_rkme, x_train, y_train): hetero_learnware_list = [] for learnware in learnware_list: hetero_learnware = HeteroMapAlignLearnware(learnware, mode="regression", **align_model_params) @@ -68,7 +67,7 @@ class HeteroMethods: @staticmethod def multiple_aug_score(x_train, y_train, test_info): - user_rkme, multiple_learnwares = test_info["user_rkme"], test_info["multiple_learnwares"] + user_rkme, multiple_learnwares = test_info["user_rkme"], test_info["learnwares"] hetero_learnware_list = HeteroMethods.create_hetero_learnware_list(multiple_learnwares, user_rkme, x_train, y_train) reuse_multiple_augment = FeatureAugmentReuser(hetero_learnware_list, mode="regression") reuse_multiple_augment.fit(x_train=x_train, y_train=y_train) @@ -76,7 +75,7 @@ class HeteroMethods: @staticmethod def multiple_ensemble_pruning_score(x_train, y_train, test_info): - user_rkme, multiple_learnwares = test_info["user_rkme"], test_info["multiple_learnwares"] + user_rkme, multiple_learnwares = test_info["user_rkme"], test_info["learnwares"] hetero_learnware_list = HeteroMethods.create_hetero_learnware_list(multiple_learnwares, user_rkme, x_train, y_train) if len(hetero_learnware_list) == 1: return hetero_learnware_list[0] @@ -86,7 +85,7 @@ class HeteroMethods: @staticmethod def multiple_avg_score(x_train, y_train, test_info): - user_rkme, multiple_learnwares = test_info["user_rkme"], test_info["multiple_learnwares"] + user_rkme, multiple_learnwares = test_info["user_rkme"], test_info["learnwares"] hetero_learnware_list = HeteroMethods.create_hetero_learnware_list(multiple_learnwares, user_rkme, x_train, y_train) reuse_multiple_avg = AveragingReuser(hetero_learnware_list, mode="mean") return reuse_multiple_avg diff --git a/examples/dataset_table_workflow/train.py b/examples/dataset_table_workflow/train.py new file mode 100644 index 0000000..7ac2b3a --- /dev/null +++ b/examples/dataset_table_workflow/train.py @@ -0,0 +1,48 @@ +import numpy as np +import lightgbm as lgb +from lightgbm import early_stopping +from sklearn.metrics import mean_squared_error + +from learnware.logger import get_module_logger +from config import user_model_params + +logger = get_module_logger("train_table", level="INFO") + + +def train_lgb(X_train, y_train, X_val, y_val, dataset): + logger.info("Training and predicting models...") + model_param = user_model_params[dataset]["lgb"] + params = model_param["params"] + + MAX_ROUNDS = model_param["MAX_ROUNDS"] + val_pred = [] + cate_vars = [] + + logger.info(f"{np.shape(X_train)}, {np.shape(y_train)}, {np.shape(X_val)}, {np.shape(y_val)}") + + dtrain = lgb.Dataset(X_train, label=y_train, categorical_feature=cate_vars) + dval = lgb.Dataset(X_val, label=y_val, reference=dtrain, categorical_feature=cate_vars) + bst = lgb.train( + params, + dtrain, + num_boost_round=MAX_ROUNDS, + valid_sets=[dtrain, dval], + callbacks=[early_stopping(model_param["early_stopping_rounds"], verbose=False)] + ) + val_pred.append(bst.predict(X_val, num_iteration=bst.best_iteration or MAX_ROUNDS)) + logger.info(f"Validation mse:{mean_squared_error(y_val, np.array(val_pred).transpose())}") + + return bst + + +def train_ridge(X_train, y_train, X_val, y_val, dataset): + pass + + +def train_model(X_train, y_train, X_val, y_val, test_info): + dataset = test_info["dataset"] + model_type = test_info["model_type"] + assert model_type in ["lgb", "ridge"] + + if model_type == "lgb": + return train_lgb(X_train, y_train, X_val, y_val, dataset) \ No newline at end of file diff --git a/examples/dataset_table_workflow/utils.py b/examples/dataset_table_workflow/utils.py index d37e4bd..c467b9b 100644 --- a/examples/dataset_table_workflow/utils.py +++ b/examples/dataset_table_workflow/utils.py @@ -1,14 +1,15 @@ import os -from collections import defaultdict - import json -import matplotlib.pyplot as plt -import numpy as np -from loguru import logger import traceback +import numpy as np +import matplotlib.pyplot as plt +from collections import defaultdict + +from learnware.logger import get_module_logger +from config import * + +logger = get_module_logger("base_table", level="INFO") -from examples.dataset_table_workflow.config import * -from benchmarks.config import default_size_list class Recorder: def __init__(self, headers=["Mean", "Std Dev"], formats=["{:.2f}", "{:.2f}"]): @@ -79,7 +80,7 @@ def analyze_performance(user, recorders): logger.info(f"{user}, {user_id}, {mean_differences[user_id]}, {single_multi_diff}") -def plot_performance_curves(user, recorders, task="Hetero", n_labeled_list=default_size_list): +def plot_performance_curves(user, recorders, task, n_labeled_list): plt.figure(figsize=(10, 6)) for method, recorder in recorders.items(): @@ -106,4 +107,8 @@ def plot_performance_curves(user, recorders, task="Hetero", n_labeled_list=defau plt.title(f'Table {task} Limited Labeled Data') plt.legend() plt.tight_layout() - plt.savefig(os.path.join('./results/figs', f"{user}_labeled_{list(recorders.keys())}.png"), bbox_inches="tight", dpi=700) \ No newline at end of file + + root_path = os.path.abspath(os.path.join(__file__, "..")) + fig_path = os.path.join(root_path, "results", "figs") + os.makedirs(fig_path, exist_ok=True) + plt.savefig(os.path.join(fig_path, f"{user}_labeled_{list(recorders.keys())}.svg"), bbox_inches="tight", dpi=700) \ No newline at end of file diff --git a/examples/dataset_table_workflow/workflow.py b/examples/dataset_table_workflow/workflow.py new file mode 100644 index 0000000..3f0a319 --- /dev/null +++ b/examples/dataset_table_workflow/workflow.py @@ -0,0 +1,35 @@ +import fire + +from learnware.logger import get_module_logger +from homo import CorporacionDatasetWorkflow +from config import homo_table_benchmark_config + +logger = get_module_logger("base_table", level="INFO") + + +class TableDatasetWorkflow: + def unlabeled_homo_table_example(self): + workflow = CorporacionDatasetWorkflow( + benchmark_config=homo_table_benchmark_config, + name="easy", + rebuild=False + ) + workflow.unlabeled_homo_table_example() + + def labeled_homo_table_example(self): + workflow = CorporacionDatasetWorkflow( + benchmark_config=homo_table_benchmark_config, + name="easy", + rebuild=False + ) + workflow.labeled_homo_table_example() + + def cross_feat_eng_hetero_table_example(self): + pass + + def cross_task_hetero_table_example(self): + pass + + +if __name__ == "__main__": + fire.Fire(TableDatasetWorkflow) \ No newline at end of file