1818import threading
1919import traceback
2020import urllib .parse
21- from typing import Any , Callable , Dict , List , Optional , Tuple
21+ from typing import Callable , List , Optional , Tuple
22+
23+ import orjson
2224
2325from . import request
2426
2527logger = logging .getLogger (__name__ )
2628
2729
30+ _STANDARD_ATTRS : frozenset [str ] = frozenset (
31+ {
32+ "args" ,
33+ "asctime" ,
34+ "created" ,
35+ "exc_info" ,
36+ "exc_text" ,
37+ "filename" ,
38+ "funcName" ,
39+ "levelname" ,
40+ "levelno" ,
41+ "lineno" ,
42+ "message" ,
43+ "module" ,
44+ "msecs" ,
45+ "msg" ,
46+ "name" ,
47+ "pathname" ,
48+ "process" ,
49+ "processName" ,
50+ "relativeCreated" ,
51+ "stack_info" ,
52+ "taskName" ,
53+ "thread" ,
54+ "threadName" ,
55+ }
56+ )
57+
58+
2859class VictoriaLogsHandler (logging .Handler ):
2960 """VictoriaLogsHandler dispatches log events to a VictoriaLogs server.
3061
@@ -34,7 +65,7 @@ class VictoriaLogsHandler(logging.Handler):
3465 buffer_size: Maximum number of logs the buffer can hold.
3566 If buffer_size <= 0, the size is unlimited (not recommended).
3667 When the buffer is full any new logs will be discarded.
37- 100 000 logs approximately consume 100 MB of RAM in the buffer .
68+ 100 000 logs consume approx. 80- 100 MB of RAM.
3869 chunk_size: Maximum number of logs send per request to the log server.
3970 record_to_stream: function that returns the stream value for a record.
4071 The default will return the name of the top package.
@@ -45,34 +76,6 @@ class VictoriaLogsHandler(logging.Handler):
4576 url: URL of the vlogs server, e.g. `"http://localhost:9428"`
4677 """
4778
48- _STANDARD_ATTRS : frozenset [str ] = frozenset (
49- {
50- "args" ,
51- "asctime" ,
52- "created" ,
53- "exc_info" ,
54- "exc_text" ,
55- "filename" ,
56- "funcName" ,
57- "levelname" ,
58- "levelno" ,
59- "lineno" ,
60- "message" ,
61- "module" ,
62- "msecs" ,
63- "msg" ,
64- "name" ,
65- "pathname" ,
66- "process" ,
67- "processName" ,
68- "relativeCreated" ,
69- "stack_info" ,
70- "taskName" ,
71- "thread" ,
72- "threadName" ,
73- }
74- )
75-
7679 def __init__ (
7780 self ,
7881 batch_size : int = 125 ,
@@ -160,47 +163,31 @@ def close(self):
160163
161164 def emit (self , record : logging .LogRecord ) -> None :
162165 try :
163- log_entry = self ._format_log_record (record )
164- self ._buffer .put_nowait (log_entry )
166+ log = _serialize_log_to_json (record , self ._record_to_stream )
167+
168+ except Exception :
169+ logger .exception ("serializing log: %s" , record )
170+ self .handleError (record )
171+ return
172+
173+ try :
174+ self ._buffer .put_nowait (log )
165175
166176 except queue .Full :
167- logger .error ("Buffer full. Discarding new log." )
177+ logger .error ("Buffer full. Discarding new log: %s" , log )
168178 return
169179
170180 except Exception :
171- logger .exception ("Emitting record" )
181+ logger .exception ("Emitting record: %s" , log )
172182 self .handleError (record )
173183 return
174184
175185 with self ._lock :
176186 self ._added_count += 1
177187 if self ._added_count > self ._batch_size :
188+ self ._added_count = 0
178189 self ._worker_run .set ()
179190
180- def _format_log_record (self , record : logging .LogRecord ) -> Dict [str , Any ]:
181- log = {
182- "stream" : self ._record_to_stream (record ),
183- "timestamp" : record .created ,
184- "level" : record .levelname ,
185- "logger" : record .name ,
186- "module" : record .module ,
187- "function" : record .funcName ,
188- "line_number" : record .lineno ,
189- "message" : record .getMessage (),
190- "process_name" : record .processName ,
191- "process" : record .process ,
192- "thread_name" : record .threadName ,
193- "thread" : record .thread ,
194- }
195- if record .exc_info :
196- log ["exception_name" ], log ["exception" ] = _format_exception (record .exc_info )
197-
198- for k , v in record .__dict__ .items ():
199- if k not in self ._STANDARD_ATTRS :
200- log [k ] = v
201-
202- return log
203-
204191 def start (self ):
205192 """Start the worker. Do nothing when the worker is already running."""
206193 with self ._lock :
@@ -215,19 +202,19 @@ def start(self):
215202 def _worker (self ):
216203 while not self ._worker_shutdown .is_set ():
217204 self ._worker_run .wait (timeout = self ._flush_interval )
218- with self ._lock :
219- self ._added_count = 0
220- self .flush ()
221205 self ._worker_run .clear ()
206+ self .flush ()
222207
223208 logger .debug ("Worker stopped" )
224209
225210 def flush (self ):
226211 """Flush the buffer and send all logs to the log server."""
227- failed : List [Dict [ str , Any ] ] = []
212+ failed : List [bytes ] = []
228213 done = False
214+ success_count = 0
215+
229216 while not done :
230- logs : List [Dict [ str , Any ] ] = []
217+ logs : List [bytes ] = []
231218 while len (logs ) < self ._chunk_size :
232219 try :
233220 logs .append (self ._buffer .get_nowait ())
@@ -237,13 +224,21 @@ def flush(self):
237224
238225 if logs :
239226 ok = request .post_ndjson (
240- url = self ._vlogs_url , data = logs , timeout = self ._request_timeout
227+ url = self ._vlogs_url , objs = logs , timeout = self ._request_timeout
241228 )
242229 if not ok :
243230 failed += logs
244- logger .warning ("Logs failed to send to log server: %d" , len (logs ))
231+ logger .warning (
232+ "Failed to transmit logs to log server: %d" , len (logs )
233+ )
245234 else :
246- logger .debug ("Logs submitted to log server: %d" , len (logs ))
235+ success_count += len (logs )
236+
237+ logger .debug (
238+ "Logs transmitted status: success: %s, failed: %d" ,
239+ success_count ,
240+ len (failed ),
241+ )
247242
248243 if not failed :
249244 return
@@ -263,6 +258,35 @@ def flush(self):
263258 logger .debug ("Saved logs to buffer after failed send: %d" , n )
264259
265260
261+ def _serialize_log_to_json (
262+ record : logging .LogRecord , record_to_stream : Callable [[logging .LogRecord ], str ]
263+ ) -> bytes :
264+ """Serialize a log record into a JSON object and return it."""
265+ obj = {
266+ "stream" : record_to_stream (record ),
267+ "timestamp" : record .created ,
268+ "level" : record .levelname ,
269+ "logger" : record .name ,
270+ "module" : record .module ,
271+ "function" : record .funcName ,
272+ "line_number" : record .lineno ,
273+ "message" : record .getMessage (),
274+ "process_name" : record .processName ,
275+ "process" : record .process ,
276+ "thread_name" : record .threadName ,
277+ "thread" : record .thread ,
278+ }
279+ if record .exc_info :
280+ obj ["exception_name" ], obj ["exception" ] = _format_exception (record .exc_info )
281+
282+ for k , v in record .__dict__ .items ():
283+ if k not in _STANDARD_ATTRS :
284+ obj [k ] = v
285+
286+ log = orjson .dumps (obj , default = str )
287+ return log
288+
289+
266290def _create_filter (name : str ):
267291 def filter_logic (record : logging .LogRecord ) -> bool :
268292 return not record .name .startswith (name )
0 commit comments