diff --git a/RecommenderSystems/dlrm/README.md b/RecommenderSystems/dlrm/README.md index f01a6dfff..88d03fc69 100644 --- a/RecommenderSystems/dlrm/README.md +++ b/RecommenderSystems/dlrm/README.md @@ -5,11 +5,19 @@ ## Directory description ``` . -|-- tools - |-- criteo1t_parquet.py # Read Criteo1T data and export it as parquet data format -|-- dlrm_train_eval.py # OneFlow DLRM training and evaluation scripts with OneEmbedding module -|-- requirements.txt # python package configuration file -└── README.md # Documentation +├── tools +│   ├── criteo1t_parquet.py # make criteo terabyte dataset for OneFlow DLRM by python +│   ├── criteo1t_parquet.scala # make criteo terabyte dataset for OneFlow DLRM by spark +│   ├── criteo1t_parquet_int32.scala # make criteo terabyte dataset for OneFlow DLRM by spark, int32 for sparse features +│   ├── launch_spark.sh # launch spark +│   ├── parquet_to_raw.py # convert dataset from parquet to raw format +│   └── split_day_23.sh # split day_23 to test and validation set +├── dlrm_train_eval.py # OneFlow DLRM training and evaluation scripts with OneEmbedding module +├── dlrm_benchmark_a100.py # OneFlow DLRM benchmark training and evaluation scripts +├── train_dlrm_benchmark.sh # OneFlow DLRM benchmark AMP training command +├── train_dlrm_benchmark_fp32.sh # OneFlow DLRM benchmark FP32 training command +├── requirements.txt # python package configuration file +└── README.md # Documentation ``` ## Arguments description @@ -101,3 +109,46 @@ python3 -m oneflow.distributed.launch \ --data_dir /path/to/dlrm_parquet \ --persistent_path /path/to/persistent ``` + +## Run OneFlow DLRM benchmark +1. make dlrm raw format dataset (sparse feature dtype = int32) + - split day_23 to test.csv and val.csv, in criteo terabyte dataset directory where extracted day_0 to day_23 files located +``` +head -n 89137319 day_23 > test.csv +tail -n +89137320 day_23 > val.csv +``` + - launch spark shell in "RecommenderSystems/dlrm/tools" directory: +``` +export SPARK_LOCAL_DIRS=/path/to/tmp_spark +spark-shell \ + --master "local[*]" \ + --conf spark.driver.maxResultSize=0 \ + --driver-memory 360G +``` + - load scala file in spark-shell, and execute `makeDlrmDatasetInt32` +``` +:load criteo1t_parquet_int32.scala +makeDlrmDatasetInt32("/path/to/criteo1t_raw", "/path/to/dlrm_parquet_int32") +``` + - convert parquet dataset to oneflow raw format +``` +python parquet_to_raw.py \ + --input_dir=/path/to/dlrm_parquet_int32 \ + --output_dir=/path/to/criteo1t_oneflow_raw +``` + +2. train OneFlow DLRM benchmark in AMP mode +``` +./train_dlrm_benchmark.sh /path/to/data_dir +``` +note: `train_dlrm_benchmark.sh` takes 3 arguments: +- $1 is data_dir pointing to criteo1t_oneflow_raw dataset +- $2 is number of GPUs, default is `8` +- $3 is to enable oneflow raw reader direct io or not, default is `0`, set `1` to enable direct io + +3. or train OneFlow DLRM benchmark in FP32 mode + +``` +./train_dlrm_benchmark_fp32.sh /path/to/data_dir +``` + diff --git a/RecommenderSystems/dlrm/dlrm_benchmark_a100.py b/RecommenderSystems/dlrm/dlrm_benchmark_a100.py new file mode 100644 index 000000000..5ab04c6f5 --- /dev/null +++ b/RecommenderSystems/dlrm/dlrm_benchmark_a100.py @@ -0,0 +1,545 @@ +""" +Copyright 2020 The OneFlow Authors. All rights reserved. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +""" +import argparse +import os +import sys +import glob +import time +import numpy as np +import psutil +import warnings +import oneflow as flow +import oneflow.nn as nn + + +sys.path.append(os.path.abspath(os.path.join(os.path.dirname(__file__), os.path.pardir))) + + +def get_args(print_args=True): + def int_list(x): + return list(map(int, x.split(","))) + + def str_list(x): + return list(map(str, x.split(","))) + + parser = argparse.ArgumentParser() + + parser.add_argument("--disable_fusedmlp", action="store_true", help="disable fused MLP or not") + parser.add_argument("--embedding_vec_size", type=int, default=128) + parser.add_argument( + "--one_embedding_key_type", + type=str, + default="int32", + help="OneEmbedding key type: int32, int64", + ) + parser.add_argument("--bottom_mlp", type=int_list, default="512,256,128") + parser.add_argument("--top_mlp", type=int_list, default="1024,1024,512,256") + parser.add_argument( + "--disable_interaction_padding", + action="store_true", + help="disable interaction padding or not", + ) + parser.add_argument( + "--disable_dense_input_padding", + action="store_true", + help="disable dense input padding or not", + ) + parser.add_argument( + "--interaction_itself", action="store_true", help="interaction itself or not" + ) + parser.add_argument("--model_load_dir", type=str, default=None) + parser.add_argument("--model_save_dir", type=str, default=None) + parser.add_argument( + "--save_initial_model", action="store_true", help="save initial model parameters or not.", + ) + parser.add_argument( + "--save_model_after_each_eval", action="store_true", help="save model after each eval.", + ) + parser.add_argument("--data_dir", type=str, required=True) + parser.add_argument("--eval_batches", type=int, default=1612, help="number of eval batches") + parser.add_argument("--eval_batch_size", type=int, default=55296) + parser.add_argument("--eval_interval", type=int, default=100000) + parser.add_argument("--eval_steps", type=int_list, default="1000000") + parser.add_argument("--train_batch_size", type=int, default=55296) + parser.add_argument("--learning_rate", type=float, default=24) + parser.add_argument("--warmup_batches", type=int, default=2750) + parser.add_argument("--decay_batches", type=int, default=15406) + parser.add_argument("--decay_start", type=int, default=64163) + parser.add_argument("--train_batches", type=int, default=75869) + parser.add_argument("--loss_print_interval", type=int, default=1000) + parser.add_argument( + "--table_size_array", + type=int_list, + default="39884406,39043,17289,7420,20263,3,7120,1543,63,38532951,2953546,403346,10,2208,11938,155,4,976,14,39979771,25641295,39664984,585935,12972,108,36", + help="Embedding table size array for sparse fields", + ) + parser.add_argument( + "--persistent_path", type=str, default="persistent", help="path for persistent kv store", + ) + parser.add_argument("--store_type", type=str, default="device_mem") + parser.add_argument("--cache_memory_budget_mb", type=int, default=8192) + parser.add_argument("--amp", action="store_true", help="Run model with amp") + parser.add_argument( + "--disable_split_allreduce", action="store_true", help="disable split bottom and top allreduce" + ) + parser.add_argument("--loss_scale_policy", type=str, default="static", help="static or dynamic") + + args = parser.parse_args() + + if print_args and flow.env.get_rank() == 0: + _print_args(args) + return args + + +def _print_args(args): + """Print arguments.""" + print("------------------------ arguments ------------------------", flush=True) + str_list = [] + for arg in vars(args): + dots = "." * (48 - len(arg)) + str_list.append(" {} {} {}".format(arg, dots, getattr(args, arg))) + for arg in sorted(str_list, key=lambda x: x.lower()): + print(arg, flush=True) + print("-------------------- end of arguments ---------------------", flush=True) + + +num_dense_fields = 13 +num_sparse_fields = 26 + + +class Dense(nn.Module): + def __init__(self, in_features: int, out_features: int, relu=True) -> None: + super(Dense, self).__init__() + self.features = ( + nn.Sequential(nn.Linear(in_features, out_features), nn.ReLU(inplace=True)) + if relu + else nn.Linear(in_features, out_features) + ) + + def forward(self, x: flow.Tensor) -> flow.Tensor: + return self.features(x) + + +class MLP(nn.Module): + def __init__( + self, in_features: int, hidden_units, skip_final_activation=False, fused=True + ) -> None: + super(MLP, self).__init__() + if fused: + self.linear_layers = nn.FusedMLP( + in_features, + hidden_units[:-1], + hidden_units[-1], + skip_final_activation=skip_final_activation, + ) + else: + units = [in_features] + hidden_units + num_layers = len(hidden_units) + denses = [ + Dense(units[i], units[i + 1], not skip_final_activation or (i + 1) < num_layers) + for i in range(num_layers) + ] + self.linear_layers = nn.Sequential(*denses) + + for name, param in self.linear_layers.named_parameters(): + if "weight" in name: + nn.init.normal_(param, 0.0, np.sqrt(2 / sum(param.shape))) + elif "bias" in name: + nn.init.normal_(param, 0.0, np.sqrt(1 / param.shape[0])) + + def forward(self, x: flow.Tensor) -> flow.Tensor: + return self.linear_layers(x) + + +class Interaction(nn.Module): + def __init__( + self, + dense_feature_size, + num_embedding_fields, + interaction_itself=False, + interaction_padding=True, + ): + super(Interaction, self).__init__() + self.interaction_itself = interaction_itself + n_cols = num_embedding_fields + 2 if self.interaction_itself else num_embedding_fields + 1 + output_size = dense_feature_size + sum(range(n_cols)) + self.output_size = ((output_size + 8 - 1) // 8 * 8) if interaction_padding else output_size + self.output_padding = self.output_size - output_size + + def forward(self, x: flow.Tensor, y: flow.Tensor) -> flow.Tensor: + (bsz, d) = x.shape + return flow._C.fused_dot_feature_interaction( + [x.view(bsz, 1, d), y], + output_concat=x, + self_interaction=self.interaction_itself, + output_padding=self.output_padding, + ) + + +class OneEmbedding(nn.Module): + def __init__( + self, + embedding_vec_size, + persistent_path, + table_size_array, + store_type, + cache_memory_budget_mb, + key_type, + ): + assert table_size_array is not None + vocab_size = sum(table_size_array) + + scales = np.sqrt(1 / np.array(table_size_array)) + tables = [ + flow.one_embedding.make_table_options( + flow.one_embedding.make_uniform_initializer(low=-scale, high=scale) + ) + for scale in scales + ] + if store_type == "device_mem": + store_options = flow.one_embedding.make_device_mem_store_options( + persistent_path=persistent_path, capacity=vocab_size + ) + elif store_type == "cached_host_mem": + assert cache_memory_budget_mb > 0 + store_options = flow.one_embedding.make_cached_host_mem_store_options( + cache_budget_mb=cache_memory_budget_mb, + persistent_path=persistent_path, + capacity=vocab_size, + ) + elif store_type == "cached_ssd": + assert cache_memory_budget_mb > 0 + store_options = flow.one_embedding.make_cached_ssd_store_options( + cache_budget_mb=cache_memory_budget_mb, + persistent_path=persistent_path, + capacity=vocab_size, + ) + else: + raise NotImplementedError("not support", store_type) + + super(OneEmbedding, self).__init__() + self.one_embedding = flow.one_embedding.MultiTableEmbedding( + "sparse_embedding", + embedding_dim=embedding_vec_size, + dtype=flow.float, + key_type=getattr(flow, key_type), + tables=tables, + store_options=store_options, + ) + + def forward(self, ids): + return self.one_embedding.forward(ids) + + +class DLRMModule(nn.Module): + def __init__( + self, + embedding_vec_size=128, + bottom_mlp=[512, 256, 128], + top_mlp=[1024, 1024, 512, 256], + use_fusedmlp=True, + persistent_path=None, + table_size_array=None, + one_embedding_store_type="cached_host_mem", + one_embedding_key_type="int64", + cache_memory_budget_mb=8192, + interaction_itself=True, + interaction_padding=True, + dense_input_padding=True, + ): + super(DLRMModule, self).__init__() + assert ( + embedding_vec_size == bottom_mlp[-1] + ), "Embedding vector size must equle to bottom MLP output size" + self.num_dense_fields = ( + ((num_dense_fields + 8 - 1) // 8 * 8) if dense_input_padding else num_dense_fields + ) + self.pad = ( + [0, self.num_dense_fields - num_dense_fields] + if self.num_dense_fields > num_dense_fields + else None + ) + + self.bottom_mlp = MLP(self.num_dense_fields, bottom_mlp, fused=use_fusedmlp) + self.embedding = OneEmbedding( + embedding_vec_size, + persistent_path, + table_size_array, + one_embedding_store_type, + cache_memory_budget_mb, + one_embedding_key_type, + ) + self.interaction = Interaction( + bottom_mlp[-1], + num_sparse_fields, + interaction_itself, + interaction_padding=interaction_padding, + ) + self.top_mlp = MLP( + self.interaction.output_size, + top_mlp + [1], + skip_final_activation=True, + fused=use_fusedmlp, + ) + + def forward(self, dense_fields, sparse_fields) -> flow.Tensor: + if self.pad: + dense_fields = flow.nn.functional.pad(dense_fields, self.pad, "constant") + # dense_fields = flow.log(dense_fields + 1.0) + dense_fields = self.bottom_mlp(dense_fields) + embedding = self.embedding(sparse_fields) + features = self.interaction(dense_fields, embedding) + return self.top_mlp(features) + + +def make_dlrm_module(args): + model = DLRMModule( + embedding_vec_size=args.embedding_vec_size, + bottom_mlp=args.bottom_mlp, + top_mlp=args.top_mlp, + use_fusedmlp=not args.disable_fusedmlp, + persistent_path=args.persistent_path, + table_size_array=args.table_size_array, + one_embedding_store_type=args.store_type, + one_embedding_key_type=args.one_embedding_key_type, + cache_memory_budget_mb=args.cache_memory_budget_mb, + interaction_itself=args.interaction_itself, + interaction_padding=not args.disable_interaction_padding, + dense_input_padding=not args.disable_dense_input_padding, + ) + return model + + +def make_raw_dataloader(data_path, batch_size, shuffle=True): + def make_reader(data_file, length, dtype): + return flow.nn.RawReader( + [data_file], + (length,), + dtype, + batch_size, + random_shuffle=shuffle, + random_seed=1234, + placement=flow.env.all_device_placement("cpu"), + sbp=flow.sbp.split(0), + ) + + label_loader = make_reader(f"{data_path}/label.bin", 1, flow.float32) + dense_loader = make_reader(f"{data_path}/dense_norm.bin", num_dense_fields, flow.float32) + sparse_loader = make_reader(f"{data_path}/sparse.bin", num_sparse_fields, flow.int32) + + return label_loader, dense_loader, sparse_loader + + +def make_lr_scheduler(args, optimizer): + warmup_lr = flow.optim.lr_scheduler.LinearLR( + optimizer, start_factor=0, total_iters=args.warmup_batches, + ) + poly_decay_lr = flow.optim.lr_scheduler.PolynomialLR( + optimizer, decay_batch=args.decay_batches, end_learning_rate=0, power=2.0, cycle=False, + ) + sequential_lr = flow.optim.lr_scheduler.SequentialLR( + optimizer=optimizer, + schedulers=[warmup_lr, poly_decay_lr], + milestones=[args.decay_start], + interval_rescaling=True, + ) + return sequential_lr + + +class DLRMValGraph(flow.nn.Graph): + def __init__(self, dlrm_module, data_dir, batch_size, amp=False): + super(DLRMValGraph, self).__init__() + self.module = dlrm_module + self.label_loader, self.dense_loader, self.sparse_loader = make_raw_dataloader( + f"{data_dir}/test", batch_size, shuffle=False + ) + + if amp: + self.config.enable_amp(True) + + def build(self): + labels = self.label_loader() + dense_fields = self.dense_loader() + sparse_fields = self.sparse_loader() + predicts = self.module(dense_fields.to("cuda"), sparse_fields.to("cuda")) + return labels, predicts.sigmoid() + + +class DLRMTrainGraph(flow.nn.Graph): + def __init__( + self, + dlrm_module, + loss, + optimizer, + data_dir, + batch_size, + lr_scheduler=None, + grad_scaler=None, + amp=False, + ): + super(DLRMTrainGraph, self).__init__() + self.module = dlrm_module + self.loss = loss + self.add_optimizer(optimizer, lr_sch=lr_scheduler) + self.label_loader, self.dense_loader, self.sparse_loader = make_raw_dataloader( + f"{data_dir}/train", batch_size + ) + self.config.allow_fuse_model_update_ops(True) + self.config.allow_fuse_add_to_output(True) + self.config.allow_fuse_cast_scale(True) + if amp: + self.config.enable_amp(True) + self.set_grad_scaler(grad_scaler) + + def build(self): + labels = self.label_loader() + dense_fields = self.dense_loader() + sparse_fields = self.sparse_loader() + logits = self.module(dense_fields.to("cuda"), sparse_fields.to("cuda")) + loss = self.loss(logits, labels.to("cuda")) + loss.backward() + return loss.to("cpu") + + +def train(args): + rank = flow.env.get_rank() + + dlrm_module = make_dlrm_module(args) + dlrm_module.to_global(flow.env.all_device_placement("cuda"), flow.sbp.broadcast) + + if args.model_load_dir: + print(f"Loading model from {args.model_load_dir}") + state_dict = flow.load(args.model_load_dir, global_src_rank=0) + dlrm_module.load_state_dict(state_dict, strict=False) + + def save_model(subdir): + if not args.model_save_dir: + return + save_path = os.path.join(args.model_save_dir, subdir) + if rank == 0: + print(f"Saving model to {save_path}") + state_dict = dlrm_module.state_dict() + flow.save(state_dict, save_path, global_dst_rank=0) + + if args.save_initial_model: + save_model("initial_checkpoint") + + opt = flow.optim.SGD(dlrm_module.parameters(), lr=args.learning_rate) + lr_scheduler = make_lr_scheduler(args, opt) + loss = flow.nn.BCEWithLogitsLoss(reduction="mean").to("cuda") + + if args.loss_scale_policy == "static": + grad_scaler = flow.amp.StaticGradScaler(1024) + else: + grad_scaler = flow.amp.GradScaler( + init_scale=1073741824, growth_factor=2.0, backoff_factor=0.5, growth_interval=2000, + ) + + eval_graph = DLRMValGraph(dlrm_module, args.data_dir, args.eval_batch_size, args.amp) + train_graph = DLRMTrainGraph( + dlrm_module, + loss, + opt, + args.data_dir, + args.train_batch_size, + lr_scheduler, + grad_scaler, + args.amp, + ) + + dlrm_module.train() + step, last_step, last_time = -1, 0, time.time() + for step in range(1, args.train_batches + 1): + loss = train_graph() + if step % args.loss_print_interval == 0: + loss = loss.numpy() + if rank == 0: + latency = (time.time() - last_time) / (step - last_step) + throughput = args.train_batch_size / latency + last_step, last_time = step, time.time() + print( + f"Rank[{rank}], Step {step}, Loss {loss:0.4f}, Latency " + + f"{(latency * 1000):0.3f} ms, Throughput {throughput:0.1f}, {last_time}" + ) + if np.isnan(loss): + exit(1) + + if (args.eval_interval > 0 and step % args.eval_interval == 0) or (step in args.eval_steps): + auc = eval(eval_graph, step) + if args.save_model_after_each_eval: + save_model(f"step_{step}_val_auc_{auc:0.5f}") + dlrm_module.train() + last_time = time.time() + + if args.eval_interval > 0 and step % args.eval_interval != 0: + auc = eval(eval_graph, step) + if args.save_model_after_each_eval: + save_model(f"step_{step}_val_auc_{auc:0.5f}") + + +def tensor_list_to_local(tensors): + return ( + flow.cat(tensors, dim=0) + .to_global(placement=flow.env.all_device_placement("cpu"), sbp=flow.sbp.split(0)) + .to_global(sbp=flow.sbp.broadcast()) + .to_local() + ) + + +def eval(eval_graph, cur_step=0): + eval_graph.module.eval() + labels, preds = [], [] + eval_start_time = time.time() + for i in range(args.eval_batches): + label, pred = eval_graph() + labels.append(label.to_local()) + preds.append(pred.to_local()) + + labels = tensor_list_to_local(labels) + preds = tensor_list_to_local(preds) + + flow.comm.barrier() + eval_time = time.time() - eval_start_time + + rank = flow.env.get_rank() + auc = 0 + if rank == 0: + auc_start_time = time.time() + auc = flow.roc_auc_score(labels, preds).numpy()[0] + auc_time = time.time() - auc_start_time + host_mem_mb = psutil.Process().memory_info().rss // (1024 * 1024) + stream = os.popen("nvidia-smi --query-gpu=memory.used --format=csv") + device_mem_str = stream.read().split("\n")[rank + 1] + + strtime = time.strftime("%Y-%m-%d %H:%M:%S") + print( + f"Rank[{rank}], Step {cur_step}, AUC {auc:0.5f}, Eval_time {eval_time:0.2f} s, " + + f"AUC_time {auc_time:0.2f} s, Eval_samples {labels.shape[0]}, " + + f"GPU_Memory {device_mem_str}, Host_Memory {host_mem_mb} MiB, {strtime}" + ) + + return auc + + +if __name__ == "__main__": + os.system(sys.executable + " -m oneflow --doctor") + #os.system("env") + flow.boxing.nccl.enable_all_to_all(True) + args = get_args() + if not args.disable_split_allreduce: + flow.boxing.nccl.set_fusion_max_ops_num(10) + + train(args) diff --git a/RecommenderSystems/dlrm/docker/Dockerfile b/RecommenderSystems/dlrm/docker/Dockerfile new file mode 100644 index 000000000..1a7160174 --- /dev/null +++ b/RecommenderSystems/dlrm/docker/Dockerfile @@ -0,0 +1,7 @@ +FROM nvcr.io/nvidia/pytorch:22.10-py3 +RUN python3 -m pip config set global.index-url https://pypi.tuna.tsinghua.edu.cn/simple +RUN apt update && DEBIAN_FRONTEND=noninteractive apt install -y libopenblas-dev nasm g++ gcc python3-pip cmake autoconf libtool libjemalloc2 google-perftools +RUN python3 -m pip install --pre oneflow -f https://staging.oneflow.info/commit/65fd4d99ad435baefe500e6ea9aaebcbb36c905e/cu116 psutil petastorm +RUN git clone --branch dlrm_benchmark_test --depth 1 https://github.com/Oneflow-Inc/models.git /models +RUN cd /models/RecommenderSystems/dlrm +ENV LD_PRELOAD /usr/lib/x86_64-linux-gnu/libjemalloc.so.2 diff --git a/RecommenderSystems/dlrm/docker/Dockerfile_make b/RecommenderSystems/dlrm/docker/Dockerfile_make new file mode 100644 index 000000000..c41ccf1fc --- /dev/null +++ b/RecommenderSystems/dlrm/docker/Dockerfile_make @@ -0,0 +1,14 @@ +FROM nvcr.io/nvidia/pytorch:22.10-py3 +RUN python3 -m pip config set global.index-url https://pypi.tuna.tsinghua.edu.cn/simple +RUN apt update && DEBIAN_FRONTEND=noninteractive apt install -y libopenblas-dev nasm g++ gcc python3-pip cmake autoconf libtool libjemalloc2 google-perftools +ENV NO_CACHE 9 +RUN git clone https://github.com/Oneflow-Inc/oneflow /oneflow && cd /oneflow && git checkout bbe0b16 +RUN python3 -m pip install -r /oneflow/dev-requirements.txt psutil petastorm +RUN mkdir /oneflow/build +RUN cd /oneflow/build +RUN cd /oneflow/build && cmake -DCMAKE_CUDA_ARCHITECTURES="80-real" -DCUDA_STATIC=OFF -DWITH_MLIR=YES -DWITH_MLIR_CUDA_CODEGEN=YES .. -C ../cmake/caches/cn/cuda.cmake +RUN cd /oneflow/build && make -j48 +ENV PYTHONPATH /oneflow/python +RUN git clone --branch dlrm_benchmark_test --depth 1 https://github.com/Oneflow-Inc/models.git /models +RUN cd /models/RecommenderSystems/dlrm +ENV LD_PRELOAD /usr/lib/x86_64-linux-gnu/libjemalloc.so.2 diff --git a/RecommenderSystems/dlrm/docker/build.sh b/RecommenderSystems/dlrm/docker/build.sh new file mode 100755 index 000000000..3a5f00d90 --- /dev/null +++ b/RecommenderSystems/dlrm/docker/build.sh @@ -0,0 +1,3 @@ +docker build \ + --rm \ + -t oneflow-dlrm:0.1 -f Dockerfile . diff --git a/RecommenderSystems/dlrm/docker/launch.sh b/RecommenderSystems/dlrm/docker/launch.sh new file mode 100755 index 000000000..c0a36b508 --- /dev/null +++ b/RecommenderSystems/dlrm/docker/launch.sh @@ -0,0 +1,7 @@ +docker run -it --rm --runtime=nvidia --privileged \ + --network host --gpus=all \ + --ipc=host \ + -v /RAID0:/RAID0 \ + -w /models/RecommenderSystems/dlrm \ + oneflow-dlrm:0.1 \ + bash diff --git a/RecommenderSystems/dlrm/tools/criteo1t_parquet.scala b/RecommenderSystems/dlrm/tools/criteo1t_parquet.scala index 0dd15e6c3..5a4d78e69 100644 --- a/RecommenderSystems/dlrm/tools/criteo1t_parquet.scala +++ b/RecommenderSystems/dlrm/tools/criteo1t_parquet.scala @@ -6,7 +6,6 @@ def makeDlrmDataset(srcDir: String, dstDir:String, tmpDir:String, modIdx:Long = val integer_names = Seq("label") ++ dense_names val col_names = integer_names ++ categorical_names - val day_23 = s"${srcDir}/day_23" val test_csv = s"${tmpDir}/test.csv" val val_csv = s"${tmpDir}/val.csv" diff --git a/RecommenderSystems/dlrm/tools/criteo1t_parquet_int32.scala b/RecommenderSystems/dlrm/tools/criteo1t_parquet_int32.scala new file mode 100644 index 000000000..7ab2501c0 --- /dev/null +++ b/RecommenderSystems/dlrm/tools/criteo1t_parquet_int32.scala @@ -0,0 +1,30 @@ +import org.apache.spark.sql.functions.udf + +def makeDlrmDatasetInt32(srcDir: String, dstDir:String) = { + val categorical_names = (1 to 26).map{id=>s"C$id"} + val dense_names = (1 to 13).map{id=>s"I$id"} + val integer_names = Seq("label") ++ dense_names + val col_names = integer_names ++ categorical_names + + val test_csv = s"${srcDir}/test.csv" + val val_csv = s"${srcDir}/val.csv" + + val make_label = udf((str:String) => str.toFloat) + val make_dense = udf((str:String) => if (str == null) 1 else str.toFloat + 1) + val make_sparse = udf((str:String, i:Int) => (if (str == null) 0 else Math.floorMod(Integer.parseUnsignedInt(str, 16).toInt, 40000000)) + i * 40000000) + val label_cols = Seq(make_label($"label").as("label")) + val dense_cols = 1.to(13).map{i=>make_dense(col(s"I$i")).as(s"I${i}")} + val sparse_cols = 1.to(26).map{i=>make_sparse(col(s"C$i"), lit(i)).as(s"C${i}")} + val cols = label_cols ++ dense_cols ++ sparse_cols + + spark.read.option("delimiter", "\t").csv(test_csv).toDF(col_names: _*).select(cols:_*).repartition(256).write.parquet(s"${dstDir}/test") + spark.read.option("delimiter", "\t").csv(val_csv).toDF(col_names: _*).select(cols:_*).repartition(256).write.parquet(s"${dstDir}/val") + + val day_files = 0.until(23).map{day=>s"${srcDir}/day_${day}"} + spark.read.option("delimiter", "\t").csv(day_files:_*).toDF(col_names: _*).select(cols:_*).orderBy(rand()).repartition(2560).write.parquet(s"${dstDir}/train") + + // print table size array + val df = spark.read.parquet(s"${dstDir}/*/*.parquet") + println(1.to(26).map{i=>df.select(s"C$i").as[Int].distinct.count}.mkString(",")) +} + diff --git a/RecommenderSystems/dlrm/tools/parquet_to_raw.py b/RecommenderSystems/dlrm/tools/parquet_to_raw.py new file mode 100644 index 000000000..f86c0ca4c --- /dev/null +++ b/RecommenderSystems/dlrm/tools/parquet_to_raw.py @@ -0,0 +1,47 @@ +import os +import glob +import time +import argparse +import numpy as np +from petastorm.reader import make_batch_reader + + +def parquet_to_raw(files, output_dir): + print(output_dir) + fields = ['label'] + fields += ["I{}".format(i + 1) for i in range(13)] + fields += ["C{}".format(i + 1) for i in range(26)] + with make_batch_reader(files, workers_count=1, shuffle_row_groups=False) as reader: + lf = open(f'{output_dir}/label.bin', 'wb') + sf = open(f'{output_dir}/sparse.bin', 'wb') + dnf = open(f'{output_dir}/dense_norm.bin', 'wb') + for rg in reader: + rgdict = rg._asdict() + rglist = [rgdict[field] for field in fields] + label = rglist[0] + lf.write(label.tobytes()) + dense = np.stack(rglist[1:14], axis=-1) + dense_norm = np.log(dense + 1.0) + dnf.write(dense_norm.tobytes()) + sparse = np.stack(rglist[14:40], axis=-1) + sf.write(sparse.tobytes()) + print(label.shape, dense_norm.shape, sparse.shape) + lf.close(), sf.close(), dnf.close() + + +if __name__ == "__main__": + def str_list(x): + return list(map(str, x.split(","))) + parser = argparse.ArgumentParser(description="convert parquet dataset to oneflow row") + parser.add_argument("--input_dir", type=str, required=True) + parser.add_argument("--output_dir", type=str, required=True) + parser.add_argument("--sub_sets", type=str_list, default="test,val,train") + args = parser.parse_args() + + for sub_set in args.sub_sets: + files = ['file://' + name for name in glob.glob(f'{args.input_dir}/{sub_set}/*.parquet')] + files.sort() + output_dir = os.path.join(args.output_dir, sub_set) + os.system(f"mkdir -p {output_dir}") + parquet_to_raw(files, output_dir) + diff --git a/RecommenderSystems/dlrm/train_dlrm_benchmark.sh b/RecommenderSystems/dlrm/train_dlrm_benchmark.sh new file mode 100755 index 000000000..31a491939 --- /dev/null +++ b/RecommenderSystems/dlrm/train_dlrm_benchmark.sh @@ -0,0 +1,28 @@ +data_dir=${1} +ngpus=${2:-8} +enable_direct_io=${3:-0} + +export ONEFLOW_RAW_READER_FORCE_DIRECT_IO=$enable_direct_io + +export NCCL_CHECKS_DISABLE=1 +export ONEFLOW_FUSE_MODEL_UPDATE_CAST=1 +export ONEFLOW_ENABLE_MULTI_TENSOR_MODEL_UPDATE=1 +export ONEFLOW_KERNEL_ENABLE_CUDA_GRAPH=1 +export ONEFLOW_EAGER_LOCAL_TO_GLOBAL_BALANCED_OVERRIDE=1 +export ONEFLOW_ONE_EMBEDDING_FUSED_MLP_ASYNC_GRAD=1 +export ONEFLOW_ONE_EMBEDDING_FUSE_EMBEDDING_INTERACTION=1 +export ONEFLOW_EP_CUDA_DEVICE_FLAGS=2 +export ONEFLOW_FUSE_BCE_REDUCE_MEAN_FW_BW=1 + +export ONEFLOW_ONE_EMBEDDING_ID_SHUFFLE_USE_P2P=1 +export ONEFLOW_ONE_EMBEDDING_EMBEDDING_SHUFFLE_USE_P2P=1 + +python3 -m oneflow.distributed.launch \ + --nproc_per_node $ngpus \ + --nnodes 1 \ + --node_rank 0 \ + --master_addr 127.0.0.1 \ + dlrm_benchmark_a100.py \ + --data_dir $data_dir \ + --amp + diff --git a/RecommenderSystems/dlrm/train_dlrm_benchmark_fp32.sh b/RecommenderSystems/dlrm/train_dlrm_benchmark_fp32.sh new file mode 100755 index 000000000..01385e196 --- /dev/null +++ b/RecommenderSystems/dlrm/train_dlrm_benchmark_fp32.sh @@ -0,0 +1,24 @@ +data_dir=${1} +ngpus=${2:-8} +enable_direct_io=${3:-0} + +export ONEFLOW_RAW_READER_FORCE_DIRECT_IO=$enable_direct_io + +export NCCL_CHECKS_DISABLE=1 +export ONEFLOW_FUSE_MODEL_UPDATE_CAST=1 +export ONEFLOW_ENABLE_MULTI_TENSOR_MODEL_UPDATE=1 +export ONEFLOW_KERNEL_ENABLE_CUDA_GRAPH=1 +export ONEFLOW_EAGER_LOCAL_TO_GLOBAL_BALANCED_OVERRIDE=1 +export ONEFLOW_ONE_EMBEDDING_FUSED_MLP_ASYNC_GRAD=1 +export ONEFLOW_ONE_EMBEDDING_FUSE_EMBEDDING_INTERACTION=1 +export ONEFLOW_EP_CUDA_DEVICE_FLAGS=2 +export ONEFLOW_FUSE_BCE_REDUCE_MEAN_FW_BW=1 + +python3 -m oneflow.distributed.launch \ + --nproc_per_node $ngpus \ + --nnodes 1 \ + --node_rank 0 \ + --master_addr 127.0.0.1 \ + dlrm_benchmark_a100.py \ + --data_dir $data_dir +