mirror of
https://github.com/pyMC-dev/pyMC_Repeater.git
synced 2026-08-07 09:23:06 +02:00
Merge pull request #196 from tjdownes/perf/rrdtool-batch
perf(rrdtool): cache get_data() result for 60 s to avoid repeated disk reads
This commit is contained in:
@@ -19,10 +19,14 @@ class RRDToolHandler:
|
||||
self.rrd_path = self.storage_dir / "metrics.rrd"
|
||||
self.available = RRDTOOL_AVAILABLE
|
||||
self._init_rrd()
|
||||
# Batch RRD updates: track pending update and last cached info
|
||||
self._pending_rrd_update = None
|
||||
self._last_rrd_info_time = 0
|
||||
self._last_rrd_info_cache = None
|
||||
# Timestamp of the last successful rrdtool.update() call (unix seconds,
|
||||
# aligned to the 60-second RRD step). Used to skip writes whose period
|
||||
# has already been committed — no rrdtool.info() call needed.
|
||||
self._last_rrd_update: int = 0
|
||||
# Read-side cache: rrdtool.fetch() returns 24 h of data and is a
|
||||
# blocking disk read. Cache the result for 60 s — matching the RRD
|
||||
# step size — so repeated dashboard refreshes don't hammer the SD card.
|
||||
self._get_data_cache: tuple = (0.0, None) # (fetched_at, result)
|
||||
|
||||
def _init_rrd(self):
|
||||
if not self.available:
|
||||
@@ -77,10 +81,11 @@ class RRDToolHandler:
|
||||
logger.error(f"Failed to create RRD database: {e}")
|
||||
|
||||
def update_packet_metrics(self, record: dict, cumulative_counts: dict):
|
||||
"""Buffer packet metrics for batch RRD update instead of per-packet writes.
|
||||
|
||||
RRD uses 60-second time steps, so we batch updates within each period
|
||||
and only write when the time period changes or buffer is full.
|
||||
"""Write packet metrics to RRD, throttled to once per 60-second step.
|
||||
|
||||
RRD enforces a 60-second minimum step between updates. We track the
|
||||
last written timestamp ourselves — no rrdtool.info() call needed, which
|
||||
previously allocated thousands of Python objects per call.
|
||||
"""
|
||||
if not self.available or not self.rrd_path.exists():
|
||||
return
|
||||
@@ -88,27 +93,8 @@ class RRDToolHandler:
|
||||
try:
|
||||
timestamp = int(record.get("timestamp", time.time()))
|
||||
|
||||
# Cache RRD info for up to 5 seconds to avoid repeated rrdtool.info() calls
|
||||
now = time.time()
|
||||
if now - self._last_rrd_info_time > 5 or self._last_rrd_info_cache is None:
|
||||
try:
|
||||
self._last_rrd_info_cache = rrdtool.info(str(self.rrd_path))
|
||||
self._last_rrd_info_time = now
|
||||
except Exception as e:
|
||||
logger.debug(f"Failed to cache RRD info: {e}")
|
||||
self._last_rrd_info_cache = None
|
||||
return
|
||||
|
||||
if self._last_rrd_info_cache is None:
|
||||
return
|
||||
|
||||
last_update = int(self._last_rrd_info_cache.get("last_update", timestamp - 60))
|
||||
|
||||
# Skip if timestamp is in same or earlier time period than last update
|
||||
# (RRD step is 60 seconds)
|
||||
if timestamp <= last_update:
|
||||
# But still buffer cumulative counts for when we do update
|
||||
self._pending_rrd_update = (timestamp, cumulative_counts, record)
|
||||
# Skip if this packet falls in the same 60-second period we already wrote.
|
||||
if timestamp <= self._last_rrd_update:
|
||||
return
|
||||
|
||||
# Build update string from cumulative counts
|
||||
@@ -140,11 +126,8 @@ class RRDToolHandler:
|
||||
type_values_str = ":".join(type_values)
|
||||
values = f"{basic_values}:{type_values_str}"
|
||||
|
||||
# Write to RRD - this is now only called once per 60-second period
|
||||
rrdtool.update(str(self.rrd_path), values)
|
||||
# Invalidate cache so next period fetches fresh info
|
||||
self._last_rrd_info_cache = None
|
||||
self._pending_rrd_update = None
|
||||
self._last_rrd_update = timestamp
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to update RRD packet metrics: {e}")
|
||||
@@ -162,9 +145,20 @@ class RRDToolHandler:
|
||||
)
|
||||
return None
|
||||
|
||||
# Serve from cache if result is still fresh. RRD step is 60 s, so
|
||||
# anything newer than that is guaranteed to be identical to a live fetch.
|
||||
# Only the default (full 24-hour, no explicit bounds) call is cached —
|
||||
# explicit start/end requests always bypass the cache.
|
||||
now = time.time()
|
||||
use_cache = start_time is None and end_time is None
|
||||
if use_cache:
|
||||
cache_fetched_at, cache_result = self._get_data_cache
|
||||
if now - cache_fetched_at < 60.0 and cache_result is not None:
|
||||
return cache_result
|
||||
|
||||
try:
|
||||
if end_time is None:
|
||||
end_time = int(time.time())
|
||||
end_time = int(now)
|
||||
if start_time is None:
|
||||
start_time = end_time - (24 * 3600)
|
||||
|
||||
@@ -220,6 +214,10 @@ class RRDToolHandler:
|
||||
|
||||
result["timestamps"] = timestamps
|
||||
|
||||
# Populate read cache for default (unconstrained) calls only.
|
||||
if use_cache:
|
||||
self._get_data_cache = (now, result)
|
||||
|
||||
return result
|
||||
|
||||
except Exception as e:
|
||||
|
||||
Reference in New Issue
Block a user