| 30 | import numpy as np |
| 31 | |
| 32 | class LogLoader(object): |
| 33 | |
| 34 | def __init__(self, logformat, n_workers=1): |
| 35 | if not logformat: |
| 36 | raise RuntimeError('Logformat is required!') |
| 37 | self.logformat = logformat.strip() |
| 38 | self.headers, self.regex = self._generate_logformat_regex(self.logformat) |
| 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 |
| 73 | """ |
| 74 | headers = [] |
| 75 | splitters = re.split(r'(<[^<>]+>)', logformat) |
| 76 | regex = '' |
| 77 | for k in range(len(splitters)): |
| 78 | if k % 2 == 0: |
| 79 | splitter = re.sub(' +', '\s+', splitters[k]) |
| 80 | regex += splitter |
| 81 | else: |
| 82 | header = splitters[k].strip('<').strip('>') |
| 83 | regex += '(?P<%s>.*?)' % header |
| 84 | headers.append(header) |
| 85 | regex = re.compile('^' + regex + '$') |
| 86 | return headers, regex |
| 87 | |
| 88 | |
| 89 | def formalize_message(enumerated_lines, regex, headers): |
nothing calls this directly
no outgoing calls
no test coverage detected