SIENTIAPDE-1083: add prometheus metrics to ingestor.
This commit is contained in:
@@ -7,6 +7,8 @@ from sientia_do.notifications.handlers import NotificationHandler
|
||||
|
||||
from ingestor.managers.ingestor_manager import IngestorManager
|
||||
|
||||
import ingestor.metrics as metrics
|
||||
|
||||
|
||||
class Ingestor:
|
||||
def __init__(self):
|
||||
@@ -33,13 +35,13 @@ class Ingestor:
|
||||
|
||||
kafka_servers = getenv("KAFKA_SERVERS", "localhost:9092")
|
||||
self.redis_host = getenv("REDIS_HOST", "localhost")
|
||||
self.redis_port = int(getenv("REDIS_PORT", '6379'))
|
||||
self.redis_port = int(getenv("REDIS_PORT", "6379"))
|
||||
self.redis_username = getenv("REDIS_USERNAME", None)
|
||||
self.redis_password = getenv("REDIS_PASSWORD", None)
|
||||
self.lease_ttl = int(getenv("LEASE_TTL", '10'))
|
||||
self.heartbeat_ttl = int(getenv("HEARTBEAT_TTL", '20'))
|
||||
self.lease_ttl = int(getenv("LEASE_TTL", "10"))
|
||||
self.heartbeat_ttl = int(getenv("HEARTBEAT_TTL", "20"))
|
||||
self.pod_id = getenv("HOSTNAME", "localhost")
|
||||
self.poll_interval = int(getenv("POLL_INTERVAL", '5'))
|
||||
self.poll_interval = int(getenv("POLL_INTERVAL", "5"))
|
||||
|
||||
self.kafka_servers = kafka_servers.split(",")
|
||||
self.logger = None
|
||||
@@ -51,7 +53,7 @@ class Ingestor:
|
||||
pipeline_name="-",
|
||||
trigger_name="-",
|
||||
model_name="-",
|
||||
model="-"
|
||||
model="-",
|
||||
)
|
||||
|
||||
self.ingestor_manager = None
|
||||
@@ -77,8 +79,7 @@ class Ingestor:
|
||||
logger = getLogger(__name__)
|
||||
logger.setLevel(getenv("LOG_LEVEL", "INFO"))
|
||||
handler = StreamHandler()
|
||||
formatter = Formatter(
|
||||
'%(asctime)s - %(name)s - %(levelname)s - %(message)s')
|
||||
formatter = Formatter("%(asctime)s - %(name)s - %(levelname)s - %(message)s")
|
||||
handler.setFormatter(formatter)
|
||||
|
||||
logger.addHandler(handler)
|
||||
@@ -128,9 +129,17 @@ class Ingestor:
|
||||
"""
|
||||
|
||||
self.ingestor_manager = IngestorManager(
|
||||
self.kafka_servers, self.redis_host, self.redis_port, self.lease_ttl,
|
||||
self.heartbeat_ttl, self.pod_id, self.poll_interval, self.logger,
|
||||
self.notification_handler, self.redis_username, self.redis_password
|
||||
self.kafka_servers,
|
||||
self.redis_host,
|
||||
self.redis_port,
|
||||
self.lease_ttl,
|
||||
self.heartbeat_ttl,
|
||||
self.pod_id,
|
||||
self.poll_interval,
|
||||
self.logger,
|
||||
self.notification_handler,
|
||||
self.redis_username,
|
||||
self.redis_password,
|
||||
)
|
||||
|
||||
# Declare ingestor ative
|
||||
@@ -142,6 +151,10 @@ class Ingestor:
|
||||
|
||||
self.handle_acquired_tags(acquired)
|
||||
|
||||
metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set(
|
||||
len(self.ingestor_manager.managed_tags)
|
||||
) # Set initial
|
||||
|
||||
def manage_no_slots(self, number_of_slots: int):
|
||||
"""
|
||||
Manages the scenario where there are no slots assigned to the ingestor.
|
||||
@@ -159,7 +172,9 @@ class Ingestor:
|
||||
# Get slot lease
|
||||
self.ingestor_manager.get_slot_leases(1)
|
||||
|
||||
def manage_leases(self, available_slots: int, lacking_ingestors: int, slot_diff: int):
|
||||
def manage_leases(
|
||||
self, available_slots: int, lacking_ingestors: int, slot_diff: int
|
||||
):
|
||||
"""
|
||||
Manages the allocation and deallocation of slot leases for ingestors based on
|
||||
the number of available slots, lacking ingestors, and slot differences.
|
||||
@@ -198,6 +213,11 @@ class Ingestor:
|
||||
self.ingestor_manager.unsubscribe_slot(lease)
|
||||
self.ingestor_manager.managed_tags.pop(lease)
|
||||
|
||||
# Update metric after removal
|
||||
metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set(
|
||||
len(self.ingestor_manager.managed_tags)
|
||||
)
|
||||
|
||||
def update_ingestor_manager(self, old_managed_tags: Dict[str, Any]):
|
||||
"""
|
||||
Updates the ingestor manager with the new managed tags.
|
||||
@@ -205,15 +225,18 @@ class Ingestor:
|
||||
old_managed_tags (Dict[str, Any]): The old managed tags.
|
||||
"""
|
||||
|
||||
self.logger.debug("Current managed tags: %s",
|
||||
self.ingestor_manager.managed_tags)
|
||||
self.logger.debug(
|
||||
"Current managed tags: %s", self.ingestor_manager.managed_tags
|
||||
)
|
||||
|
||||
self.ingestor_manager.update_opc_servers()
|
||||
new_managed_tags = deepcopy(self.ingestor_manager.managed_tags)
|
||||
|
||||
self.logger.debug("Comparing new managed tags %s"
|
||||
"with old managed tags %s",
|
||||
new_managed_tags, old_managed_tags)
|
||||
self.logger.debug(
|
||||
"Comparing new managed tags %s with old managed tags %s",
|
||||
new_managed_tags,
|
||||
old_managed_tags,
|
||||
)
|
||||
|
||||
for slot, config in new_managed_tags.items():
|
||||
if slot not in old_managed_tags:
|
||||
@@ -228,6 +251,11 @@ class Ingestor:
|
||||
if slot not in new_managed_tags:
|
||||
self.ingestor_manager.unsubscribe_slot(slot)
|
||||
|
||||
# Ensure the gauge is updated after any potential changes here
|
||||
metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set(
|
||||
len(self.ingestor_manager.managed_tags)
|
||||
)
|
||||
|
||||
def loop(self):
|
||||
"""
|
||||
Executes the main loop for managing ingestors and slots.
|
||||
@@ -259,6 +287,9 @@ class Ingestor:
|
||||
number_of_leases = self.ingestor_manager.get_number_of_leases()
|
||||
number_of_slots = self.ingestor_manager.get_number_of_slots()
|
||||
|
||||
# Update active ingestors gauge
|
||||
metrics.ACTIVE_INGESTORS.set(number_of_ingestors)
|
||||
|
||||
# Handle no slots
|
||||
self.logger.debug("Managing no slots...")
|
||||
self.manage_no_slots(number_of_slots)
|
||||
@@ -270,6 +301,11 @@ class Ingestor:
|
||||
self.logger.debug("Managing leases...")
|
||||
self.manage_leases(available_slots, lacking_ingestors, slot_diff)
|
||||
|
||||
# Update managed slots gauge
|
||||
metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set(
|
||||
len(self.ingestor_manager.managed_tags)
|
||||
)
|
||||
|
||||
self.logger.debug(
|
||||
"Active ingestors: %s, "
|
||||
"Number of slots: %s, "
|
||||
@@ -280,7 +316,7 @@ class Ingestor:
|
||||
number_of_slots,
|
||||
number_of_leases,
|
||||
self.ingestor_manager.managed_tags,
|
||||
self.ingestor_manager.opc_managers
|
||||
self.ingestor_manager.opc_managers,
|
||||
)
|
||||
if not self.ingestor_manager.managed_tags:
|
||||
# No slots acquired
|
||||
@@ -295,6 +331,4 @@ class Ingestor:
|
||||
self.ingestor_manager.check_opc_servers_integrity()
|
||||
|
||||
self.logger.debug("Updating managed tags...")
|
||||
self.update_ingestor_manager(
|
||||
current_managed_tags
|
||||
)
|
||||
self.update_ingestor_manager(current_managed_tags)
|
||||
|
||||
@@ -34,17 +34,17 @@ APP_UP = Gauge(
|
||||
)
|
||||
|
||||
# --- Ingestor Manager Metrics ---
|
||||
INGESTORS_ATIVOS = Gauge(
|
||||
ACTIVE_INGESTORS = Gauge(
|
||||
"ingestor_active_total",
|
||||
"Number of active ingestors reported by Redis",
|
||||
# No pod_id here, as it's a global view from Redis
|
||||
)
|
||||
SLOTS_TOTAIS = Gauge(
|
||||
SLOTS_TOTAL = Gauge(
|
||||
"ingestor_slots_total",
|
||||
"Total number of slots configured in Redis",
|
||||
# No pod_id here, as it's a global view from Redis
|
||||
)
|
||||
LEASES_TOTAIS = Gauge(
|
||||
LEASES_TOTAL = Gauge(
|
||||
"ingestor_leases_total",
|
||||
"Total number of leases (allocated slots) in Redis",
|
||||
# No pod_id here, as it's a global view from Redis
|
||||
|
||||
Reference in New Issue
Block a user