From cb55f9e31c0494fda0866ad3dd28d1007ca4c824 Mon Sep 17 00:00:00 2001 From: Alexander Movsesyan Date: Fri, 27 Sep 2024 00:20:14 -0400 Subject: [PATCH 01/11] initial draft of base reader --- pipit/readers/base_reader.py | 126 +++++++++++++++++++++++++++++++++++ 1 file changed, 126 insertions(+) create mode 100644 pipit/readers/base_reader.py diff --git a/pipit/readers/base_reader.py b/pipit/readers/base_reader.py new file mode 100644 index 00000000..caefad5a --- /dev/null +++ b/pipit/readers/base_reader.py @@ -0,0 +1,126 @@ +from abc import ABC, abstractmethod +from typing import List, Dict +from ..graph import Graph, Node +import numpy + + +class BaseTraceReader(ABC): + + + # The following methods should be called by each reader class + def create_empty_trace(self, num_processes: int, create_cct: bool) -> None: + # keep track if we want to create a CCT + self.create_cct = create_cct + + # keep track of a unique id for each event + self.unique_id = -1 + + # events are indexed by process number, then thread number + self.events: List[Dict[List[Dict]]] = [{}] * num_processes + + # stacks are indexed by process number, then thread number + self.stacks: List[Dict[List[int]]] = [{}] * num_processes + + self.ccts: List[Dict[Graph]] = [{}] * num_processes + + + + def add_event(self, event: Dict) -> None: + + # get process number -- if not present, set to 0 + if "process" in event: + process = event["process"] + else: + process = 0 + + # get process number -- if not present, set to 0 + if "thread" in event: + thread = event["thread"] + else: + thread = 0 + + # assign a unique id to the event + event["id"] = self.__get_unique_id() + + + # get event list + if thread not in self.events[process]: + self.events[process][thread] = [] + event_list: List[Dict] = self.events[process][thread] + + # get stack + if thread not in self.stacks[process]: + self.stacks[process][thread] = [] + stack: List[int] = self.stacks[process][thread] + + # if the event is an enter event, add the event to the stack and CCT + if event["Event Type"] == "Enter": + cct = None + # if we are creating a CCT, get the correct CCT + if self.create_cct: + if thread not in self.ccts[process]: + self.ccts[process][thread] = Graph() + cct = self.ccts[process][thread] + self.__update_cct_and_parent_child_relationships(event, self.stacks[process][thread], event_list, cct) + elif event["Event Type"] == "Leave": + self.__update_match_event(event, self.stacks[process][thread], event_list) + + + def finalize(self) -> None: + pass + + # Helper methods + + # This method can be thought of the update upon an "Enter" event + # It adds to the stack and CCT + def __update_cct_and_parent_child_relationships(self, event: Dict, stack: List[int], event_list: List[Dict], cct: Graph) -> None: + if len(stack) == 0: + # root event + event["parent"] = numpy.nan + if self.create_cct: + new_graph_node = Node(event["id"], None) + cct.add_root(new_graph_node) + event["Node"] = new_graph_node + else: + parent_event = event_list[stack[-1]] + event["parent"] = parent_event["id"] + if self.create_cct: + new_graph_node = Node(event["id"], parent_event["Node"]) + parent_event["Node"].add_child(new_graph_node) + event["Node"] = new_graph_node + + # update stack and event list + stack.append(len(event_list) - 1) + # event_list.append(event) + + + # def __update_cct(self, event: Dict, process: int, thread: int) -> None: + # pass + + # This method can be thought of the update upon a "Leave" event + # It pops from the stack and updates the event list + # We should look into using this function to add artificial "Leave" events for unmatched "Enter" events + def __update_match_event(self, leave_event: Dict, stack: List[int], event_list: List[Dict]) -> None: + + while len(stack) > 0: + enter_event = event_list[stack[-1]] + + if enter_event["Name"] == leave_event["Name"]: + # matching event found + + # update matching event ids + leave_event["_matching_event"] = enter_event["id"] + enter_event["_matching_event"] = leave_event["id"] + + # popping matched events from the stack + stack.pop() + break + else: + # popping unmatched events from the stack + stack.pop() + + # event_list.append(leave_event) + + def __get_unique_id(self) -> int: + self.unique_id += 1 + return self.unique_id \ No newline at end of file From 8d7d9e2667c839a91694acee77433b92fa130a0d Mon Sep 17 00:00:00 2001 From: Alexander Movsesyan Date: Mon, 30 Sep 2024 13:39:36 -0400 Subject: [PATCH 02/11] updates to base reader --- pipit/readers/base_reader.py | 31 +++++++++++++++++++++---------- 1 file changed, 21 insertions(+), 10 deletions(-) diff --git a/pipit/readers/base_reader.py b/pipit/readers/base_reader.py index caefad5a..419ea572 100644 --- a/pipit/readers/base_reader.py +++ b/pipit/readers/base_reader.py @@ -1,5 +1,8 @@ from abc import ABC, abstractmethod from typing import List, Dict + +import pandas as pd + from ..graph import Graph, Node import numpy @@ -16,12 +19,12 @@ def create_empty_trace(self, num_processes: int, create_cct: bool) -> None: self.unique_id = -1 # events are indexed by process number, then thread number - self.events: List[Dict[List[Dict]]] = [{}] * num_processes + self.events: List[Dict[int, List[Dict]]] = [{}] * num_processes # stacks are indexed by process number, then thread number - self.stacks: List[Dict[List[int]]] = [{}] * num_processes + self.stacks: List[Dict[int, List[int]]] = [{}] * num_processes - self.ccts: List[Dict[Graph]] = [{}] * num_processes + self.ccts: List[Dict[int, Graph]] = [{}] * num_processes @@ -67,6 +70,12 @@ def add_event(self, event: Dict) -> None: def finalize(self) -> None: + # first step put everything in one list + all_events = [] + for process in self.events: + for thread in process: + all_events.extend(process[thread]) + self.dataframe = pd.DataFrame(all_events) pass # Helper methods @@ -77,17 +86,19 @@ def __update_cct_and_parent_child_relationships(self, event: Dict, stack: List[i if len(stack) == 0: # root event event["parent"] = numpy.nan - if self.create_cct: - new_graph_node = Node(event["id"], None) - cct.add_root(new_graph_node) - event["Node"] = new_graph_node + # if self.create_cct: + # new_graph_node = Node(event["id"], None) + # cct.add_root(new_graph_node) + # event["Node"] = new_graph_node else: parent_event = event_list[stack[-1]] event["parent"] = parent_event["id"] if self.create_cct: - new_graph_node = Node(event["id"], parent_event["Node"]) - parent_event["Node"].add_child(new_graph_node) - event["Node"] = new_graph_node + parent_graph_node = parent_event["Node"] + # if + # new_graph_node = Node(event["id"], parent_event["Node"]) + # parent_event["Node"].add_child(new_graph_node) + # event["Node"] = new_graph_node # update stack and event list stack.append(len(event_list) - 1) From ceae939419f8f86c1019fe477121d93a7413a1e3 Mon Sep 17 00:00:00 2001 From: Alexander Movsesyan Date: Thu, 3 Oct 2024 13:13:56 -0400 Subject: [PATCH 03/11] base reader working without cct creation --- pipit/readers/base_reader.py | 14 +++++++++++--- 1 file changed, 11 insertions(+), 3 deletions(-) diff --git a/pipit/readers/base_reader.py b/pipit/readers/base_reader.py index 419ea572..a33e1e29 100644 --- a/pipit/readers/base_reader.py +++ b/pipit/readers/base_reader.py @@ -3,12 +3,17 @@ import pandas as pd +from .. import Trace from ..graph import Graph, Node import numpy class BaseTraceReader(ABC): + @abstractmethod + def read(self) -> Trace: + pass + # The following methods should be called by each reader class def create_empty_trace(self, num_processes: int, create_cct: bool) -> None: @@ -68,6 +73,8 @@ def add_event(self, event: Dict) -> None: elif event["Event Type"] == "Leave": self.__update_match_event(event, self.stacks[process][thread], event_list) + event_list.append(event) + def finalize(self) -> None: # first step put everything in one list @@ -75,8 +82,9 @@ def finalize(self) -> None: for process in self.events: for thread in process: all_events.extend(process[thread]) - self.dataframe = pd.DataFrame(all_events) - pass + # create a dataframe + self.events_dataframe = pd.DataFrame(all_events) + self.trace = Trace(None, self.events_dataframe, None) # Helper methods @@ -101,7 +109,7 @@ def __update_cct_and_parent_child_relationships(self, event: Dict, stack: List[i # event["Node"] = new_graph_node # update stack and event list - stack.append(len(event_list) - 1) + stack.append(len(event_list)) # event_list.append(event) From 8f7880af4e3389f7ee438cff8374c4fb898bf210 Mon Sep 17 00:00:00 2001 From: Alexander Movsesyan Date: Sun, 13 Oct 2024 18:36:02 -0400 Subject: [PATCH 04/11] Base Reader for single threaded reading --- pipit/readers/base_reader.py | 145 ++++++++++++++++++++--------------- 1 file changed, 83 insertions(+), 62 deletions(-) diff --git a/pipit/readers/base_reader.py b/pipit/readers/base_reader.py index a33e1e29..89057ab3 100644 --- a/pipit/readers/base_reader.py +++ b/pipit/readers/base_reader.py @@ -1,6 +1,7 @@ from abc import ABC, abstractmethod from typing import List, Dict +import pandas import pandas as pd from .. import Trace @@ -16,105 +17,128 @@ def read(self) -> Trace: # The following methods should be called by each reader class - def create_empty_trace(self, num_processes: int, create_cct: bool) -> None: - # keep track if we want to create a CCT - self.create_cct = create_cct - + def create_empty_trace(self, num_processes: int) -> None: # keep track of a unique id for each event self.unique_id = -1 # events are indexed by process number, then thread number - self.events: List[Dict[int, List[Dict]]] = [{}] * num_processes + # stores a list of events + self.events: List[Dict[int, List[Dict]]] = [] + for i in range(num_processes): + self.events.append({}) # stacks are indexed by process number, then thread number + # stores indices of events in the event list self.stacks: List[Dict[int, List[int]]] = [{}] * num_processes - self.ccts: List[Dict[int, Graph]] = [{}] * num_processes - def add_event(self, event: Dict) -> None: # get process number -- if not present, set to 0 - if "process" in event: - process = event["process"] + if "Process" in event: + process = event["Process"] else: + print("something is wrong") process = 0 - # get process number -- if not present, set to 0 - if "thread" in event: - thread = event["thread"] + # get thread number -- if not present, set to 0 + if "Thread" in event: + print("something is wrong") + thread = event["Thread"] else: thread = 0 + # event["Thread"] = 0 # assign a unique id to the event - event["id"] = self.__get_unique_id() + event["unique_id"] = self.__get_unique_id() + process_events = self.events[process] + process_stacks = self.stacks[process] + # get event list - if thread not in self.events[process]: - self.events[process][thread] = [] - event_list: List[Dict] = self.events[process][thread] + if thread not in process_events: + process_events[thread] = [] + event_list = process_events[thread] # get stack - if thread not in self.stacks[process]: - self.stacks[process][thread] = [] - stack: List[int] = self.stacks[process][thread] + if thread not in process_stacks: + process_stacks[thread] = [] + stack: List[int] = process_stacks[thread] - # if the event is an enter event, add the event to the stack and CCT + # if the event is an enter event, add the event to the stack and update the parent-child relationships if event["Event Type"] == "Enter": - cct = None - # if we are creating a CCT, get the correct CCT - if self.create_cct: - if thread not in self.ccts[process]: - self.ccts[process][thread] = Graph() - cct = self.ccts[process][thread] - self.__update_cct_and_parent_child_relationships(event, self.stacks[process][thread], event_list, cct) + self.__update_parent_child_relationships(event, stack, event_list) + # if the event is a leave event, update the matching event and pop from the stack elif event["Event Type"] == "Leave": - self.__update_match_event(event, self.stacks[process][thread], event_list) + self.__update_match_event(event, stack, event_list) event_list.append(event) + x = 0 - - def finalize(self) -> None: + def finalize_process(self, process: int) -> pd.DataFrame: # first step put everything in one list + # all_events = [] + # for process in self.events: + # for thread in process: + # all_events.extend(process[thread]) + + # convert 3d list of events to 1d list all_events = [] - for process in self.events: - for thread in process: - all_events.extend(process[thread]) + for thread_id in self.events[process]: + all_events.extend(self.events[process][thread_id]) + # df = pd.DataFrame(self.events[proc_id][thread_id]) + # just_for_break = 0 + + # for i in range(len(self.events[proc_id][thread_id])): + # all_events.append(self.events[proc_id][thread_id][i]) + # pass + # print(self.events[i][j][k]['Process']) + + # print('all_events has length: ' + str(len(all_events))) + # for i in range(len(self.events)): + # print(f'self.events[{i}] has length', len(self.events[i].keys())) + # # print(type(self.events[i])) + # # print (self.events[i].keys()) + # for j in self.events[i]: + # # print(j) + # # print (j in self.events[i]) + # # print(self.events[i][j]) + # print(f'self.events[{i}][{j}] has length', len(self.events[i][j])) + + # df_list = [] + # for process in self.events: + # for thread in process: + # df_list.append(pd.DataFrame(process[thread])) + # all_events = pd.concat(df_list) + # create a dataframe - self.events_dataframe = pd.DataFrame(all_events) - self.trace = Trace(None, self.events_dataframe, None) + df = pd.DataFrame(all_events) + # print(df.head()) + # print(self.events_dataframe["unique_id"].value_counts()) + return df + # self.events_dataframe = pandas.DataFrame(all_events) + # print number of events per id + # print(self.events_dataframe.sort_values(by=["unique_id"])) + # self.trace = Trace(None, self.events_dataframe, None) # Helper methods # This method can be thought of the update upon an "Enter" event # It adds to the stack and CCT - def __update_cct_and_parent_child_relationships(self, event: Dict, stack: List[int], event_list: List[Dict], cct: Graph) -> None: + def __update_parent_child_relationships(self, event: Dict, stack: List[int], event_list: List[Dict]) -> None: if len(stack) == 0: # root event event["parent"] = numpy.nan - # if self.create_cct: - # new_graph_node = Node(event["id"], None) - # cct.add_root(new_graph_node) - # event["Node"] = new_graph_node else: parent_event = event_list[stack[-1]] - event["parent"] = parent_event["id"] - if self.create_cct: - parent_graph_node = parent_event["Node"] - # if - # new_graph_node = Node(event["id"], parent_event["Node"]) - # parent_event["Node"].add_child(new_graph_node) - # event["Node"] = new_graph_node - - # update stack and event list + event["parent"] = parent_event["unique_id"] + + + # update stack stack.append(len(event_list)) - # event_list.append(event) - - # def __update_cct(self, event: Dict, process: int, thread: int) -> None: - # pass # This method can be thought of the update upon a "Leave" event # It pops from the stack and updates the event list @@ -122,23 +146,20 @@ def __update_cct_and_parent_child_relationships(self, event: Dict, stack: List[i def __update_match_event(self, leave_event: Dict, stack: List[int], event_list: List[Dict]) -> None: while len(stack) > 0: - enter_event = event_list[stack[-1]] + + # popping matched events from the stack + enter_event = event_list[stack.pop(-1)] + if enter_event["Name"] == leave_event["Name"]: # matching event found # update matching event ids - leave_event["_matching_event"] = enter_event["id"] - enter_event["_matching_event"] = leave_event["id"] + leave_event["_matching_event"] = enter_event["unique_id"] + enter_event["_matching_event"] = leave_event["unique_id"] - # popping matched events from the stack - stack.pop() break - else: - # popping unmatched events from the stack - stack.pop() - # event_list.append(leave_event) def __get_unique_id(self) -> int: self.unique_id += 1 From 65e433232c7f7c205c2cca0472f5ace8ac94425a Mon Sep 17 00:00:00 2001 From: Alexander Movsesyan Date: Sun, 13 Oct 2024 19:47:02 -0400 Subject: [PATCH 05/11] changes to make core reader more parallel-friendly --- pipit/readers/base_reader.py | 166 ----------------------------------- pipit/readers/core_reader.py | 124 ++++++++++++++++++++++++++ 2 files changed, 124 insertions(+), 166 deletions(-) delete mode 100644 pipit/readers/base_reader.py create mode 100644 pipit/readers/core_reader.py diff --git a/pipit/readers/base_reader.py b/pipit/readers/base_reader.py deleted file mode 100644 index 89057ab3..00000000 --- a/pipit/readers/base_reader.py +++ /dev/null @@ -1,166 +0,0 @@ -from abc import ABC, abstractmethod -from typing import List, Dict - -import pandas -import pandas as pd - -from .. import Trace -from ..graph import Graph, Node -import numpy - - -class BaseTraceReader(ABC): - - @abstractmethod - def read(self) -> Trace: - pass - - - # The following methods should be called by each reader class - def create_empty_trace(self, num_processes: int) -> None: - # keep track of a unique id for each event - self.unique_id = -1 - - # events are indexed by process number, then thread number - # stores a list of events - self.events: List[Dict[int, List[Dict]]] = [] - for i in range(num_processes): - self.events.append({}) - - # stacks are indexed by process number, then thread number - # stores indices of events in the event list - self.stacks: List[Dict[int, List[int]]] = [{}] * num_processes - - - - def add_event(self, event: Dict) -> None: - - # get process number -- if not present, set to 0 - if "Process" in event: - process = event["Process"] - else: - print("something is wrong") - process = 0 - - # get thread number -- if not present, set to 0 - if "Thread" in event: - print("something is wrong") - thread = event["Thread"] - else: - thread = 0 - # event["Thread"] = 0 - - # assign a unique id to the event - event["unique_id"] = self.__get_unique_id() - - - process_events = self.events[process] - process_stacks = self.stacks[process] - - # get event list - if thread not in process_events: - process_events[thread] = [] - event_list = process_events[thread] - - # get stack - if thread not in process_stacks: - process_stacks[thread] = [] - stack: List[int] = process_stacks[thread] - - # if the event is an enter event, add the event to the stack and update the parent-child relationships - if event["Event Type"] == "Enter": - self.__update_parent_child_relationships(event, stack, event_list) - # if the event is a leave event, update the matching event and pop from the stack - elif event["Event Type"] == "Leave": - self.__update_match_event(event, stack, event_list) - - event_list.append(event) - x = 0 - - def finalize_process(self, process: int) -> pd.DataFrame: - # first step put everything in one list - # all_events = [] - # for process in self.events: - # for thread in process: - # all_events.extend(process[thread]) - - # convert 3d list of events to 1d list - all_events = [] - for thread_id in self.events[process]: - all_events.extend(self.events[process][thread_id]) - # df = pd.DataFrame(self.events[proc_id][thread_id]) - # just_for_break = 0 - - # for i in range(len(self.events[proc_id][thread_id])): - # all_events.append(self.events[proc_id][thread_id][i]) - # pass - # print(self.events[i][j][k]['Process']) - - # print('all_events has length: ' + str(len(all_events))) - # for i in range(len(self.events)): - # print(f'self.events[{i}] has length', len(self.events[i].keys())) - # # print(type(self.events[i])) - # # print (self.events[i].keys()) - # for j in self.events[i]: - # # print(j) - # # print (j in self.events[i]) - # # print(self.events[i][j]) - # print(f'self.events[{i}][{j}] has length', len(self.events[i][j])) - - # df_list = [] - # for process in self.events: - # for thread in process: - # df_list.append(pd.DataFrame(process[thread])) - # all_events = pd.concat(df_list) - - # create a dataframe - df = pd.DataFrame(all_events) - # print(df.head()) - # print(self.events_dataframe["unique_id"].value_counts()) - return df - # self.events_dataframe = pandas.DataFrame(all_events) - # print number of events per id - # print(self.events_dataframe.sort_values(by=["unique_id"])) - # self.trace = Trace(None, self.events_dataframe, None) - - # Helper methods - - # This method can be thought of the update upon an "Enter" event - # It adds to the stack and CCT - def __update_parent_child_relationships(self, event: Dict, stack: List[int], event_list: List[Dict]) -> None: - if len(stack) == 0: - # root event - event["parent"] = numpy.nan - else: - parent_event = event_list[stack[-1]] - event["parent"] = parent_event["unique_id"] - - - # update stack - stack.append(len(event_list)) - - - # This method can be thought of the update upon a "Leave" event - # It pops from the stack and updates the event list - # We should look into using this function to add artificial "Leave" events for unmatched "Enter" events - def __update_match_event(self, leave_event: Dict, stack: List[int], event_list: List[Dict]) -> None: - - while len(stack) > 0: - - # popping matched events from the stack - enter_event = event_list[stack.pop(-1)] - - - if enter_event["Name"] == leave_event["Name"]: - # matching event found - - # update matching event ids - leave_event["_matching_event"] = enter_event["unique_id"] - enter_event["_matching_event"] = leave_event["unique_id"] - - break - - - def __get_unique_id(self) -> int: - self.unique_id += 1 - return self.unique_id \ No newline at end of file diff --git a/pipit/readers/core_reader.py b/pipit/readers/core_reader.py new file mode 100644 index 00000000..9c08fcb5 --- /dev/null +++ b/pipit/readers/core_reader.py @@ -0,0 +1,124 @@ +from abc import ABC, abstractmethod +from typing import List, Dict + +import pandas +import numpy + + +class CoreTraceReader: + """ + Helper Object to read traces from different sources and convert them into a common format + """ + + def __init__(self): + """ + Should be called by each process to create an empty trace per process in the reader. Creates the following + data structures to represent an empty trace: + - events: Dict[int, Dict[int, List[Dict]]] + - stacks: Dict[int, Dict[int, List[int]]] + """ + # keep track of a unique id for each event + self.unique_id = -1 + + # events are indexed by process number, then thread number + # stores a list of events + self.events: Dict[int, Dict[int, List[Dict]]] = {} + + # stacks are indexed by process number, then thread number + # stores indices of events in the event list + self.stacks: Dict[int, Dict[int, List[int]]] = {} + + def add_event(self, event: Dict) -> None: + """ + Should be called to add each event to the trace. Will update the event lists and stacks accordingly. + """ + # get process number -- if not present, set to 0 + if "Process" in event: + process = event["Process"] + else: + process = 0 + + # get thread number -- if not present, set to 0 + if "Thread" in event: + thread = event["Thread"] + else: + thread = 0 + # event["Thread"] = 0 + + # assign a unique id to the event + event["unique_id"] = self.__get_unique_id() + + # get event list + if process not in self.events: + self.events[process] = {} + if thread not in self.events[process]: + self.events[process][thread] = [] + event_list = self.events[process][thread] + + # get stack + if process not in self.stacks: + self.stacks[process] = {} + if thread not in self.stacks[process]: + self.stacks[process][thread] = [] + stack: List[int] = self.stacks[process][thread] + + # if the event is an enter event, add the event to the stack and update the parent-child relationships + if event["Event Type"] == "Enter": + self.__update_parent_child_relationships(event, stack, event_list) + # if the event is a leave event, update the matching event and pop from the stack + elif event["Event Type"] == "Leave": + self.__update_match_event(event, stack, event_list) + + # Finally add the event to the event list + event_list.append(event) + + def finalize(self): + """ + Converts the events data structure into a pandas dataframe and returns it + """ + all_events = [] + for process in self.events: + for thread in self.events[process]: + all_events.extend(self.events[process][thread]) + + # create a dataframe + events_dataframe = pandas.DataFrame(all_events) + return events_dataframe + + def __update_parent_child_relationships(self, event: Dict, stack: List[int], event_list: List[Dict]) -> None: + """ + This method can be thought of the update upon an "Enter" event. It adds to the stack and CCT + """ + if len(stack) == 0: + # root event + event["parent"] = numpy.nan + else: + parent_event = event_list[stack[-1]] + event["parent"] = parent_event["unique_id"] + + # update stack + stack.append(len(event_list)) + + def __update_match_event(self, leave_event: Dict, stack: List[int], event_list: List[Dict]) -> None: + """ + This method can be thought of the update upon a "Leave" event. It pops from the stack and updates the event list. + We should look into using this function to add artificial "Leave" events for unmatched "Enter" events + """ + + while len(stack) > 0: + + # popping matched events from the stack + enter_event = event_list[stack.pop(-1)] + + if enter_event["Name"] == leave_event["Name"]: + # matching event found + + # update matching event ids + leave_event["_matching_event"] = enter_event["unique_id"] + enter_event["_matching_event"] = leave_event["unique_id"] + + break + + def __get_unique_id(self) -> int: + self.unique_id += 1 + return self.unique_id From 2195cffbe25d0505c5df81fede693d79c69bcdf0 Mon Sep 17 00:00:00 2001 From: Alexander Movsesyan Date: Sun, 13 Oct 2024 19:53:04 -0400 Subject: [PATCH 06/11] minor updates to remove re-reading of sts file --- pipit/readers/projections_reader.py | 186 ++++++++++++++-------------- 1 file changed, 92 insertions(+), 94 deletions(-) diff --git a/pipit/readers/projections_reader.py b/pipit/readers/projections_reader.py index 38153e8f..d0b0b5b6 100644 --- a/pipit/readers/projections_reader.py +++ b/pipit/readers/projections_reader.py @@ -83,8 +83,6 @@ class ProjectionsConstants: class STSReader: def __init__(self, file_location): - self.sts_file = open(file_location, "r") # self.chares = {} - # In 'self.entries', each entry stores (entry_name: str, chare_id: int) self.entries = {} @@ -94,7 +92,7 @@ def __init__(self, file_location): # Stores user stat names: {user_event_id: user stat name} self.user_stats = {} - self.read_sts_file() + self.read_sts_file(file_location) # to get name of entry print > def get_entry_name(self, entry_id): @@ -132,95 +130,94 @@ def get_num_perf_counts(self): def get_event_name(self, event_id): return self.user_events[event_id] - def read_sts_file(self): - for line in self.sts_file: - line_arr = line.split() - - # Note: I'm disregarding TOTAL_STATS and TOTAL_EVENTS, because - # projections reader disregards them - - # Note: currently not reading/storing VERSION, MACHINE, SMPMODE, - # COMMANDLINE, CHARMVERSION, USERNAME, HOSTNAME - - # create chares array - # In 'self.chares', each entry stores (chare_name: str, dimension: int) - if line_arr[0] == "TOTAL_CHARES": - total_chares = int(line_arr[1]) - self.chares = [None] * total_chares - - elif line_arr[0] == "TOTAL_EPS": - self.num_eps = int(line_arr[1]) - - # get num processors - elif line_arr[0] == "PROCESSORS": - self.num_pes = int(line_arr[1]) - - # create message array - elif line_arr[0] == "TOTAL_MSGS": - total_messages = int(line_arr[1]) - self.message_table = [None] * total_messages - elif line_arr[0] == "TIMESTAMP": - self.timestamp_string = line_arr[1] - - # Add to self.chares - elif line_arr[0] == "CHARE": - id = int(line_arr[1]) - name = line_arr[2][1 : len(line_arr[2]) - 1] - dimensions = int(line_arr[3]) - self.chares[id] = (name, dimensions) - # print(int(line_arr[1]), line_arr[2][1:len(line_arr[2]) - 1]) - - # add to self.entries - elif line_arr[0] == "ENTRY": - # Need to concat entry_name - while not line_arr[3].endswith('"'): - line_arr[3] = line_arr[3] + " " + line_arr[4] - del line_arr[4] - - id = int(line_arr[2]) - entry_name = line_arr[3][1 : len(line_arr[3]) - 1] - chare_id = int(line_arr[4]) - # name = self.chares[chare_id][0] + '::' + entry_name - self.entries[id] = (entry_name, chare_id) - - # Add to message_table - # Need clarification on this, as message_table is never referenced in - # projections - elif line_arr[0] == "MESSAGE": - id = int(line_arr[1]) - message_size = int(line_arr[2]) - self.message_table[id] = message_size - - # Read/store event - elif line_arr[0] == "EVENT": - id = int(line_arr[1]) - event_name = "" - # rest of line is the event name - for i in range(2, len(line_arr)): - event_name = event_name + line_arr[i] + " " - self.user_events[id] = event_name - - # Read/store user stat - elif line_arr[0] == "STAT": - id = int(line_arr[1]) - event_name = "" - # rest of line is the stat - for i in range(2, len(line_arr)): - event_name = event_name + line_arr[i] + " " - self.user_stats[id] = event_name - - # create papi array - elif line_arr[0] == "TOTAL_PAPI_EVENTS": - num_papi_events = int(line_arr[1]) - self.papi_event_names = [None] * num_papi_events - - # Unsure of what these are for - elif line_arr[0] == "PAPI_EVENT": - id = int(line_arr[1]) - papi_event = line_arr[2] - self.papi_event_names[id] = papi_event - - self.sts_file.close() + def read_sts_file(self, file_path): + with open(file_path, "r") as sts_file: + for line in sts_file: + line_arr = line.split() + + # Note: I'm disregarding TOTAL_STATS and TOTAL_EVENTS, because + # projections reader disregards them + + # Note: currently not reading/storing VERSION, MACHINE, SMPMODE, + # COMMANDLINE, CHARMVERSION, USERNAME, HOSTNAME + + # create chares array + # In 'self.chares', each entry stores (chare_name: str, dimension: int) + if line_arr[0] == "TOTAL_CHARES": + total_chares = int(line_arr[1]) + self.chares = [None] * total_chares + + elif line_arr[0] == "TOTAL_EPS": + self.num_eps = int(line_arr[1]) + + # get num processors + elif line_arr[0] == "PROCESSORS": + self.num_pes = int(line_arr[1]) + + # create message array + elif line_arr[0] == "TOTAL_MSGS": + total_messages = int(line_arr[1]) + self.message_table = [None] * total_messages + elif line_arr[0] == "TIMESTAMP": + self.timestamp_string = line_arr[1] + + # Add to self.chares + elif line_arr[0] == "CHARE": + id = int(line_arr[1]) + name = line_arr[2][1 : len(line_arr[2]) - 1] + dimensions = int(line_arr[3]) + self.chares[id] = (name, dimensions) + # print(int(line_arr[1]), line_arr[2][1:len(line_arr[2]) - 1]) + + # add to self.entries + elif line_arr[0] == "ENTRY": + # Need to concat entry_name + while not line_arr[3].endswith('"'): + line_arr[3] = line_arr[3] + " " + line_arr[4] + del line_arr[4] + + id = int(line_arr[2]) + entry_name = line_arr[3][1 : len(line_arr[3]) - 1] + chare_id = int(line_arr[4]) + # name = self.chares[chare_id][0] + '::' + entry_name + self.entries[id] = (entry_name, chare_id) + + # Add to message_table + # Need clarification on this, as message_table is never referenced in + # projections + elif line_arr[0] == "MESSAGE": + id = int(line_arr[1]) + message_size = int(line_arr[2]) + self.message_table[id] = message_size + + # Read/store event + elif line_arr[0] == "EVENT": + id = int(line_arr[1]) + event_name = "" + # rest of line is the event name + for i in range(2, len(line_arr)): + event_name = event_name + line_arr[i] + " " + self.user_events[id] = event_name + + # Read/store user stat + elif line_arr[0] == "STAT": + id = int(line_arr[1]) + event_name = "" + # rest of line is the stat + for i in range(2, len(line_arr)): + event_name = event_name + line_arr[i] + " " + self.user_stats[id] = event_name + + # create papi array + elif line_arr[0] == "TOTAL_PAPI_EVENTS": + num_papi_events = int(line_arr[1]) + self.papi_event_names = [None] * num_papi_events + + # Unsure of what these are for + elif line_arr[0] == "PAPI_EVENT": + id = int(line_arr[1]) + papi_event = line_arr[2] + self.papi_event_names[id] = papi_event class ProjectionsReader: @@ -247,7 +244,8 @@ def __init__( if not hasattr(self, "executable_location"): raise ValueError("Invalid directory for projections - no sts files found.") - self.num_pes = STSReader(self.executable_location + ".sts").num_pes + self.sts_reader = STSReader(self.executable_location + ".sts") + self.num_pes = self.sts_reader.num_pes # make sure all the log files exist for i in range(self.num_pes): @@ -328,7 +326,7 @@ def read(self): def _read_log_file(self, rank_size) -> pd.DataFrame: # has information needed in sts file - sts_reader = STSReader(self.executable_location + ".sts") + sts_reader = self.sts_reader rank, size = rank_size[0], rank_size[1] per_process = int(self.num_pes // size) From dac0b9c42875927e63fd06b5458e6050d620d044 Mon Sep 17 00:00:00 2001 From: Alexander Movsesyan Date: Sun, 13 Oct 2024 21:46:50 -0400 Subject: [PATCH 07/11] updates to have strided unique_id, take care of instant events, and to concat trace dfs --- pipit/readers/core_reader.py | 59 ++++++++++++++++++++++++++++-------- 1 file changed, 47 insertions(+), 12 deletions(-) diff --git a/pipit/readers/core_reader.py b/pipit/readers/core_reader.py index 9c08fcb5..8fd8aa73 100644 --- a/pipit/readers/core_reader.py +++ b/pipit/readers/core_reader.py @@ -3,6 +3,7 @@ import pandas import numpy +from pipit.trace import Trace class CoreTraceReader: @@ -10,15 +11,18 @@ class CoreTraceReader: Helper Object to read traces from different sources and convert them into a common format """ - def __init__(self): + def __init__(self, start: int = 0, stride: int = 1): """ Should be called by each process to create an empty trace per process in the reader. Creates the following data structures to represent an empty trace: - events: Dict[int, Dict[int, List[Dict]]] - stacks: Dict[int, Dict[int, List[int]]] """ + # keep stride for how much unique id should be incremented + self.stride = stride + # keep track of a unique id for each event - self.unique_id = -1 + self.unique_id = start - self.stride # events are indexed by process number, then thread number # stores a list of events @@ -64,7 +68,9 @@ def add_event(self, event: Dict) -> None: # if the event is an enter event, add the event to the stack and update the parent-child relationships if event["Event Type"] == "Enter": - self.__update_parent_child_relationships(event, stack, event_list) + self.__update_parent_child_relationships(event, stack, event_list, False) + elif event["Event Type"] == "Instant": + self.__update_parent_child_relationships(event, stack, event_list, True) # if the event is a leave event, update the matching event and pop from the stack elif event["Event Type"] == "Leave": self.__update_match_event(event, stack, event_list) @@ -82,22 +88,35 @@ def finalize(self): all_events.extend(self.events[process][thread]) # create a dataframe - events_dataframe = pandas.DataFrame(all_events) - return events_dataframe - - def __update_parent_child_relationships(self, event: Dict, stack: List[int], event_list: List[Dict]) -> None: + trace_df = pandas.DataFrame(all_events) + + # categorical for memory savings + trace_df = trace_df.astype( + { + "Name": "category", + "Event Type": "category", + "Process": "category", + "_matching_event": "Int32", + "_parent": "Int32", + "_matching_timestamp": "Int32", + } + ) + return trace_df + + def __update_parent_child_relationships(self, event: Dict, stack: List[int], event_list: List[Dict],is_instant: bool) -> None: """ This method can be thought of the update upon an "Enter" event. It adds to the stack and CCT """ if len(stack) == 0: # root event - event["parent"] = numpy.nan + event["_parent"] = numpy.nan else: parent_event = event_list[stack[-1]] - event["parent"] = parent_event["unique_id"] + event["_parent"] = parent_event["unique_id"] # update stack - stack.append(len(event_list)) + if not is_instant: + stack.append(len(event_list)) def __update_match_event(self, leave_event: Dict, stack: List[int], event_list: List[Dict]) -> None: """ @@ -108,7 +127,7 @@ def __update_match_event(self, leave_event: Dict, stack: List[int], event_list: while len(stack) > 0: # popping matched events from the stack - enter_event = event_list[stack.pop(-1)] + enter_event = event_list[stack.pop()] if enter_event["Name"] == leave_event["Name"]: # matching event found @@ -117,8 +136,24 @@ def __update_match_event(self, leave_event: Dict, stack: List[int], event_list: leave_event["_matching_event"] = enter_event["unique_id"] enter_event["_matching_event"] = leave_event["unique_id"] + # update matching timestamps + leave_event["_matching_timestamp"] = enter_event["Timestamp (ns)"] + enter_event["_matching_timestamp"] = leave_event["Timestamp (ns)"] + break def __get_unique_id(self) -> int: - self.unique_id += 1 + self.unique_id += self.stride return self.unique_id + +def concat_trace_data(data_list): + """ + Concatenates the data from multiple trace readers into a single trace reader + """ + trace_data = pandas.concat(data_list, ignore_index=True) + # set index to unique_id + trace_data.set_index("unique_id", inplace=True) + trace_data.sort_values( + by="Timestamp (ns)", axis=0, ascending=True, inplace=True, ignore_index=True + ) + return Trace(None, trace_data, None) From 635cd2d59ef94c536fe017d829218453dc94a12e Mon Sep 17 00:00:00 2001 From: Alexander Movsesyan Date: Sun, 13 Oct 2024 21:45:39 -0400 Subject: [PATCH 08/11] updates to use core reader --- pipit/readers/projections_reader.py | 151 ++++++++++++---------------- 1 file changed, 67 insertions(+), 84 deletions(-) diff --git a/pipit/readers/projections_reader.py b/pipit/readers/projections_reader.py index d0b0b5b6..a91d1800 100644 --- a/pipit/readers/projections_reader.py +++ b/pipit/readers/projections_reader.py @@ -5,10 +5,15 @@ import os import gzip + +from numba.cuda import event + import pipit.trace import pandas as pd import multiprocessing as mp +from pipit.readers.core_reader import CoreTraceReader, concat_trace_data + class ProjectionsConstants: """ @@ -292,39 +297,18 @@ def read(self): pool_size, pool = self.num_processes, mp.Pool(self.num_processes) # Read each log file and store as list of dataframes - dataframes_list = pool.map( + data_list = pool.map( self._read_log_file, [(rank, pool_size) for rank in range(pool_size)] ) pool.close() # Concatenate the dataframes list into dataframe containing entire trace - trace_df = pd.concat(dataframes_list, ignore_index=True) - trace_df.sort_values( - by="Timestamp (ns)", axis=0, ascending=True, inplace=True, ignore_index=True - ) - - # categorical for memory savings - trace_df = trace_df.astype( - { - "Name": "category", - "Event Type": "category", - "Process": "category", - } - ) - # re-order columns - trace_df = trace_df[ - ["Timestamp (ns)", "Event Type", "Name", "Process", "Attributes"] - ] - - trace = pipit.trace.Trace(None, trace_df) - if self.create_cct: - trace.create_cct() - - return trace + return concat_trace_data(data_list) def _read_log_file(self, rank_size) -> pd.DataFrame: + # has information needed in sts file sts_reader = self.sts_reader @@ -332,6 +316,9 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: per_process = int(self.num_pes // size) remainder = int(self.num_pes % size) + # Start Core Reader + core_reader = CoreTraceReader(rank, size) + if rank < remainder: begin_int = rank * (per_process + 1) end_int = (rank + 1) * (per_process + 1) @@ -339,10 +326,7 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: begin_int = (rank * per_process) + remainder end_int = ((rank + 1) * per_process) + remainder - dfs = [] for pe_num in range(begin_int, end_int, 1): - # create an empty dict to append to - data = self._create_empty_dict() # opening the log file we need to read log_file = gzip.open( @@ -363,7 +347,7 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: details = {"From PE": pe} - _add_to_trace_dict(data, "Idle", "Enter", time, pe_num, details) + _add_to_trace(core_reader, "Idle", "Enter", time, pe_num, details) elif int(line_arr[0]) == ProjectionsConstants.END_IDLE: time = int(line_arr[1]) * 1000 @@ -371,7 +355,7 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: details = {"From PE": pe} - _add_to_trace_dict(data, "Idle", "Leave", time, pe_num, details) + _add_to_trace(core_reader, "Idle", "Leave", time, pe_num, details) # Pack message to be sent elif int(line_arr[0]) == ProjectionsConstants.BEGIN_PACK: @@ -380,7 +364,7 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: details = {"From PE": pe} - _add_to_trace_dict(data, "Pack", "Enter", time, pe_num, details) + _add_to_trace(core_reader, "Pack", "Enter", time, pe_num, details) elif int(line_arr[0]) == ProjectionsConstants.END_PACK: time = int(line_arr[1]) * 1000 @@ -388,7 +372,7 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: details = {"From PE": pe} - _add_to_trace_dict(data, "Pack", "Leave", time, pe_num, details) + _add_to_trace(core_reader, "Pack", "Leave", time, pe_num, details) # Unpacking a received message elif int(line_arr[0]) == ProjectionsConstants.BEGIN_UNPACK: @@ -397,7 +381,7 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: details = {"From PE": pe} - _add_to_trace_dict(data, "Unpack", "Enter", time, pe_num, details) + _add_to_trace(core_reader, "Unpack", "Enter", time, pe_num, details) elif int(line_arr[0]) == ProjectionsConstants.END_UNPACK: time = int(line_arr[1]) * 1000 @@ -405,14 +389,14 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: details = {"From PE": pe} - _add_to_trace_dict(data, "Unpack", "Leave", time, pe_num, details) + _add_to_trace(core_reader, "Unpack", "Leave", time, pe_num, details) elif int(line_arr[0]) == ProjectionsConstants.USER_SUPPLIED: user_supplied = line_arr[1] details = {"User Supplied": user_supplied} - _add_to_trace_dict( - data, "User Supplied", "Instant", -1, pe_num, details + _add_to_trace( + core_reader, "User Supplied", "Instant", -1, pe_num, details ) elif int(line_arr[0]) == ProjectionsConstants.USER_SUPPLIED_NOTE: @@ -423,8 +407,8 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: details = {"Note": note} - _add_to_trace_dict( - data, "User Supplied Note", "Instant", time, pe_num, details + _add_to_trace( + core_reader, "User Supplied Note", "Instant", time, pe_num, details ) # Not sure if this should be instant or enter/leave @@ -446,8 +430,8 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: "Note": note, } - _add_to_trace_dict( - data, + _add_to_trace( + core_reader, "User Supplied Bracketed Note", "Enter", time, @@ -455,8 +439,8 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: details, ) - _add_to_trace_dict( - data, + _add_to_trace( + core_reader, "User Supplied Bracketed Note", "Leave", end_time, @@ -471,8 +455,8 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: details = {"Memory Usage": memory_usage} - _add_to_trace_dict( - data, "Memory Usage", "Instant", time, pe_num, details + _add_to_trace( + core_reader, "Memory Usage", "Instant", time, pe_num, details ) # New chare create message being sent @@ -494,8 +478,8 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: "Send Time": send_time, } - _add_to_trace_dict( - data, + _add_to_trace( + core_reader, sts_reader.get_entry_name(entry), "Instant", time, @@ -526,8 +510,8 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: "Destinatopn PEs": destPEs, } - _add_to_trace_dict( - data, + _add_to_trace( + core_reader, sts_reader.get_entry_name(entry), "Instant", time, @@ -567,8 +551,8 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: "perf counts list": perf_counts, } - _add_to_trace_dict( - data, + _add_to_trace( + core_reader, sts_reader.get_entry_name(entry), "Enter", time, @@ -599,8 +583,8 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: "perf counts list": perf_counts, } - _add_to_trace_dict( - data, + _add_to_trace( + core_reader, sts_reader.get_entry_name(entry), "Leave", time, @@ -612,12 +596,12 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: elif int(line_arr[0]) == ProjectionsConstants.BEGIN_TRACE: time = int(line_arr[1]) * 1000 - _add_to_trace_dict(data, "Trace", "Enter", time, pe_num, None) + _add_to_trace(core_reader, "Trace", "Enter", time, pe_num, None) elif int(line_arr[0]) == ProjectionsConstants.END_TRACE: time = int(line_arr[1]) * 1000 - _add_to_trace_dict(data, "Trace", "Leave", time, pe_num, None) + _add_to_trace(core_reader, "Trace", "Leave", time, pe_num, None) # Message Receive ? elif int(line_arr[0]) == ProjectionsConstants.MESSAGE_RECV: @@ -634,8 +618,8 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: "Message Length": message_length, } - _add_to_trace_dict( - data, "Message Receive", "Instant", time, pe_num, details + _add_to_trace( + core_reader, "Message Receive", "Instant", time, pe_num, details ) # queueing creation ? @@ -647,7 +631,7 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: details = {"From PE": pe, "Message Type": mtype, "Event ID": event} - _add_to_trace_dict(data, "Enque", "Instant", time, pe_num, details) + _add_to_trace(core_reader, "Enque", "Instant", time, pe_num, details) elif int(line_arr[0]) == ProjectionsConstants.DEQUEUE: mtype = int(line_arr[1]) @@ -657,7 +641,7 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: details = {"From PE": pe, "Message Type": mtype, "Event ID": event} - _add_to_trace_dict(data, "Deque", "Instant", time, pe_num, details) + _add_to_trace(core_reader, "Deque", "Instant", time, pe_num, details) # Interrupt from different chare ? elif int(line_arr[0]) == ProjectionsConstants.BEGIN_INTERRUPT: @@ -667,8 +651,8 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: details = {"From PE": pe, "Event ID": event} - _add_to_trace_dict( - data, "Interrupt", "Enter", time, pe_num, details + _add_to_trace( + core_reader, "Interrupt", "Enter", time, pe_num, details ) elif int(line_arr[0]) == ProjectionsConstants.END_INTERRUPT: @@ -678,20 +662,20 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: details = {"From PE": pe, "Event ID": event} - _add_to_trace_dict( - data, "Interrupt", "Leave", time, pe_num, details + _add_to_trace( + core_reader, "Interrupt", "Leave", time, pe_num, details ) # Very start of the program - encapsulates every other event elif int(line_arr[0]) == ProjectionsConstants.BEGIN_COMPUTATION: time = int(line_arr[1]) * 1000 - _add_to_trace_dict(data, "Computation", "Enter", time, pe_num, None) + _add_to_trace(core_reader, "Computation", "Enter", time, pe_num, None) elif int(line_arr[0]) == ProjectionsConstants.END_COMPUTATION: time = int(line_arr[1]) * 1000 - _add_to_trace_dict(data, "Computation", "Leave", time, pe_num, None) + _add_to_trace(core_reader, "Computation", "Leave", time, pe_num, None) # User event (in code) elif int(line_arr[0]) == ProjectionsConstants.USER_EVENT: @@ -708,8 +692,8 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: "Event Type": "User Event", } - _add_to_trace_dict( - data, user_event_name, "Instant", time, pe_num, details + _add_to_trace( + core_reader, user_event_name, "Instant", time, pe_num, details ) elif int(line_arr[0]) == ProjectionsConstants.USER_EVENT_PAIR: @@ -728,8 +712,8 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: "Event Type": "User Event Pair", } - _add_to_trace_dict( - data, user_event_name, "Instant", time, pe_num, details + _add_to_trace( + core_reader, user_event_name, "Instant", time, pe_num, details ) elif int(line_arr[0]) == ProjectionsConstants.BEGIN_USER_EVENT_PAIR: @@ -746,8 +730,8 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: "User Event Name": sts_reader.get_user_event(user_event_id), } - _add_to_trace_dict( - data, "User Event Pair", "Enter", time, pe_num, details + _add_to_trace( + core_reader, "User Event Pair", "Enter", time, pe_num, details ) elif int(line_arr[0]) == ProjectionsConstants.END_USER_EVENT_PAIR: @@ -764,7 +748,7 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: "User Event Name": sts_reader.get_user_event(user_event_id), } - _add_to_trace_dict( + _add_to_trace( "User Event Pair", "Leave", time, pe_num, details ) @@ -785,24 +769,23 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: "Event Type": "User Stat", } - _add_to_trace_dict( - data, user_stat_name, "Instant", time, pe_num, details + _add_to_trace( + core_reader, user_stat_name, "Instant", time, pe_num, details ) - # Making sure that the log file ends with END_COMPUTATION - if len(data["Name"]) > 0 and data["Name"][-1] != "Computation": - time = data["Timestamp (ns)"][-1] * 1000 - _add_to_trace_dict(data, "Computation", "Leave", time, pe_num, None) log_file.close() - dfs.append(pd.DataFrame(data)) - return pd.concat(dfs) + + return core_reader.finalize() -def _add_to_trace_dict(data, name, evt_type, time, process, attributes): - data["Name"].append(name) - data["Event Type"].append(evt_type) - data["Timestamp (ns)"].append(time) - data["Process"].append(process) - data["Attributes"].append(attributes) +def _add_to_trace(core_reader: CoreTraceReader, name, evt_type, time, process, attributes): + new_event = { + "Name": name, + "Event Type": evt_type, + "Timestamp (ns)": time, + "Process": process, + "Attributes": attributes, + } + core_reader.add_event(new_event) From 3daf5b92e9991236ffa5b07e67455585b30957fc Mon Sep 17 00:00:00 2001 From: Alexander Movsesyan Date: Sun, 13 Oct 2024 21:51:51 -0400 Subject: [PATCH 09/11] removed unneeded function --- pipit/readers/projections_reader.py | 11 ----------- 1 file changed, 11 deletions(-) diff --git a/pipit/readers/projections_reader.py b/pipit/readers/projections_reader.py index a91d1800..4850116c 100644 --- a/pipit/readers/projections_reader.py +++ b/pipit/readers/projections_reader.py @@ -276,17 +276,6 @@ def __init__( self.create_cct = create_cct - # Returns an empty dict, used for reading log file into dataframe - @staticmethod - def _create_empty_dict() -> dict: - return { - "Name": [], - "Event Type": [], - "Timestamp (ns)": [], - "Process": [], - "Attributes": [], - } - def read(self): if self.num_pes < 1: return None From 964a9b79ea17b5fc454f4360bf7a519c4dbaee0b Mon Sep 17 00:00:00 2001 From: Thomas Li <47963215+lithomas1@users.noreply.github.com> Date: Tue, 18 Feb 2025 23:25:29 -0500 Subject: [PATCH 10/11] cleanup --- pipit/readers/projections_reader.py | 3 --- 1 file changed, 3 deletions(-) diff --git a/pipit/readers/projections_reader.py b/pipit/readers/projections_reader.py index 1f5da13c..ded6bb12 100644 --- a/pipit/readers/projections_reader.py +++ b/pipit/readers/projections_reader.py @@ -6,9 +6,6 @@ import os import gzip -from numba.cuda import event - -import pipit.trace import pandas as pd import multiprocessing as mp From ad6349f27bf67d4fbab124b37ca3891eb2c7e363 Mon Sep 17 00:00:00 2001 From: Thomas Li <47963215+lithomas1@users.noreply.github.com> Date: Tue, 18 Feb 2025 23:32:05 -0500 Subject: [PATCH 11/11] lint --- pipit/readers/core_reader.py | 1 - pipit/readers/projections_reader.py | 33 +++++++++++++++++++---------- 2 files changed, 22 insertions(+), 12 deletions(-) diff --git a/pipit/readers/core_reader.py b/pipit/readers/core_reader.py index 4a636eb0..3f826449 100644 --- a/pipit/readers/core_reader.py +++ b/pipit/readers/core_reader.py @@ -110,7 +110,6 @@ def finalize(self): ) return trace_df - def __update_parent_child_relationships( self, event: Dict, stack: List[int], event_list: List[Dict], is_instant: bool ) -> None: diff --git a/pipit/readers/projections_reader.py b/pipit/readers/projections_reader.py index ded6bb12..09f55933 100644 --- a/pipit/readers/projections_reader.py +++ b/pipit/readers/projections_reader.py @@ -394,7 +394,12 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: details = {"Note": note} _add_to_trace( - core_reader, "User Supplied Note", "Instant", time, pe_num, details + core_reader, + "User Supplied Note", + "Instant", + time, + pe_num, + details, ) # Not sure if this should be instant or enter/leave @@ -617,7 +622,9 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: details = {"From PE": pe, "Message Type": mtype, "Event ID": event} - _add_to_trace(core_reader, "Enque", "Instant", time, pe_num, details) + _add_to_trace( + core_reader, "Enque", "Instant", time, pe_num, details + ) elif int(line_arr[0]) == ProjectionsConstants.DEQUEUE: mtype = int(line_arr[1]) @@ -627,7 +634,9 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: details = {"From PE": pe, "Message Type": mtype, "Event ID": event} - _add_to_trace(core_reader, "Deque", "Instant", time, pe_num, details) + _add_to_trace( + core_reader, "Deque", "Instant", time, pe_num, details + ) # Interrupt from different chare ? elif int(line_arr[0]) == ProjectionsConstants.BEGIN_INTERRUPT: @@ -656,12 +665,16 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: elif int(line_arr[0]) == ProjectionsConstants.BEGIN_COMPUTATION: time = int(line_arr[1]) * 1000 - _add_to_trace(core_reader, "Computation", "Enter", time, pe_num, None) + _add_to_trace( + core_reader, "Computation", "Enter", time, pe_num, None + ) elif int(line_arr[0]) == ProjectionsConstants.END_COMPUTATION: time = int(line_arr[1]) * 1000 - _add_to_trace(core_reader, "Computation", "Leave", time, pe_num, None) + _add_to_trace( + core_reader, "Computation", "Leave", time, pe_num, None + ) # User event (in code) elif int(line_arr[0]) == ProjectionsConstants.USER_EVENT: @@ -734,9 +747,7 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: "User Event Name": sts_reader.get_user_event(user_event_id), } - _add_to_trace( - "User Event Pair", "Leave", time, pe_num, details - ) + _add_to_trace("User Event Pair", "Leave", time, pe_num, details) # User stat (in code) elif int(line_arr[0]) == ProjectionsConstants.USER_STAT: @@ -759,14 +770,14 @@ def _read_log_file(self, rank_size) -> pd.DataFrame: core_reader, user_stat_name, "Instant", time, pe_num, details ) - log_file.close() - return core_reader.finalize() -def _add_to_trace(core_reader: CoreTraceReader, name, evt_type, time, process, attributes): +def _add_to_trace( + core_reader: CoreTraceReader, name, evt_type, time, process, attributes +): new_event = { "Name": name, "Event Type": evt_type,