Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -24,13 +24,33 @@
###############################################################################
from contextlib import contextmanager
from typing import List, Optional, Tuple, Union
import inspect

# Third Party
from peft import PeftConfig, PeftModel, PeftType, get_peft_model
from peft.mapping import PEFT_TYPE_TO_CONFIG_MAPPING
from peft.peft_model import PEFT_TYPE_TO_MODEL_MAPPING

try:
# Third Party
from peft.peft_model import PEFT_TYPE_TO_MODEL_MAPPING
except ImportError:
# peft >= 0.19 renamed this mapping and provides no back-compat alias
from peft.peft_model import (
PEFT_TYPE_TO_TUNER_MAPPING as PEFT_TYPE_TO_MODEL_MAPPING,
)

# Third Party
from peft.tuners.lora import LoraConfig, LoraModel
from peft.tuners.lora.gptq import GPTQLoraLinear

# peft >= 0.19 made `config` a required positional argument of
# GPTQLoraLinear.__init__. On earlier versions there is no such parameter and
# the value was absorbed (and discarded) by **kwargs, so only pass it when the
# installed peft actually declares it.
_GPTQ_LORA_LINEAR_TAKES_CONFIG = (
"config" in inspect.signature(GPTQLoraLinear.__init__).parameters
)
# Third Party
import torch

