|
1 | 1 | import concurrent.futures |
2 | 2 | import multiprocessing |
3 | | -from typing import Dict, Optional, List |
| 3 | +from typing import Dict, Optional, Union |
4 | 4 |
|
5 | 5 | import click |
6 | 6 | import numpy as np |
|
11 | 11 | from wta.helpers import print_section_boundaries, convert_timestamp_columns_to_datetime, log_ids_non_nil, \ |
12 | 12 | EventLogIDs, TRANSITION_COLUMN_KEY |
13 | 13 | from wta.waiting_time import analysis as wt_analysis |
| 14 | +from wta.calendars.calendars import make as make_calendar |
| 15 | + |
| 16 | + |
| 17 | +CONVERT_COLUMNS = ['wt_total', 'wt_contention', 'wt_batching', 'wt_prioritization', 'wt_unavailability', 'wt_extraneous'] |
| 18 | +ORDERED_COLUMNS = [ |
| 19 | + 'start_time', |
| 20 | + 'end_time', |
| 21 | + 'source_activity', |
| 22 | + 'source_resource', |
| 23 | + 'destination_activity', |
| 24 | + 'destination_resource', |
| 25 | + 'case_id', |
| 26 | + 'wt_total', |
| 27 | + 'wt_contention', |
| 28 | + 'wt_batching', |
| 29 | + 'wt_prioritization', |
| 30 | + 'wt_unavailability', |
| 31 | + 'wt_extraneous' |
| 32 | +] |
14 | 33 |
|
15 | 34 |
|
16 | 35 | @print_section_boundaries('Activity Transitions Analysis') |
17 | | -def identify( |
18 | | - log: pd.DataFrame, |
19 | | - parallel_activities: Dict[str, set], |
20 | | - parallel_run=True, |
21 | | - log_ids: Optional[EventLogIDs] = None, |
22 | | - calendar: Optional[Dict] = None, |
23 | | - group_results: bool = True) -> Optional[pd.DataFrame]: |
24 | | - from wta.calendars.calendars import make as make_calendar |
25 | | - |
| 36 | +def identify(log: pd.DataFrame, parallel_activities: Dict[str, set], parallel_run: bool = True, |
| 37 | + log_ids: Optional[EventLogIDs] = None, calendar: Optional[Dict] = None) -> Optional[pd.DataFrame]: |
26 | 38 | click.echo(f'Parallel run: {parallel_run}') |
27 | | - |
28 | 39 | log_ids = log_ids_non_nil(log_ids) |
| 40 | + log_calendar = make_calendar_if_none(log, log_ids, calendar) |
| 41 | + run_func = __multiprocess_run if parallel_run else __sequential_run |
| 42 | + all_items = run_func(log, log_ids, log_calendar, parallel_activities) |
| 43 | + return None if len(all_items) == 0 else process_all_items(all_items) |
29 | 44 |
|
30 | | - if not calendar: |
31 | | - log_calendar = make_calendar(log, granularity=GRANULARITY_MINUTES, log_ids=log_ids) |
32 | | - else: |
33 | | - log_calendar = calendar |
34 | 45 |
|
35 | | - if parallel_run: |
36 | | - all_items = __multiprocess_run(log, log_ids, log_calendar, parallel_activities) |
37 | | - else: |
38 | | - all_items = __sequential_run(log, log_ids, log_calendar, parallel_activities) |
| 46 | +def process_all_items(all_items: pd.DataFrame) -> pd.DataFrame: |
| 47 | + # Convert time columns to seconds |
| 48 | + columns_to_convert = [col for col in CONVERT_COLUMNS if col in all_items.columns] |
| 49 | + all_items[columns_to_convert] = all_items[columns_to_convert].applymap(lambda x: pd.to_timedelta(x).total_seconds()) |
39 | 50 |
|
40 | | - if len(all_items) == 0: |
41 | | - return None |
| 51 | + # Return the dataframe in the order of ORDERED_COLUMNS |
| 52 | + return all_items[ORDERED_COLUMNS] |
42 | 53 |
|
43 | | - if group_results: |
44 | | - result = __join_per_case_items(all_items, log_ids=log_ids) |
45 | | - if result is not None: |
46 | | - result['wt_total_seconds'] = result[log_ids.wt_total] / np.timedelta64(1, 's') |
47 | | - else: |
48 | | - result = __create_single_dataframe(all_items) |
49 | 54 |
|
50 | | - return result |
| 55 | +def make_calendar_if_none(log, log_ids, calendar): |
| 56 | + return make_calendar(log, granularity=GRANULARITY_MINUTES, log_ids=log_ids) if not calendar else calendar |
51 | 57 |
|
52 | 58 |
|
53 | 59 | def __sequential_run(log, log_ids, calendar, parallel_activities): |
54 | | - results = [] |
55 | 60 | log_grouped = log.groupby(by=log_ids.case) |
56 | | - |
57 | | - for (case_id, case) in log_grouped: |
58 | | - case = case.sort_values(by=[log_ids.end_time, log_ids.start_time]) |
59 | | - result = __identify_transitions_per_case_and_make_report( |
60 | | - case, |
61 | | - parallel_activities=parallel_activities, |
62 | | - case_id=case_id, |
63 | | - log_calendar=calendar, |
64 | | - log=log, |
65 | | - log_ids=log_ids) |
66 | | - |
67 | | - if result is not None: |
68 | | - results.append(result) |
69 | | - |
70 | | - return results |
| 61 | + results_transitions = [identify_transitions_and_report(sort_case(case, log_ids), parallel_activities, case_id, calendar, log, log_ids) |
| 62 | + for case_id, case in log_grouped] |
| 63 | + return concatenate_transitions_if_exists(results_transitions) |
71 | 64 |
|
72 | 65 |
|
73 | 66 | def __multiprocess_run(log, log_ids, calendar, parallel_activities): |
74 | | - all_items = [] |
75 | 67 | n_cores = multiprocessing.cpu_count() - 1 |
76 | 68 | handles = [] |
77 | 69 | log_grouped = log.groupby(by=log_ids.case) |
78 | 70 |
|
79 | 71 | with concurrent.futures.ProcessPoolExecutor(max_workers=n_cores) as executor: |
80 | | - for (case_id, case) in tqdm(log_grouped, desc='Submitting tasks for concurrent execution'): |
81 | | - case = case.sort_values(by=[log_ids.end_time, log_ids.start_time]) |
82 | | - handle = executor.submit(__identify_transitions_per_case_and_make_report, |
83 | | - case, |
84 | | - parallel_activities=parallel_activities, |
85 | | - case_id=case_id, |
86 | | - log_calendar=calendar, |
87 | | - log=log, |
88 | | - log_ids=log_ids) |
89 | | - handles.append(handle) |
90 | | - |
91 | | - for h in tqdm(handles, desc='Waiting for tasks to finish'): |
92 | | - done = h.done() |
93 | | - result = h.result() |
94 | | - if done and not result.empty: |
95 | | - all_items.append(result) |
| 72 | + handles = [submit_task(executor, sort_case(case, log_ids), parallel_activities, case_id, calendar, log, log_ids) |
| 73 | + for case_id, case in tqdm(log_grouped, desc='Submitting tasks for concurrent execution')] |
96 | 74 |
|
97 | | - return all_items |
| 75 | + all_transitions = [h.result() for h in tqdm(handles, desc='Waiting for tasks to finish') if not h.result().empty] |
| 76 | + return concatenate_transitions_if_exists(all_transitions) |
98 | 77 |
|
99 | 78 |
|
100 | | -def __identify_transitions_per_case_and_make_report(case: pd.DataFrame, **kwargs) -> pd.DataFrame: |
101 | | - parallel_activities = kwargs['parallel_activities'] |
102 | | - case_id = kwargs['case_id'] |
103 | | - log_calendar = kwargs['log_calendar'] |
104 | | - log = kwargs['log'] |
105 | | - log_ids = log_ids_non_nil(kwargs.get('log_ids')) |
| 79 | +def sort_case(case, log_ids): |
| 80 | + return case.sort_values(by=[log_ids.end_time, log_ids.start_time]) |
106 | 81 |
|
107 | | - case = case.sort_values(by=[log_ids.end_time, log_ids.start_time]).copy() |
108 | 82 |
|
109 | | - # converting timestamps to datetime |
110 | | - log = convert_timestamp_columns_to_datetime(log, log_ids) |
111 | | - case = convert_timestamp_columns_to_datetime(case, log_ids) |
112 | | - |
113 | | - __mark_activity_transitions(case, parallel_activities, log_ids=log_ids) |
114 | | - |
115 | | - transitions = wt_analysis.run(case, log_calendar, log, log_ids=log_ids) |
| 83 | +def submit_task(executor, case, parallel_activities, case_id, calendar, log, log_ids): |
| 84 | + return executor.submit(identify_transitions_and_report, case, parallel_activities, case_id, calendar, log, log_ids) |
116 | 85 |
|
117 | | - transitions_with_frequency = __calculate_frequency_and_duration(transitions, log_ids=log_ids) |
118 | 86 |
|
119 | | - # dropping edge cases with Start and End as an activity |
120 | | - starts_ends_values = ['Start', 'End'] |
121 | | - starts_and_ends = (transitions_with_frequency['source_activity'].isin(starts_ends_values) |
122 | | - & transitions_with_frequency['source_resource'].isin(starts_ends_values)) \ |
123 | | - | (transitions_with_frequency['destination_activity'].isin(starts_ends_values) |
124 | | - & transitions_with_frequency['destination_resource'].isin(starts_ends_values)) |
125 | | - transitions_with_frequency = transitions_with_frequency[starts_and_ends == False] |
| 87 | +def concatenate_transitions_if_exists(results_transitions): |
| 88 | + return pd.concat(results_transitions, ignore_index=True) if results_transitions else None |
126 | 89 |
|
127 | | - # attaching case ID as additional information |
128 | | - transitions_with_frequency['case_id'] = case_id |
129 | 90 |
|
130 | | - return transitions_with_frequency |
131 | | - |
132 | | - |
133 | | -def __mark_activity_transitions( |
134 | | - case: pd.DataFrame, |
135 | | - parallel_activities: Optional[Dict[str, set]] = None, |
136 | | - log_ids: Optional[EventLogIDs] = None): |
137 | | - log_ids = log_ids_non_nil(log_ids) |
138 | | - |
139 | | - # NOTE: we assume (a) the case was sorted by end time |
| 91 | +def identify_transitions_and_report(case, parallel_activities, case_id, log_calendar, log, log_ids): |
| 92 | + case = convert_timestamp_columns_to_datetime(case, log_ids) |
| 93 | + log = convert_timestamp_columns_to_datetime(log, log_ids) |
| 94 | + mark_activity_transitions(case, parallel_activities, log_ids=log_ids) |
| 95 | + transitions = wt_analysis.run(case, log_calendar, log, log_ids=log_ids) |
| 96 | + transitions['case_id'] = case_id |
| 97 | + return transitions |
140 | 98 |
|
141 | | - if not parallel_activities: |
142 | | - parallel_activities = {} |
143 | 99 |
|
| 100 | +def mark_activity_transitions(case, parallel_activities, log_ids): |
144 | 101 | case[log_ids.transition_source_index] = np.NAN |
145 | | - |
146 | | - # processing the case backwards |
147 | 102 | reversed_index = list(reversed(case.index)) |
148 | 103 | for i in range(len(reversed_index)): |
149 | 104 | non_concurrent_previous_event_found = False |
150 | | - |
151 | 105 | index = reversed_index[i] |
152 | 106 | current_event = case.loc[index] |
153 | 107 | parallel_activities_for_current_event = parallel_activities.get(current_event[log_ids.activity], []) |
154 | | - |
155 | 108 | previous_event_index_delta = i + 1 |
156 | | - while not non_concurrent_previous_event_found: |
157 | | - if previous_event_index_delta > len(reversed_index) - 1: |
158 | | - break |
159 | | - |
| 109 | + while not non_concurrent_previous_event_found and previous_event_index_delta <= len(reversed_index) - 1: |
160 | 110 | previous_event_index = reversed_index[previous_event_index_delta] |
161 | 111 | previous_event = case.loc[previous_event_index] |
162 | | - # Check if they are overlapping |
163 | 112 | overlapping_activity_instances = previous_event[log_ids.end_time] > current_event[log_ids.start_time] |
164 | | - if previous_event[ |
165 | | - log_ids.activity] in parallel_activities_for_current_event or overlapping_activity_instances: |
166 | | - # If they are concurrent activities, or overlapping instances, jump to the previous event |
| 113 | + if previous_event[log_ids.activity] in parallel_activities_for_current_event or overlapping_activity_instances: |
167 | 114 | previous_event_index_delta += 1 |
168 | 115 | else: |
169 | | - # If they are not concurrent nor overlapping, transition! |
170 | 116 | case.at[index, TRANSITION_COLUMN_KEY] = previous_event_index |
171 | 117 | non_concurrent_previous_event_found = True |
172 | | - |
173 | | - |
174 | | -def __calculate_frequency_and_duration(transitions: pd.DataFrame, |
175 | | - log_ids: Optional[EventLogIDs] = None) -> pd.DataFrame: |
176 | | - log_ids = log_ids_non_nil(log_ids) |
177 | | - |
178 | | - # calculating frequency per case of the transitions with the same activities and resources |
179 | | - columns = transitions.columns.tolist() |
180 | | - transition_with_frequency = pd.DataFrame(columns=columns) |
181 | | - for (pair, records) in transitions.groupby(by=['source_activity', 'source_resource', |
182 | | - 'destination_activity', 'destination_resource']): |
183 | | - transition_with_frequency = pd.concat([transition_with_frequency, pd.DataFrame({ |
184 | | - 'source_activity': [pair[0]], |
185 | | - 'source_resource': [pair[1]], |
186 | | - 'destination_activity': [pair[2]], |
187 | | - 'destination_resource': [pair[3]], |
188 | | - 'frequency': [len(records)], |
189 | | - 'transition_type': [records['transition_type'].iloc[0]], |
190 | | - log_ids.wt_total: [records[log_ids.wt_total].sum()], |
191 | | - log_ids.wt_batching: [records[log_ids.wt_batching].sum()], |
192 | | - log_ids.wt_prioritization: [records[log_ids.wt_prioritization].sum()], |
193 | | - log_ids.wt_contention: [records[log_ids.wt_contention].sum()], |
194 | | - log_ids.wt_unavailability: [records[log_ids.wt_unavailability].sum()], |
195 | | - log_ids.wt_extraneous: [records[log_ids.wt_extraneous].sum()], |
196 | | - })], ignore_index=True) |
197 | | - return transition_with_frequency |
198 | | - |
199 | | - |
200 | | -def __create_single_dataframe(items: List[pd.DataFrame]) -> Optional[pd.DataFrame]: |
201 | | - items = list(filter(lambda df: not df.empty, items)) |
202 | | - if len(items) == 0: |
203 | | - return None |
204 | | - else: |
205 | | - return pd.concat(items).reset_index(drop=True) |
206 | | - |
207 | | - |
208 | | -def __join_per_case_items(items: List[pd.DataFrame], log_ids: Optional[EventLogIDs] = None) -> Optional[pd.DataFrame]: |
209 | | - """Joins a list of items summing up frequency and duration.""" |
210 | | - |
211 | | - log_ids = log_ids_non_nil(log_ids) |
212 | | - |
213 | | - items = list(filter(lambda df: not df.empty, items)) |
214 | | - |
215 | | - if len(items) == 0: |
216 | | - return None |
217 | | - |
218 | | - columns = ['source_activity', 'source_resource', 'destination_activity', 'destination_resource'] |
219 | | - grouped = pd.concat(items).groupby(columns) |
220 | | - result = pd.DataFrame(columns=columns) |
221 | | - for pair_index, group in grouped: |
222 | | - source_activity, source_resource, destination_activity, destination_resource = pair_index |
223 | | - group_wt_total: pd.Timedelta = group[log_ids.wt_total].sum() |
224 | | - |
225 | | - group_wt_batching = pd.Timedelta(0) |
226 | | - if log_ids.wt_batching in group.columns: |
227 | | - group_wt_batching = group[log_ids.wt_batching].sum() |
228 | | - |
229 | | - group_wt_prioritization = pd.Timedelta(0) |
230 | | - if log_ids.wt_prioritization in group.columns: |
231 | | - group_wt_prioritization = group[log_ids.wt_prioritization].sum() |
232 | | - |
233 | | - group_wt_contention = pd.Timedelta(0) |
234 | | - if log_ids.wt_contention in group.columns: |
235 | | - group_wt_contention = pd.to_timedelta(group[log_ids.wt_contention]).sum() |
236 | | - |
237 | | - group_wt_unavailability = pd.Timedelta(0) |
238 | | - if log_ids.wt_unavailability in group.columns: |
239 | | - group_wt_unavailability = pd.to_timedelta(group[log_ids.wt_unavailability]).sum() |
240 | | - |
241 | | - group_wt_extraneous = pd.Timedelta(0) |
242 | | - if log_ids.wt_extraneous in group.columns: |
243 | | - group_wt_extraneous = pd.to_timedelta(group[log_ids.wt_extraneous]).sum() |
244 | | - |
245 | | - group_frequency: float = group['frequency'].sum() |
246 | | - group_case_id: str = ','.join(group['case_id'].astype(str).unique()) |
247 | | - result = pd.concat([result, pd.DataFrame({ |
248 | | - 'source_activity': [source_activity], |
249 | | - 'source_resource': [source_resource], |
250 | | - 'destination_activity': [destination_activity], |
251 | | - 'destination_resource': [destination_resource], |
252 | | - 'frequency': [group_frequency], |
253 | | - 'cases': [group_case_id], |
254 | | - log_ids.wt_total: [group_wt_total], |
255 | | - log_ids.wt_batching: [group_wt_batching], |
256 | | - log_ids.wt_prioritization: [group_wt_prioritization], |
257 | | - log_ids.wt_contention: [group_wt_contention], |
258 | | - log_ids.wt_unavailability: [group_wt_unavailability], |
259 | | - log_ids.wt_extraneous: [group_wt_extraneous] |
260 | | - })], ignore_index=True) |
261 | | - result.reset_index(drop=True, inplace=True) |
262 | | - return result |
0 commit comments