MCPcopy Create free account
hub / github.com/WraySmith/log-anomaly / load_to_dataframe

Method load_to_dataframe

parse/logparser/utils/logloader.py:41–69  ·  view source on GitHub ↗

Function to transform log file to dataframe

(self, log_filepath)

Source from the content-addressed store, hash-verified

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

Callers 2

parseMethod · 0.95
matchMethod · 0.95

Calls 1

formalize_messageFunction · 0.85

Tested by

no test coverage detected