Skip to content

Commit aa876b1

Browse files
committed
perf: store logs as bytes in buffer
1 parent 2ac9fcd commit aa876b1

6 files changed

Lines changed: 109 additions & 115 deletions

File tree

src/vlogs_handler/__init__.py

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,5 @@
66
__version__ = "0.1.0dev7"
77

88
# TODO
9-
# - Consider storing logs as strings or bytes to reduce memory footprint
10-
# - Consider using a persistent queue
119
# - Enable logging and switch to full debug logging in prod for testing
10+
# - Consider using a persistent queue

src/vlogs_handler/handler.py

Lines changed: 80 additions & 66 deletions
Original file line numberDiff line numberDiff line change
@@ -18,13 +18,44 @@
1818
import threading
1919
import traceback
2020
import 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

2325
from . import request
2426

2527
logger = 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+
2859
class 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,17 @@ 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
229214
while not done:
230-
logs: List[Dict[str, Any]] = []
215+
logs: List[bytes] = []
231216
while len(logs) < self._chunk_size:
232217
try:
233218
logs.append(self._buffer.get_nowait())
@@ -237,13 +222,13 @@ def flush(self):
237222

238223
if logs:
239224
ok = request.post_ndjson(
240-
url=self._vlogs_url, data=logs, timeout=self._request_timeout
225+
url=self._vlogs_url, objs=logs, timeout=self._request_timeout
241226
)
242227
if not ok:
243228
failed += logs
244-
logger.warning("Logs failed to send to log server: %d", len(logs))
229+
logger.warning("Failed to send logs to log server: %d", len(logs))
245230
else:
246-
logger.debug("Logs submitted to log server: %d", len(logs))
231+
logger.debug("Logs send to log server: %d", len(logs))
247232

248233
if not failed:
249234
return
@@ -263,6 +248,35 @@ def flush(self):
263248
logger.debug("Saved logs to buffer after failed send: %d", n)
264249

265250

251+
def _serialize_log_to_json(
252+
record: logging.LogRecord, record_to_stream: Callable[[logging.LogRecord], str]
253+
) -> bytes:
254+
"""Serialize a log record into a JSON object and return it."""
255+
obj = {
256+
"stream": record_to_stream(record),
257+
"timestamp": record.created,
258+
"level": record.levelname,
259+
"logger": record.name,
260+
"module": record.module,
261+
"function": record.funcName,
262+
"line_number": record.lineno,
263+
"message": record.getMessage(),
264+
"process_name": record.processName,
265+
"process": record.process,
266+
"thread_name": record.threadName,
267+
"thread": record.thread,
268+
}
269+
if record.exc_info:
270+
obj["exception_name"], obj["exception"] = _format_exception(record.exc_info)
271+
272+
for k, v in record.__dict__.items():
273+
if k not in _STANDARD_ATTRS:
274+
obj[k] = v
275+
276+
log = orjson.dumps(obj, default=str)
277+
return log
278+
279+
266280
def _create_filter(name: str):
267281
def filter_logic(record: logging.LogRecord) -> bool:
268282
return not record.name.startswith(name)

src/vlogs_handler/request.py

Lines changed: 7 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -4,9 +4,7 @@
44
import urllib.error
55
import urllib.parse
66
import urllib.request
7-
from typing import Any, List, Optional
8-
9-
import orjson # significantly better performance than standard json library
7+
from typing import List, Optional
108

119
logger = logging.getLogger(__name__)
1210

@@ -20,30 +18,21 @@ def is_url(url: str) -> bool:
2018
return False
2119

2220

23-
def post_ndjson(*, url: str, data: List[Any], timeout: Optional[float] = None) -> bool:
21+
def post_ndjson(
22+
*, url: str, objs: List[bytes], timeout: Optional[float] = None
23+
) -> bool:
2424
"""Send a POST request with the ndjson protocol
2525
and report whether it was successful.
2626
2727
Args:
2828
url: request URL
29-
data: list of objects to send
29+
data: list of JSON objects to send
3030
timeout: request timeout in seconds.
3131
Settings it to None will disable the timeout.
3232
"""
3333

34-
chunks: List[bytes] = []
35-
for obj in data:
36-
try:
37-
chunks.append(orjson.dumps(obj, option=orjson.OPT_APPEND_NEWLINE))
38-
except (TypeError, ValueError):
39-
logger.exception(
40-
"Could not convert obj to JSON. Discarded", extra={"log": obj}
41-
)
42-
continue
43-
44-
data_bytes = b"".join(chunks)
45-
46-
req = urllib.request.Request(url, data=data_bytes, method="POST")
34+
data = b"\n".join(objs)
35+
req = urllib.request.Request(url, data=data, method="POST")
4736
req.add_header("Content-Type", "application/x-ndjson")
4837

4938
try:

tests/test_encoder.py

Lines changed: 1 addition & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -31,18 +31,11 @@ def test_should_encode(self):
3131
got = orjson.loads(orjson.dumps(data, default=str))
3232

