Support cluster step trace.

This commit is contained in:
liuchuting 2022-03-02 15:23:15 +08:00
parent 5b1772d2b3
commit cf2877d86e
4 changed files with 384 additions and 83 deletions

View File

@ -116,6 +116,7 @@ void Profiler::RecordOneStepStartEndInfo() {
step_start_end_info_vector_.erase(step_start_end_info_vector_.begin(),
step_start_end_info_vector_.begin() + iter_end_op_index + 1);
} else {
step_start_end_info_.fp_start_op_name = step_start_end_info_vector_[1];
step_start_end_info_.iter_end_op_name = step_start_end_info_vector_[step_start_end_info_vector_.size() - 1];
step_start_end_info_vector_.clear();
}

View File

@ -351,7 +351,9 @@ class Integrator:
if not (reduce_start and reduce_duration):
logger.info("Reduce event missing value.")
continue
cur_stream_id = reduce_field.split('_', 2)[1]
cur_stream_id = reduce_field.split('_', 3)[1]
if reduce_field.split('_', 2)[1] == 'ops':
cur_stream_id = reduce_field.split('_', 3)[2]
reduce_meta = [reduce_field, int(cur_stream_id), reduce_start,
reduce_duration, reduce_pid]
reduce_info.append(reduce_meta)
@ -515,6 +517,10 @@ class BaseTimelineGenerator:
_HOST_CPU_PID = 11000
_OP_OVERLAP_PID = 12000
_OP_GPU_ACTIVITY_PID = 13000
_RECEIVE_ALONE = 7997
_ALLREDUCE_ALONE = 7998
_MERGED_COMPUTATION_TID = 7999
_PURE_COMMUNICATION_TID = 8000
_MERGED_COMMUNICATION_TID = 8001
@ -559,6 +565,8 @@ class BaseTimelineGenerator:
{"name": "process_labels", "ph": "M", "pid": self._HOST_CPU_PID, "args": {"labels": "Host CPU Op"}},
{"name": "process_labels", "ph": "M", "pid": self._OP_OVERLAP_PID,
"args": {"labels": "Op Overlap Analyse"}},
{"name": "process_labels", "ph": "M", "pid": self._OP_GPU_ACTIVITY_PID,
"args": {"labels": "Activity Op"}},
{"name": "process_sort_index", "ph": "M", "pid": self._device_id, "args": {"sort_index": 0}},
{"name": "process_sort_index", "ph": "M", "pid": self._AI_CPU_PID, "args": {"sort_index": 10}},
@ -613,8 +621,8 @@ class BaseTimelineGenerator:
# merged_display_list data used for ui page.
merged_display_list = [
[display_name, tid, time_merged_segment_list[i * 2],
(time_merged_segment_list[i * 2 + 1] - time_merged_segment_list[i * 2]) * factor, pid]
for i in range(len(time_merged_segment_list) // 2)
(time_merged_segment_list[i * 2 + 1] - time_merged_segment_list[i * 2]) * factor, pid] for i in \
range(len(time_merged_segment_list) // 2)
]
if get_interval_time:
@ -622,8 +630,8 @@ class BaseTimelineGenerator:
# merged_res_list data used to compute overlap with other time_list.
merged_res_list = [
[display_name, tid, time_merged_segment_list[i * 2], time_merged_segment_list[i * 2 + 1], pid]
for i in range(len(time_merged_segment_list) // 2)
[display_name, tid, time_merged_segment_list[i * 2], time_merged_segment_list[i * 2 + 1], pid] for i in \
range(len(time_merged_segment_list) // 2)
]
# interval_display_list is interval time used for ui page.
@ -787,7 +795,8 @@ class BaseTimelineGenerator:
or time_item[op_full_name_idx] == self._step_end_op_name:
start_time = timeline_list[val['start_item_idx']][self._start_time_idx]
duration = (float(timeline_list[val['end_item_idx']][self._start_time_idx]) - float(start_time)) * \
factor_start_time_to_duration + float(timeline_list[val['end_item_idx']][self._duration_idx])
factor_start_time_to_duration + \
float(timeline_list[val['end_item_idx']][self._duration_idx])
scope_name_time_item = [key, tid, start_time, duration]
scope_name_time_list.append(scope_name_time_item)
scope_name_start_duration_dict[key]['start_item_idx'] = invalid_idx
@ -829,7 +838,7 @@ class BaseTimelineGenerator:
cur_step_start_time = time_item[self._start_time_idx]
if time_item[self._op_name_idx] == self._step_end_op_name:
cur_step_duration_time = (float(time_item[self._start_time_idx]) - float(cur_step_start_time)) * \
factor_start_time_to_duration + float(time_item[self._duration_idx])
float(factor_start_time_to_duration) + float(time_item[self._duration_idx])
step_time_item = [str(step_num), tid, float(cur_step_start_time), cur_step_duration_time]
step_time_list.append(step_time_item)
step_num += 1
@ -845,15 +854,26 @@ class GpuTimelineGenerator(BaseTimelineGenerator):
_output_activity_execute_time_file_path = "activity_execute_timestamp_{}.txt"
_output_gpu_activity_info_file_path = "gpu_activity_data_{}.csv"
_step_trace_original_filename = 'step_trace_profiling_{}.txt'
_cluster_analyse_filename = 'gpu_cluster_analyse_{}_{}_{}_{}.csv'
_activity_keys_list = []
def __init__(self, profiling_dir, device_id):
def __init__(self, profiling_dir, device_id, rank_size):
super().__init__()
self._device_id = device_id
self._rank_size = rank_size
self._profiling_dir = profiling_dir
self._device_id = device_id
self._timeline_meta = []
self._display_filename = self._display_filename.format(device_id)
self._timeline_summary_filename = self._timeline_summary_filename.format(device_id)
self._tid_dict = {
"receive_op_not_overlapped": (self._RECEIVE_ALONE, self._OP_OVERLAP_PID),
"exclude_receive_op": (self._ALLREDUCE_ALONE, self._OP_OVERLAP_PID),
"computation_op": (self._MERGED_COMPUTATION_TID, self._OP_OVERLAP_PID),
"communication_not_overlapped": (self._PURE_COMMUNICATION_TID, self._OP_OVERLAP_PID),
"communication": (self._MERGED_COMMUNICATION_TID, self._OP_OVERLAP_PID),
"free_time": (self._FREE_TIME_TID, self._OP_OVERLAP_PID)
}
def _get_and_validate_path(self, file_name):
"""Generate op or activity file path from file name, and validate this path."""
@ -880,7 +900,10 @@ class GpuTimelineGenerator(BaseTimelineGenerator):
timeline_dict['ts'] = (op_meta.start_time - min_cycle_counter) / factor
dur = op_meta.duration
timeline_dict['dur'] = dur
timeline_dict['pid'] = int(self._device_id)
if op_meta.pid is None:
timeline_dict['pid'] = int(self._device_id)
else:
timeline_dict['pid'] = op_meta.pid
if op_meta.stream_id == "Scope Name":
# remove the level of scope name which has a format like "0-conv2-Conv2d".
timeline_dict['name'] = "-".join(op_meta.op_name.split('-')[1:])
@ -904,7 +927,7 @@ class GpuTimelineGenerator(BaseTimelineGenerator):
self._update_format_meta_data(timeline_dict)
self._timeline_meta.append(timeline_dict)
def _load_timeline_data(self):
def _load_timeline_data(self, reduce_op_type):
"""Load timeline data from file."""
op_file_path = self._get_and_validate_path(
self._output_op_execute_time_file_path)
@ -913,10 +936,10 @@ class GpuTimelineGenerator(BaseTimelineGenerator):
activity_args_file_path = self._get_and_validate_path(
self._output_gpu_activity_info_file_path)
timeline_list = self._load_op_data(op_file_path)
timeline_list, communication_info = self._load_op_data(op_file_path, reduce_op_type)
communication_info.sort(key=lambda x: float(x[2]))
# Add host cpu op timeline.
cpu_timeline_generator = CpuTimelineGenerator(self._profiling_dir, self._device_id)
cpu_timeline_generator = CpuTimelineGenerator(self._profiling_dir, self._device_id, self._rank_size)
cpu_timeline_list = cpu_timeline_generator.load_cpu_op_data()
if cpu_timeline_list:
self._clock_synchronize_to_gpu(cpu_timeline_list)
@ -941,15 +964,27 @@ class GpuTimelineGenerator(BaseTimelineGenerator):
factor_start_time_uint_to_duration)
recompute_scope_name_time_list = self._get_scope_name_time_list(timeline_list, "recompute_Default",
factor_start_time_uint_to_duration)
activity_timeline_list, cuda_compute_ops_timeline_list = self._load_activity_data( \
activity_file_path, activity_args_file_path)
# Add AllReduce info to timeline temp list and sort by start time.
if communication_info:
logger.debug('Allreduce info found, Start adding info to timeline...')
cluster_related_timeline = self._get_cluster_timeline(
timeline_list, cuda_compute_ops_timeline_list, communication_info, step_time_list)
timeline_list.extend(cluster_related_timeline)
timeline_list.extend(communication_info)
timeline_list.sort(key=lambda x: float(x[self._start_time_idx]))
timeline_list.extend(default_scope_name_time_list)
timeline_list.extend(gradient_scope_name_time_list)
timeline_list.extend(recompute_scope_name_time_list)
timeline_list.extend(step_time_list)
timeline_list.sort(key=lambda x: (float(x[self._start_time_idx]), x[self._tid_idx]))
timeline_list.sort(key=lambda x: (float(x[self._start_time_idx])))
# Add cuda activity timeline.
activity_timeline_list = self._load_activity_data(activity_file_path, activity_args_file_path)
timeline_list.extend(activity_timeline_list)
timeline_list.sort(key=lambda x: float(x[2]))
@ -974,9 +1009,10 @@ class GpuTimelineGenerator(BaseTimelineGenerator):
for idx, time_item in enumerate(timeline_list):
timeline_list[idx][self._start_time_idx] = int(time_item[self._start_time_idx]) + time_diff
def _load_op_data(self, op_file_path):
def _load_op_data(self, op_file_path, reduce_op_type):
"""Load operator data from file"""
op_timeline_list = []
communication_info = []
try:
with open(op_file_path, 'r') as f_obj:
for line in f_obj:
@ -986,23 +1022,22 @@ class GpuTimelineGenerator(BaseTimelineGenerator):
time_arr = time_arr.split(" ")
for time in time_arr:
time = time.split(",")
if len(time) == 3:
# for time value is [start_timestamp, duration, tid]
# line_list[1] would be like "HostCpuOps" + str(tid)
line_list = op_list[:1] + [op_list[1] + str(time[-1])] + time[:-1]
line_list = op_list[:2] + time
communication_op_name = line_list[0].strip().split('/')[-1]
if communication_op_name not in reduce_op_type:
op_timeline_list.append(line_list)
else:
# for time value is [start_timestamp, duration]
line_list = op_list[:2] + time
op_timeline_list.append(line_list)
communication_info.append(line_list)
except (IOError, OSError) as err:
logger.critical('Error occurred when load operator timeline data intermediate file: %s', err)
raise ProfilerIOException()
return op_timeline_list
return op_timeline_list, communication_info
def _load_activity_data(self, activity_file_path, activity_args_file_path):
"""Load activity data from file"""
activity_timeline_list = []
cuda_compute_ops_timeline_list = []
try:
args_dict = {}
with open(activity_args_file_path, 'r') as args_file:
@ -1017,16 +1052,18 @@ class GpuTimelineGenerator(BaseTimelineGenerator):
line_list = line.strip('\n').split(';')
# concat activity args info.
line_list += args_dict[line_list[0]]
if not line_list[0].startswith('nccl'):
cuda_compute_ops_timeline_list.append(line_list)
activity_timeline_list.append(line_list)
except (IOError, OSError) as err:
logger.critical('Error occurred when load activity timeline data intermediate file: %s', err)
raise ProfilerIOException()
return activity_timeline_list
return activity_timeline_list, cuda_compute_ops_timeline_list
def init_timeline(self):
def init_timeline(self, reduce_op_type):
"""Init timeline metadata, adding all collected info."""
timeline_list = self._load_timeline_data()
timeline_list = self._load_timeline_data(reduce_op_type)
# Init a dict for counting the num of streams.
stream_count_dict = {}
@ -1108,6 +1145,224 @@ class GpuTimelineGenerator(BaseTimelineGenerator):
logger.critical(f'Error occurred when read {step_trace_profiling_path}: {err}')
raise ProfilerIOException()
def _get_cluster_timeline(self, timeline, activity_info, comm_info, step_info):
"""
Analyse the cluster communication and computation data, and write result to file.
To analyse the cluster performance bottleneck based on timeline, define the time of a training
step as "t_total", propose five metrics as follows:
1) The time that "receive" operators not overlapped by others(t1)
2) The time that is consumed inside the stage(t_total - t1)
3) The time that "communication" operators not overlapped by others(t2)
4) The time that consumed by computation(t_total - t2)
5) The time that "collective communication" operators not overlapped by others(t3)
In pipeline parallel mode, we can locate slow stage based on t_total - t1. Inside each stage,
we can locate slow card based on t_total - t2. The value of t1 indicates the degree that
communication time between stages slow down the training. The value of t3 indicates the degree
that communication inside each stage slow down the training.
"""
step_num = len(step_info)
is_pipeline_parallel = False
comm_merged_timeline, _, comm_display_timeline = self._get_merged_time_list(
comm_info,
display_name="communication",
factor=1e-3
)
compute_op_timeline = timeline + activity_info
compute_op_timeline.sort(key=lambda x: float(x[self._start_time_idx]))
compute_op_timeline_interval, _, compute_op_display_timeline = self._get_merged_time_list(
compute_op_timeline,
get_interval_time=True,
factor=1e-3
)
# Consider if the overlap will be 0 or not.
comm_not_overlapped_timeline = self._get_intersection_time(
compute_op_timeline_interval,
comm_merged_timeline
)
# Process receive part.
all_timeline = timeline + comm_info
all_timeline.sort(key=lambda x: float(x[self._start_time_idx]))
receive_op_timeline = self._produce_two_separated_timeline(
all_timeline,
"Receive-op"
)[0]
if receive_op_timeline:
is_pipeline_parallel = True
receive_op_merged_timeline = self._get_merged_time_list(receive_op_timeline,
factor=1e-3)[0]
receive_op_not_overlapped_timeline = self._get_intersection_time(
compute_op_timeline_interval,
receive_op_merged_timeline,
display_name="receive_op_not_overlapped"
)
# Process collective communication part.
collective_comm_timeline = self._produce_two_separated_timeline(
comm_info,
"Receive-op"
)[-1]
collective_comm_merged_timeline = self._get_merged_time_list(collective_comm_timeline,
factor=1e-3)[0]
collective_comm_not_overlapped_timeline = self._get_intersection_time(
compute_op_timeline_interval,
collective_comm_merged_timeline,
display_name="exclude_receive_op"
)
# Generate free time that exclude computation and communication time.
all_timeline = compute_op_timeline + comm_info
all_timeline.sort(key=lambda x: float(x[self._start_time_idx]))
free_timeline = self._get_merged_time_list(
all_timeline,
get_interval_time=True,
display_name="free_time",
factor=1e-3
)[1]
# Compute these five metrics mentioned above per step.
recieve_alone_time = self._compute_time_inside_step(receive_op_not_overlapped_timeline, step_info)
stage_time, computation_time = [], []
comm_alone_time = self._compute_time_inside_step(comm_not_overlapped_timeline, step_info)
collective_comm_alone_time = self._compute_time_inside_step(
collective_comm_not_overlapped_timeline,
step_info
)
for step in range(step_num):
try:
if is_pipeline_parallel:
stage_time.append(step_info[step][self._duration_idx] - recieve_alone_time[step])
computation_time.append(step_info[step][self._duration_idx] - comm_alone_time[step])
except IndexError as e:
logger.error(e)
metrices_per_step_list = [computation_time, comm_alone_time, stage_time,
recieve_alone_time, collective_comm_alone_time]
if step_num > 1:
for metric in metrices_per_step_list:
metric.append(sum(metric[1:]) / (step_num - 1))
self._write_cluster_metrices(metrices_per_step_list, is_pipeline_parallel)
res_timeline = []
res_timeline.extend(comm_not_overlapped_timeline)
res_timeline.extend(compute_op_display_timeline)
res_timeline.extend(comm_display_timeline)
res_timeline.extend(free_timeline)
return res_timeline
def _write_cluster_metrices(self, metrices, is_pipeline_parallel):
"""Write cluster metric."""
# Note that the feature of cluster bottleneck analyse is not supported in offline parse mode,
# due to that parallel context is not set.
try:
parallel_mode = get_auto_parallel_context("parallel_mode")
stage_num = get_auto_parallel_context("pipeline_stages")
except RuntimeError:
logger.warning("[profiler] the feature of cluster bottleneck analyse "
"is not supported in offline parse mode.")
parallel_mode = "data_parallel"
stage_num = 1
if stage_num > 1:
parallel_mode = "pipeline-parallel"
elif parallel_mode != "data_parallel":
parallel_mode = "model-parallel"
else:
parallel_mode = "data-parallel"
cluster_analyse_file_path = os.path.join(
self._profiling_dir,
self._cluster_analyse_filename.format(parallel_mode, stage_num, self._rank_size, self._device_id))
cluster_analyse_file_path = validate_and_normalize_path(cluster_analyse_file_path)
try:
with open(cluster_analyse_file_path, 'w') as file_handle:
csv_writer = csv.writer(file_handle)
if is_pipeline_parallel:
header = ['computation_time', 'communication_alone_time', 'stage_time',
'receive_alone_time', 'collective_communication_alone_time']
zip_metrices = zip(metrices[0], metrices[1], metrices[2], metrices[3], metrices[4])
else:
header = ['computation_time', 'communication_alone_time']
zip_metrices = zip(metrices[0], metrices[1])
csv_writer.writerow(header)
for row_data in zip_metrices:
row_data = [round(val / 1e3, 4) for val in row_data]
csv_writer.writerow(row_data)
os.chmod(cluster_analyse_file_path, stat.S_IREAD | stat.S_IWRITE)
except (IOError, OSError) as err:
logger.warning(f'Failed to save {cluster_analyse_file_path}. {err}')
raise ProfilerIOException
def _compute_time_inside_step(self, metric_timeline, step_time_list):
"""Compute per step time of metric_timeline."""
per_step_time_list = []
step = 0
cur_step_metric_time = 0
factor_us_to_ns = 1e3
step_end_time = step_time_list[step][self._start_time_idx] + \
step_time_list[step][self._duration_idx] * factor_us_to_ns
for time_item in metric_timeline:
start_time = time_item[self._start_time_idx]
if start_time > step_end_time:
per_step_time_list.append(cur_step_metric_time)
step += 1
if step >= len(step_time_list):
logger.warning("Compute profiler compute_time_inside_step time, "
"find the data length is more than step count, "
"maybe current graph has multi sub graph, skip the last data.")
break
step_end_time = step_time_list[step][self._start_time_idx] + \
step_time_list[step][self._duration_idx] * factor_us_to_ns
cur_step_metric_time = 0
cur_step_metric_time += time_item[self._duration_idx]
per_step_time_list.append(cur_step_metric_time)
return per_step_time_list
def _get_intersection_time(self, first_time_list, second_time_list,
display_name="communication_not_overlapped"):
"""Get intersection time of two time list."""
first_list_idx, second_list_idx = 0, 0
first_list_len = len(first_time_list)
second_list_len = len(second_time_list)
intersection_segment_display_list = []
factor_ns_to_us = 1e-3
while first_list_idx < first_list_len and second_list_idx < second_list_len:
intersection_start = max(
first_time_list[first_list_idx][self._start_time_idx],
second_time_list[second_list_idx][self._start_time_idx]
)
intersection_end = min(
first_time_list[first_list_idx][self._duration_idx],
second_time_list[second_list_idx][self._duration_idx]
)
if intersection_start < intersection_end:
intersection_segment_display_list.append(
[display_name, self._tid_dict[display_name][0],
intersection_start, (intersection_end - intersection_start) * factor_ns_to_us,
self._tid_dict[display_name][1]]
)
if first_time_list[first_list_idx][self._duration_idx] >= \
second_time_list[second_list_idx][self._duration_idx]:
second_list_idx += 1
else:
first_list_idx += 1
return intersection_segment_display_list
def _produce_two_separated_timeline(self, timeline, op_name):
"""Produce two separated timeline based on op_name."""
timeline_include_op_name = []
timeline_exclude_op_name = []
for time_item in timeline:
if op_name in time_item[self._op_name_idx]:
timeline_include_op_name.append(time_item)
else:
timeline_exclude_op_name.append(time_item)
return timeline_include_op_name, timeline_exclude_op_name
class AscendTimelineGenerator(BaseTimelineGenerator):
"""Generate ascend Timeline data from file."""
@ -1199,12 +1454,12 @@ class AscendTimelineGenerator(BaseTimelineGenerator):
all_reduce_names = []
for info in communication_info:
# stream_{stream_id}_{stream_op_index}_{opname}
all_reduce_name = info[0][info[0].rindex('_')+1:]
all_reduce_name = info[0][info[0].rindex('_') + 1:]
if all_reduce_name not in all_reduce_names:
all_reduce_names.append(all_reduce_name)
timeline_list = self._load_timeline_data(all_reduce_names)
cpu_timeline_generator = CpuTimelineGenerator(self._profiling_dir, self._rank_id)
cpu_timeline_generator = CpuTimelineGenerator(self._profiling_dir, self._rank_id, self._rank_size)
cpu_timeline_list = cpu_timeline_generator.get_timeline_data()
if cpu_timeline_list:
self._clock_synchronize_to_device(cpu_timeline_list, source_path)
@ -1237,7 +1492,7 @@ class AscendTimelineGenerator(BaseTimelineGenerator):
# Add AllReduce info to timeline temp list and sort by start time.
if communication_info:
logger.debug('AllReduce info found. Start adding info into timeline...')
cluster_related_timeline = self._analyse_and_write_cluster_profiling_data(
cluster_related_timeline = self._get_cluster_timeline(
timeline_list, communication_info, step_time_list)
timeline_list.extend(cluster_related_timeline)
timeline_list.extend(communication_info)
@ -1345,7 +1600,7 @@ class AscendTimelineGenerator(BaseTimelineGenerator):
timeline_exclude_op_name.append(time_item)
return timeline_include_op_name, timeline_exclude_op_name
def _analyse_and_write_cluster_profiling_data(self, aicore_timeline, communication_timeline, step_time_list):
def _get_cluster_timeline(self, aicore_info, comm_info, step_info):
"""
Analyse the cluster communication and computation data, and write result to file.
@ -1361,24 +1616,24 @@ class AscendTimelineGenerator(BaseTimelineGenerator):
communication time between stages slow down the training. The value of t3 indicates the degree
that communication inside each stage slow down the training.
"""
step_num = len(step_time_list)
step_num = len(step_info)
is_pipeline_parallel = False
comm_merged_timeline, _, comm_display_timeline = self._get_merged_time_list(
communication_timeline,
comm_info,
display_name="communication"
)
aicore_timeline_interval, _, aicore_display_timeline = self._get_merged_time_list(
aicore_timeline,
aicore_info,
get_interval_time=True
)
# Consider if the overlap will be 0 or not.
comm_not_overlaped_timeline = self._get_intersection_time(
comm_not_overlapped_timeline = self._get_intersection_time(
aicore_timeline_interval,
comm_merged_timeline
)
# Process receive part.
all_timeline = aicore_timeline + communication_timeline
all_timeline = aicore_info + comm_info
all_timeline.sort(key=lambda x: float(x[self._start_time_idx]))
receive_op_timeline, timeline_exclude_receive_op = self._produce_two_separated_timeline(
all_timeline,
@ -1391,18 +1646,18 @@ class AscendTimelineGenerator(BaseTimelineGenerator):
timeline_exclude_receive_op,
get_interval_time=True
)[0]
receive_op_not_overlaped_timeline = self._get_intersection_time(
receive_op_not_overlapped_timeline = self._get_intersection_time(
timeline_exclude_receive_op_interval,
receive_op_merged_timeline
)
# Process collective communication part.
collective_comm_timeline = self._produce_two_separated_timeline(
communication_timeline,
comm_info,
"Receive-op"
)[-1]
collective_comm_merged_timeline = self._get_merged_time_list(collective_comm_timeline)[0]
collective_comm_not_overlaped_timeline = self._get_intersection_time(
collective_comm_not_overlapped_timeline = self._get_intersection_time(
aicore_timeline_interval,
collective_comm_merged_timeline
)
@ -1411,38 +1666,35 @@ class AscendTimelineGenerator(BaseTimelineGenerator):
free_timeline = self._get_merged_time_list(
all_timeline,
get_interval_time=True,
display_name="free_time"
)[1]
display_name="free_time")[1]
try:
# Compute these five metrics mentioned above per step.
recieve_alone_time = self._compute_time_inside_step(receive_op_not_overlaped_timeline, step_time_list)
stage_time, computation_time = [], []
comm_alone_time = self._compute_time_inside_step(comm_not_overlaped_timeline, step_time_list)
collective_comm_alone_time = self._compute_time_inside_step(
collective_comm_not_overlaped_timeline,
step_time_list
)
for step in range(step_num):
# Compute these five metrics mentioned above per step.
recieve_alone_time = self._compute_time_inside_step(receive_op_not_overlapped_timeline, step_info)
stage_time, computation_time = [], []
comm_alone_time = self._compute_time_inside_step(comm_not_overlapped_timeline, step_info)
collective_comm_alone_time = self._compute_time_inside_step(
collective_comm_not_overlapped_timeline,
step_info
)
for step in range(step_num):
try:
if is_pipeline_parallel:
stage_time.append(step_time_list[step][self._duration_idx] - recieve_alone_time[step])
computation_time.append(step_time_list[step][self._duration_idx] - comm_alone_time[step])
metrices_per_step_list = [computation_time, comm_alone_time, stage_time,
recieve_alone_time, collective_comm_alone_time]
if step_num > 1:
for metric in metrices_per_step_list:
metric.append(sum(metric[1:]) / (step_num - 1))
self._write_cluster_metrices(metrices_per_step_list, is_pipeline_parallel)
except IndexError as e:
logger.error(e)
stage_time.append(step_info[step][self._duration_idx] - recieve_alone_time[step])
computation_time.append(step_info[step][self._duration_idx] - comm_alone_time[step])
except IndexError as e:
logger.error(e)
metrices_per_step_list = [computation_time, comm_alone_time, stage_time, recieve_alone_time, \
collective_comm_alone_time]
if step_num > 1:
for metric in metrices_per_step_list:
metric.append(sum(metric[1:]) / (step_num - 1))
self._write_cluster_metrices(metrices_per_step_list, is_pipeline_parallel)
res_timeline = []
res_timeline.extend(comm_not_overlaped_timeline)
res_timeline.extend(comm_not_overlapped_timeline)
res_timeline.extend(aicore_display_timeline)
res_timeline.extend(comm_display_timeline)
res_timeline.extend(free_timeline)
return res_timeline
def _write_cluster_metrices(self, metrices, is_pipeline_parallel):
@ -1466,8 +1718,7 @@ class AscendTimelineGenerator(BaseTimelineGenerator):
cluster_analyse_file_path = os.path.join(
self._profiling_dir,
self._cluster_analyse_filename.format(parallel_mode, stage_num, self._rank_size, self._rank_id)
)
self._cluster_analyse_filename.format(parallel_mode, stage_num, self._rank_size, self._rank_id))
cluster_analyse_file_path = validate_and_normalize_path(cluster_analyse_file_path)
try:
@ -1494,8 +1745,7 @@ class AscendTimelineGenerator(BaseTimelineGenerator):
per_step_time_list = []
step = 0
cur_step_metric_time = 0
step_end_time = step_time_list[step][self._start_time_idx] + \
step_time_list[step][self._duration_idx]
step_end_time = step_time_list[step][self._start_time_idx] + step_time_list[step][self._duration_idx]
for time_item in metric_timeline:
start_time = time_item[self._start_time_idx]
if start_time > step_end_time:
@ -1506,8 +1756,7 @@ class AscendTimelineGenerator(BaseTimelineGenerator):
"find the data length is more than step count, "
"maybe current graph has multi sub graph, skip the last data.")
break
step_end_time = step_time_list[step][self._start_time_idx] + \
step_time_list[step][self._duration_idx]
step_end_time = step_time_list[step][self._start_time_idx] + step_time_list[step][self._duration_idx]
cur_step_metric_time = 0
cur_step_metric_time += time_item[self._duration_idx]
per_step_time_list.append(cur_step_metric_time)
@ -1521,9 +1770,7 @@ class AscendTimelineGenerator(BaseTimelineGenerator):
first_list_len = len(first_time_list)
second_list_len = len(second_time_list)
intersection_segment_display_list = []
while first_list_idx < first_list_len and \
second_list_idx < second_list_len:
while first_list_idx < first_list_len and second_list_idx < second_list_len:
intersection_start = max(
first_time_list[first_list_idx][self._start_time_idx],
second_time_list[second_list_idx][self._start_time_idx]
@ -1577,6 +1824,32 @@ class CpuTimelineGenerator(GpuTimelineGenerator):
return timeline_list
def _load_op_data(self, op_file_path):
"""Load operator data from file"""
op_timeline_list = []
try:
with open(op_file_path, 'r') as f_obj:
for line in f_obj:
self._timeline_summary['num_of_ops'] += 1
op_list = line.strip('\n').strip().split(';')
time_arr = op_list[-1]
time_arr = time_arr.split(" ")
for time in time_arr:
time = time.split(",")
if len(time) == 3:
# for time value is [start_timestamp, duration, tid]
# line_list[1] would be like "HostCpuOps" + str(tid)
line_list = op_list[:1] + [op_list[1] + str(time[-1])] + time[:-1]
else:
# for time value is [start_timestamp, duration]
line_list = op_list[:2] + time
op_timeline_list.append(line_list)
except (IOError, OSError) as err:
logger.critical('Error occurred when load operator timeline data intermediate file: %s', err)
raise ProfilerIOException()
return op_timeline_list
def get_timeline_data(self):
"""Get timeline data from file."""
timeline_list = self.load_cpu_op_data()

View File

@ -255,6 +255,7 @@ class GpuStepTraceParser(BaseStepTraceParser):
def __init__(self, *args, **kwargs):
super(GpuStepTraceParser, self).__init__(*args, **kwargs)
self._source_file_path = self._input_dir
self._reduce_op_type = []
def get_fp_bp(self, f_obj, all_step_fp, all_step_bp):
"""Parser the fp and bp."""
@ -398,6 +399,7 @@ class GpuStepTraceParser(BaseStepTraceParser):
# add communication op start and end time, time unit from ns to 10ns.
reduce_time_info.append(reduce_item.split(',')[1][:-1])
reduce_time_info.append(reduce_item.split(',')[2][:-1])
self._reduce_op_type.append(reduce_item.split(',')[0].split('/')[-1])
step_trace = {
'start': start_time,
'fp': fp_time,
@ -438,7 +440,8 @@ class GpuStepTraceParser(BaseStepTraceParser):
"""
reduce_info = {}
op_type = 'AllReduce'
index = int(field_name.split('_')[2])
op_type = self._reduce_op_type[index]
# append field name with op type.
field_name += '_' + op_type
reduce_info[field_name] = int(end_point) - int(start_point)

View File

@ -317,7 +317,7 @@ class Profiler:
logger.info("Profiling: all the data have been analyzed.")
def _ascend_analyse(self):
"""Collect and analyse ascend performance data"""
"""Collect and analyse ascend performance data."""
self._rank_size = 1
if self._profile_communication and not GlobalComm.INITED:
self._profile_communication = False
@ -559,17 +559,22 @@ class Profiler:
logger.info("Profiling: stop time: %d", self._stop_time)
def _gpu_analyse(self):
"""Collect and analyse gpu performance data"""
"""Collect and analyse gpu performance data."""
self._dev_id = context.get_context("device_id")
self._rank_size = 1
if GlobalComm.WORLD_COMM_GROUP == "nccl_world_group":
self._dev_id = str(get_rank())
if GlobalComm.INITED:
self._rank_size = get_group_size()
if self._has_started:
self.stop()
else:
logger.info("No need to stop profiler because profiler has been stopped.")
timeline_generator = self._generate_timeline()
reduce_op_type = self._get_step_reduce_op_type()
timeline_generator = self._generate_timeline(reduce_op_type)
# parse minddata pipeline operator and queue for GPU
try:
@ -603,12 +608,32 @@ class Profiler:
'otherwise, this warning can be ignored.'
)
logger.warning(
'\nProfile communication is not supported on GPU currently.\n'
'Please running on Ascend if you would like to see cluster communication analysis, '
'otherwise, this warning can be ignored.'
)
def _get_step_reduce_op_type(self):
"""Gets all communication operator names."""
step_trace_original_filename = f'step_trace_profiling_{self._dev_id}.txt'
step_trace_file_path = os.path.join(self._output_path, step_trace_original_filename)
step_trace_file_path = validate_and_normalize_path(step_trace_file_path)
reduce_op_type = []
with open(step_trace_file_path, 'r') as f_obj:
one_step_info = f_obj.readline().strip().split()
# The communication operator starts at index 4.
for reduce_item in one_step_info[4:]:
reduce_op_type.append(reduce_item.split(',')[0].split('/')[-1])
return reduce_op_type
def _cpu_analyse(self):
"""Collect and analyse cpu performance data"""
"""Collect and analyse cpu performance data."""
try:
size_limit = 100 * 1024 * 1024 # 100MB
timeline_generator = CpuTimelineGenerator(self._output_path, 0)
timeline_generator = CpuTimelineGenerator(self._output_path, 0, 1)
timeline_generator.init_timeline()
timeline_generator.write_timeline(size_limit)
timeline_generator.write_timeline_summary()
@ -617,7 +642,6 @@ class Profiler:
logger.warning('Fail to write timeline data: %s', err)
raise RuntimeError('Fail to write timeline data.')
def _analyse_step_trace(self, source_path=None, framework_parser=None, is_training_mode_flag=True,
is_gpu_kernel_async_launch_flag=False):
"""
@ -711,12 +735,12 @@ class Profiler:
timeline_analyser.write_timeline(size_limit)
timeline_analyser.write_timeline_summary()
def _generate_timeline(self):
def _generate_timeline(self, reduce_op_type):
"""Used for gpu, generate timeline info, write to json format file."""
try:
size_limit = 100 * 1024 * 1024 # 100MB
timeline_generator = GpuTimelineGenerator(self._output_path, self._dev_id)
timeline_generator.init_timeline()
timeline_generator = GpuTimelineGenerator(self._output_path, self._dev_id, self._rank_size)
timeline_generator.init_timeline(reduce_op_type)
timeline_generator.write_timeline(size_limit)
timeline_generator.write_timeline_summary()
return timeline_generator