Function to transform log file to dataframe
(self, log_filepath)
| 39 | self.n_workers = n_workers |
| 40 | |
| 41 | def load_to_dataframe(self, log_filepath): |
| 42 | """ Function to transform log file to dataframe |
| 43 | """ |
| 44 | print('Loading log messages to dataframe...') |
| 45 | lines = [] |
| 46 | with open(log_filepath, 'r') as fid: |
| 47 | lines = fid.readlines() |
| 48 | |
| 49 | log_messages = [] |
| 50 | if self.n_workers == 1: |
| 51 | log_messages = formalize_message(enumerate(lines), self.regex, self.headers) |
| 52 | else: |
| 53 | chunk_size = np.ceil(len(lines) / float(self.n_workers)) |
| 54 | chunks = groupby(enumerate(lines), key=lambda k, line=count(): next(line)//chunk_size) |
| 55 | log_chunks = [list(chunk) for _, chunk in chunks] |
| 56 | print('Read %d log chunks in parallel'%len(log_chunks)) |
| 57 | pool = mp.Pool(processes=self.n_workers) |
| 58 | result_chunks = [pool.apply_async(formalize_message, args=(chunk, self.regex, self.headers)) |
| 59 | for chunk in log_chunks] |
| 60 | pool.close() |
| 61 | pool.join() |
| 62 | log_messages = list(chain(*[result.get() for result in result_chunks])) |
| 63 | |
| 64 | if not log_messages: |
| 65 | raise RuntimeError('Logformat error or log file is empty!') |
| 66 | log_dataframe = pd.DataFrame(log_messages, columns=['LineId'] + self.headers) |
| 67 | success_rate = len(log_messages) / float(len(lines)) |
| 68 | print('Loading {} messages done, loading rate: {:.1%}'.format(len(log_messages), success_rate)) |
| 69 | return log_dataframe |
| 70 | |
| 71 | def _generate_logformat_regex(self, logformat): |
| 72 | """ Function to generate regular expression to split log messages |
no test coverage detected