3333
# then
34-
self.assertEqual(got["class"], "<class 'tests.test_encoder.MyClass'>")
34+
self.assertIn("MyClass", got["class"])
3535
self.assertEqual(got["date"], "2026-01-11")
3636
self.assertEqual(got["datetime"], "2026-01-11T12:15:42.000099+00:00")
3737
self.assertEqual(got["float"], 1.23)
3838
self.assertIn("my_func at", got["func"])
3939
self.assertEqual(got["integer"], 1)
4040
self.assertEqual(got["set"], "{1, 2, 3}")
4141
self.assertEqual(got["text"], "Alpha")
42-
self.assertEqual(got["text"], "Alpha")
43-
self.assertEqual(got["text"], "Alpha")
44-
self.assertEqual(got["text"], "Alpha")
45-
self.assertEqual(got["text"], "Alpha")
46-
self.assertEqual(got["text"], "Alpha")
47-
self.assertEqual(got["text"], "Alpha")
48-
self.assertEqual(got["text"], "Alpha")

tests/test_handler.py

Lines changed: 14 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,8 @@
55
import unittest
66
from unittest.mock import MagicMock, patch
77

8+
import orjson
9+
810
from vlogs_handler import VictoriaLogsHandler, handler
911

1012
MODULE_PATH = "vlogs_handler.handler"
@@ -93,7 +95,7 @@ def test_should_send_a_log(self, m: MagicMock):
9395
self.assertTrue(called_in_time)
9496

9597
self.assertEqual(m.call_count, 1)
96-
got = m.call_args.kwargs["data"][0]
98+
got = orjson.loads(m.call_args.kwargs["objs"][0])
9799
self.assertEqual(got["stream"], "test_logger")
98100
self.assertEqual(got["level"], "INFO")
99101
self.assertEqual(got["logger"], "test_logger.module_1.module_2")
@@ -113,10 +115,10 @@ def test_should_send_log_with_extras(self, m: MagicMock):
113115
self.assertTrue(called_in_time)
114116

115117
self.assertEqual(m.call_count, 1)
116-
got = m.call_args.kwargs["data"][0]
118+
got = orjson.loads(m.call_args.kwargs["objs"][0])
117119
self.assertEqual(got["message"], "Alpha")
118120
self.assertEqual(got["planet"], "Jupiter")
119-
self.assertEqual(got["deadline"], my_date)
121+
self.assertEqual(got["deadline"], "2026-01-11T12:15:42.000099+00:00")
120122

121123
def test_should_send_exception_log(self, m: MagicMock):
122124
# given
@@ -133,7 +135,7 @@ def test_should_send_exception_log(self, m: MagicMock):
133135
self.assertTrue(called_in_time)
134136

135137
self.assertEqual(m.call_count, 1)
136-
got = m.call_args.kwargs["data"][0]
138+
got = orjson.loads(m.call_args.kwargs["objs"][0])
137139
self.assertEqual(got["stream"], "test_logger")
138140
self.assertEqual(got["level"], "ERROR")
139141
self.assertEqual(got["logger"], "test_logger.module_1.module_2")
@@ -185,11 +187,11 @@ def test_handler_should_send_multiple_logs_in_single_request(self, m: MagicMock)
185187
self.assertTrue(called_in_time)
186188

187189
self.assertEqual(m.call_count, 1)
188-
got = m.call_args.kwargs["data"]
190+
got = m.call_args.kwargs["objs"]
189191
self.assertEqual(len(got), 2)
190192

191-
self.assertEqual(got[0]["message"], "Alpha")
192-
self.assertEqual(got[1]["message"], "Bravo")
193+
self.assertEqual(orjson.loads(got[0])["message"], "Alpha")
194+
self.assertEqual(orjson.loads(got[1])["message"], "Bravo")
193195

194196

195197
@patch(MODULE_PATH + ".request.post_ndjson")
@@ -237,7 +239,7 @@ def test_handler_should_send_immediately_when_batch_size_reached(
237239

238240
self.assertEqual(m.call_count, 1)
239241

240-
self.assertEqual(len(m.call_args.kwargs["data"]), 4)
242+
self.assertEqual(len(m.call_args.kwargs["objs"]), 4)
241243

242244
def test_handler_should_shutdown_gracefully(self, m: MagicMock):
243245
# given
@@ -251,9 +253,9 @@ def test_handler_should_shutdown_gracefully(self, m: MagicMock):
251253

252254
# then
253255
self.assertEqual(m.call_count, 1)
254-
got = m.call_args.kwargs["data"]
256+
got = m.call_args.kwargs["objs"]
255257
self.assertEqual(len(got), 1)
256-
self.assertEqual(got[0]["message"], "Alpha")
258+
self.assertEqual(orjson.loads(got[0])["message"], "Alpha")
257259

258260

259261
@patch(MODULE_PATH + ".request.post_ndjson")
@@ -299,8 +301,8 @@ def test_should_send_logs_in_chunks(self, m: MagicMock):
299301
# then
300302
self.assertEqual(handler._buffer.qsize(), 0)
301303
self.assertEqual(m.call_count, 2)
302-
self.assertEqual(len(m.call_args_list[0].kwargs["data"]), 3)
303-
self.assertEqual(len(m.call_args_list[1].kwargs["data"]), 1)
304+
self.assertEqual(len(m.call_args_list[0].kwargs["objs"]), 3)
305+
self.assertEqual(len(m.call_args_list[1].kwargs["objs"]), 1)
304306

305307

306308
def add_to_queue(queue: queue.Queue, items):
@@ -338,5 +340,3 @@ def test_empty_string_name(self):
338340
self.mock_record.name = ""
339341
result = handler._top_package_name(self.mock_record)
340342
self.assertEqual(result, "")
341-
self.assertEqual(result, "")
342-
self.assertEqual(result, "")

0 commit comments

Comments
 (0)