Files
sientia-dataops-scouter_tem…/scouter/activities/faker.py
2025-05-19 16:51:22 -03:00

81 lines
2.6 KiB
Python

import random
from datetime import datetime, timezone
from typing import Any
import json
from logging import Logger
from kafka import KafkaProducer
from temporalio import activity
from sientia_do.notifications.handlers import NotificationHandler
from scouter.activities.base import BaseActivity
class Faker(BaseActivity):
def __init__(self, bootstrap_servers: str, logger: Logger,
notification_handler: NotificationHandler):
self.producer = KafkaProducer(
bootstrap_servers=bootstrap_servers,
value_serializer=lambda v: json.dumps(v).encode('utf-8')
)
# Predefined lists for tag and name
self.tags = {
'ns=1;i=1001': 'Temperature Sensor',
'ns=1;i=1002': 'Vibration Meter',
'ns=1;i=1003': 'Pressure Gauge',
'ns=1;i=1004': 'Flow Meter',
'ns=1;i=1005': 'Voltage Sensor',
'ns=1;i=1006': 'Current Sensor'
}
BaseActivity.__init__(self, logger, notification_handler)
@activity.defn(name="generate_and_send_data")
async def generate_and_send_data(self, input_data: dict[str, Any]):
"""
Generates random data and sends it to a Kafka topic.
Args:
input_data (dict[str, Any]): The input data containing:
topic (str): The Kafka topic to send data to
num_messages (int, optional): Number of messages to generate.
Defaults to random.randint(1, len(self.tags)).
"""
topic = input_data.get('topic')
num_messages = input_data.get(
'num_messages', random.randint(1, len(self.tags))) # NOSONAR
if not topic:
raise ValueError("Topic must be specified in input_data")
self.logger.info(
f"Generating {num_messages} messages for topic {topic}")
for _ in range(num_messages):
# Select random tag and name
tag = random.choice(list(self.tags.keys())) # NOSONAR
name = self.tags[tag]
# Generate random value between 0 and 100
if random.random() < 0.1: # NOSONAR
value = None
else:
value = round(random.uniform(0, 100), 2)
# Create data dictionary
data = {
'tag': tag,
'name': name,
'timestamp': datetime.now(timezone.utc).strftime('%Y-%m-%d %H:%M:%S'),
'value': value
}
# Send to Kafka
self.producer.send(topic, value=data)
# Ensure all messages are sent
self.producer.flush()
self.logger.info("Success")