# Local
Expand Down Expand Up @@ -68,9 +88,11 @@ def _create_new_module(
# to be installed
new_module = None
if isinstance(target, target_cls):
new_module = GPTQLoraLinear(
target, adapter_name, lora_config=lora_config, **kwargs
)
if _GPTQ_LORA_LINEAR_TAKES_CONFIG:
# peft >= 0.19: `config` is a required positional argument
new_module = GPTQLoraLinear(target, adapter_name, lora_config, **kwargs)
else:
new_module = GPTQLoraLinear(target, adapter_name, **kwargs)

# if module cannot be found, return None which results in a raise in the call-stack
return new_module
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -44,8 +44,8 @@ def register_foak_model_patch_rules(
gpt_bigcode,
granite,
granitemoe,
granitemoeshared,
granitemoehybrid,
granitemoeshared,
llama,
mistral,
mixtral,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -176,11 +176,18 @@ def _cluster_embeddings(
)

try:
from cuml import KMeans # pylint: disable=import-outside-toplevel
# Third Party
from cuml import KMeans # pylint: disable=import-outside-toplevel

print("Using GPU accelerated Kmeans")
except ImportError:
print("GPU accelerated KMeans is not avaialble. Falling back to CPU based KMeans")
from sklearn.cluster import KMeans # pylint: disable=import-outside-toplevel
print(
"GPU accelerated KMeans is not avaialble. Falling back to CPU based KMeans"
)
# Third Party
from sklearn.cluster import ( # pylint: disable=import-outside-toplevel
KMeans,
)

kwargs = {"n_init": 10}
kwargs.update(self.config.cluster_kwargs)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -213,7 +213,9 @@ def __init__(
self.output_dir = output_dir
if not os.path.exists(self.output_dir):
os.makedirs(self.output_dir)
self.log_file_path = os.path.join(self.output_dir, f"odm_rank_{self.rank}.jsonl")
self.log_file_path = os.path.join(
self.output_dir, f"odm_rank_{self.rank}.jsonl"
)
logger.info(
"Logs for online data mixing to be stored at {log_file_path}".format(
log_file_path=self.log_file_path
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -124,7 +124,10 @@ def compute_reward(

handlers = {
Reward.TRAIN_LOSS: lambda: _compute_train_loss_reward(
train_loss_history, last_sampled_category, current_category, total_categories
train_loss_history,
last_sampled_category,
current_category,
total_categories,
),
Reward.VALIDATION_LOSS: lambda: _compute_validation_loss_reward(
eval_loss_history, current_category, total_categories
Expand Down
6 changes: 5 additions & 1 deletion plugins/online-data-mixing/tests/test_compute_reward.py
Original file line number Diff line number Diff line change
Expand Up @@ -191,7 +191,11 @@ def _make_batch(batch_size, seq_length, vocab_size, offset=0):
% vocab_size
)
attention_mask = torch.ones(batch_size, seq_length, dtype=torch.long)
return {"input_ids": input_ids, "labels": input_ids, "attention_mask": attention_mask}
return {
"input_ids": input_ids,
"labels": input_ids,
"attention_mask": attention_mask,
}


def test_compute_reward_learnability():
Expand Down
4 changes: 3 additions & 1 deletion plugins/online-data-mixing/tests/test_online_data.py
Original file line number Diff line number Diff line change
Expand Up @@ -176,7 +176,9 @@ def reduce(self, x, reduction): # pylint: disable=unused-argument
# update_sampling_weights() moves eval batches to accelerator.device (or
# torch.device(0), i.e. cuda:0, if no accelerator is given). Pass a
# single-process CPU stub so this test doesn't require a GPU.
dataset.update_sampling_weights(model, accelerator=CPUAccelerator(), state=DummyState())
dataset.update_sampling_weights(
model, accelerator=CPUAccelerator(), state=DummyState()
)

assert dataset.log["rewards"], "expected rewards to be logged after update"
counts = list(dataset.log["count"])
Expand Down
27 changes: 11 additions & 16 deletions scripts/benchmarks/benchmark.py
Original file line number Diff line number Diff line change
Expand Up @@ -302,8 +302,8 @@ def build_args_from_products(products: List[Dict], defaults: Dict):
argument_list = ConfigUtils.convert_keyvalue_arguments_to_list(
combined_args
)
pdtbs = combined_args.get('per_device_train_batch_size')
grad_accum = combined_args.get('gradient_accumulation_steps')
pdtbs = combined_args.get("per_device_train_batch_size")
grad_accum = combined_args.get("gradient_accumulation_steps")
if pdtbs is None and grad_accum is not None:
if grad_accum > 1:
warnings.warn(
Expand Down Expand Up @@ -389,15 +389,11 @@ def __init__(self, scenario: Dict, acceleration_config_map: Dict = None) -> None
# if acceleration_config_map is None, then do not do mapping
if acceleration_config_map:

# - we allow k to be None to indicate we do not wish to
# - we allow k to be None to indicate we do not wish to
# set a config for that matrix entry. However, we do not
# check for multiple None's, so be careful.
val = [
(
acceleration_config_map[k]
if k is not None
else None
)
(acceleration_config_map[k] if k is not None else None)
for k in val
if k in acceleration_config_map or k is None
]
Expand Down Expand Up @@ -466,8 +462,8 @@ def is_completed(self):
# return complete only if no errors
# and is not a dry run
return (
not ERROR_MESSAGES in results and
results.get(DRY_RUN_MESSAGE, False) == False
not ERROR_MESSAGES in results
and results.get(DRY_RUN_MESSAGE, False) == False
)

def run(
Expand Down Expand Up @@ -719,12 +715,11 @@ def prepare_arguments(args, benchmark_dataset: BenchmarkDataset):
# build scenario matrix
scenario = ScenarioMatrix(scenario_config, acceleration_config_map)

if (
not args.run_only_scenarios
and scenario.slow
):
if not args.run_only_scenarios and scenario.slow:
# unfiltered runs omit all "slow" marked scenarios
print(f"Skipping slow scenario '{_scn_name}' beacuse run_only_scenarios=None.")
print(
f"Skipping slow scenario '{_scn_name}' beacuse run_only_scenarios=None."
)
continue

scenario_matrices, scenario_constants = (
Expand All @@ -736,7 +731,7 @@ def prepare_arguments(args, benchmark_dataset: BenchmarkDataset):
scn_factor *= len(v)

# scenario-specific constants should overwrite any similar values in defaults
defaults = {k:v for k, v in defaults.items() if k not in scenario_constants}
defaults = {k: v for k, v in defaults.items() if k not in scenario_constants}
# update defaults with scenario constants
constants = {**defaults, **scenario_constants}
# Remove any empty variables and combine matrices to dictionary to cartesian product on
Expand Down
43 changes: 22 additions & 21 deletions scripts/benchmarks/compare_with_reference.py
Original file line number Diff line number Diff line change
Expand Up @@ -64,21 +64,19 @@ def compare_results(df, ref, plot_columns, threshold_ratio=0.1):
df_series = df[column].fillna(0)
# Extract outliers base on some threshold % difference on referance
cmp = ref_series.to_frame()
cmp['metric'] = column
cmp = cmp.join(df_series.to_frame(), lsuffix='_ref')
cmp = cmp.rename(columns={f'{column}_ref': 'reference', column: 'new'})
cmp['ds'] = cmp.apply(
lambda x: (
abs(x.reference - x.new) / (x.reference + 1e-9)
), axis=1
cmp["metric"] = column
cmp = cmp.join(df_series.to_frame(), lsuffix="_ref")
cmp = cmp.rename(columns={f"{column}_ref": "reference", column: "new"})
cmp["ds"] = cmp.apply(
lambda x: (abs(x.reference - x.new) / (x.reference + 1e-9)), axis=1
)
outliers = cmp[cmp.ds > threshold_ratio]
outliers = outliers.drop('ds', axis=1)
outliers = outliers.drop("ds", axis=1)

plot_chart(
ax,
cmp['reference'],
cmp['new'],
cmp["reference"],
cmp["new"],
title=f"Metric: {column}",
xlabel="Reference",
ylabel="New",
Expand Down Expand Up @@ -112,30 +110,33 @@ def main(

# NOTE: this is a bit of a hack, if the new bench is a smaller bench, then we
# supplement the data from the raw summary
new_benchmark_filepath = os.path.join(result_dir, BENCHMARK_FILENAME)
new_benchmark_filepath = os.path.join(result_dir, BENCHMARK_FILENAME)
try:
df, args_df = read_df(new_benchmark_filepath, indices, plot_columns)
df, args_df = read_df(new_benchmark_filepath, indices, plot_columns)
except KeyError:
raw_filepath = os.path.join(result_dir, RAW_FILENAME)
print (
print(
f"New '{new_benchmark_filepath}' is probably a partial bench. Supplementing "
f"missing columns from raw data '{raw_filepath}'."
)
df2 = pd.read_csv(new_benchmark_filepath)
df = pd.read_csv(raw_filepath)
df, args_df = read_df(
pd.concat([df, df2[[x for x in df2.columns if x not in df.columns]]], axis=1),
indices, plot_columns
pd.concat(
[df, df2[[x for x in df2.columns if x not in df.columns]]], axis=1
),
indices,
plot_columns,
)

# Analyse between both sets of results and retrieve outliers
# - this has a side effect of plotting the charts
outliers_df, outliers, charts = compare_results(
df, ref, plot_columns, threshold_ratio=threshold_ratio
)
# this logic is brittle and will not hold if new benchmark is not
# this logic is brittle and will not hold if new benchmark is not
# of the exact same format as the reference benchmark,
# so put a try-catch.
# so put a try-catch.
try:
# Find arguments that are different between ref and new
# to highlight as possible cause of anomaly
Expand All @@ -147,8 +148,8 @@ def main(
outliers_df = outliers_df.set_index(indices).merge(
diff, left_index=True, right_index=True
)
except ValueError:
print (
except ValueError:
print(
f"New '{new_benchmark_filepath}' is probably a partial bench. So unable"
"to properly compare if the arguments are consistent with old bench."
)
Expand Down Expand Up @@ -186,8 +187,8 @@ def main(
)

parser.add_argument(
"--plot_columns",
default=DEFAULT_PLOT_COLUMNS,
"--plot_columns",
default=DEFAULT_PLOT_COLUMNS,
nargs="+",
help="list of metric names in benchmark results to analyze visually",
)
Expand Down
14 changes: 8 additions & 6 deletions scripts/benchmarks/data_processing.py
Original file line number Diff line number Diff line change
Expand Up @@ -60,14 +60,15 @@ def _build_data_formatting_func(
if tokenize and response_template is not None:
loss_masking = instruction_mask_loss(tokenizer, response_template)
elif tokenize and response_template is None:
assert response_field is not None, \
"response_field must be specified if tokenize=True and response_template=None."
assert (
response_field is not None
), "response_field must be specified if tokenize=True and response_template=None."

def _format(example):
# `nonlocal` is used because the format_fn will be passed to dataset.map and
# `loss_masking` needs to be bounded by `nonlocal` otherwise the spawned
# processes will have no reference to it
nonlocal loss_masking
nonlocal loss_masking
formatted_and_maybe_tokenized = tokenizer.apply_chat_template(
[example], tokenize=tokenize
)
Expand All @@ -82,10 +83,11 @@ def _format(example):
"in the chat template."
)
# NOTE: in this case not handling attention mask
response = tokenizer(example[response_field])['input_ids']
response = tokenizer(example[response_field])["input_ids"]
return {
key: formatted_and_maybe_tokenized + response,
'labels': [ ignore_index ] * len(formatted_and_maybe_tokenized) + response
"labels": [ignore_index] * len(formatted_and_maybe_tokenized)
+ response,
}

if not loss_masking:
Expand Down Expand Up @@ -221,4 +223,4 @@ def collate_example(example):
# flatten the additional dim
return {k: v.view(-1) for k, v in collated_example.items()}

return collate_example
return collate_example
2 changes: 1 addition & 1 deletion scripts/benchmarks/display_bench_results.py
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,7 @@ def main(
try:
df = df.drop("output_dir", axis=1)
except KeyError:
pass # output_dir not found
pass # output_dir not found

df.reindex(sorted(df.columns), axis=1).to_csv(output_filename, index=False)
print("***************** Report Created ******************")
Expand Down
Loading
Loading