Browse Source

[MNT] init homo table workflow

tags/v0.3.2
Gene 2 years ago
parent
commit
5729af6b20
9 changed files with 454 additions and 456 deletions
  1. +0
    -199
      examples/base.py
  2. +105
    -0
      examples/dataset_table_workflow/base.py
  3. +87
    -19
      examples/dataset_table_workflow/config.py
  4. +152
    -181
      examples/dataset_table_workflow/homo.py
  5. +0
    -34
      examples/dataset_table_workflow/main.py
  6. +13
    -14
      examples/dataset_table_workflow/methods.py
  7. +48
    -0
      examples/dataset_table_workflow/train.py
  8. +14
    -9
      examples/dataset_table_workflow/utils.py
  9. +35
    -0
      examples/dataset_table_workflow/workflow.py

+ 0
- 199
examples/base.py View File

@@ -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


+ 105
- 0
examples/dataset_table_workflow/base.py View File

@@ -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")

+ 87
- 19
examples/dataset_table_workflow/config.py View File

@@ -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",
)

+ 152
- 181
examples/dataset_table_workflow/homo.py View File

@@ -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)
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)

+ 0
- 34
examples/dataset_table_workflow/main.py View File

@@ -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()

+ 13
- 14
examples/dataset_table_workflow/methods.py View File

@@ -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


+ 48
- 0
examples/dataset_table_workflow/train.py View File

@@ -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)

+ 14
- 9
examples/dataset_table_workflow/utils.py View File

@@ -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)
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)

+ 35
- 0
examples/dataset_table_workflow/workflow.py View File

@@ -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)

Loading…
Cancel
Save