SIENTIAPDE-1110

Refactor Activities class to remove Kafka and Druid dependencies, simplifying initialization. Update values.yaml to set replica count to 1 for reduced resource usage. Adjust Redis activity to set TTL to None for better data retention. Remove unused Kafka and Druid activity files and their associated tests, streamlining the codebase.
This commit is contained in:
vitor-aignosi
2025-07-03 16:25:55 -03:00
parent 1dcba6b48e
commit d08d1b1337
13 changed files with 105 additions and 369 deletions

View File

@@ -1,26 +1,22 @@
from temporalio import activity, workflow
from temporalio import workflow
with workflow.unsafe.imports_passed_through():
from sientia_do.temporal.activities.postgres import Postgres
from sientia_do.notifications.handlers import NotificationHandler
from sientia_do.temporal.utils.logger import Logger
from scouter.activities.redis import Redis
from scouter.activities.kafka import Kafka
from scouter.activities.gates import Gates
from scouter.activities.mongodb import MongoDB
from scouter.activities.pydruid import Druid
from typing import Any
class Activities(Postgres, Redis, Kafka, Gates, MongoDB, Druid):
class Activities(Postgres, Redis, Gates, MongoDB,):
"""Activities class that combines multiple services with proper initialization."""
def __init__(self,
postgres_config: dict[str, Any],
redis_config: dict[str, Any],
kafka_config: dict[str, Any],
mongodb_config: dict[str, Any],
druid_config: dict[str, Any],
logger: Logger,
notification_handler: NotificationHandler):
@@ -49,16 +45,6 @@ class Activities(Postgres, Redis, Kafka, Gates, MongoDB, Druid):
password=redis_config['password']
)
# Initialize Kafka
Kafka.__init__(
self,
bootstrap_servers=kafka_config['bootstrap_servers'],
polling_time=kafka_config['polling_time'],
group_id=kafka_config['group_id'],
logger=logger,
notification_handler=notification_handler
)
# Initialize Gates
Gates.__init__(
self,
@@ -75,16 +61,6 @@ class Activities(Postgres, Redis, Kafka, Gates, MongoDB, Druid):
notification_handler=notification_handler
)
# Initialize Druid
Druid.__init__(
self,
host=druid_config['host'],
port=druid_config['port'],
logger=logger,
notification_handler=notification_handler
)
def shutdown(self):
Postgres.close(self)
Kafka.close(self)
MongoDB.close(self)
MongoDB.shutdown(self)

View File

@@ -1,105 +0,0 @@
from temporalio import workflow, activity
with workflow.unsafe.imports_passed_through():
from logging import Logger
from sientia_do.notifications.handlers import NotificationHandler
from sientia_do.temporal.activities.base import BaseActivity
from sientia_do.temporal.utils.logger import Logger
from typing import Any
from aiokafka import AIOKafkaConsumer
from pandas import DataFrame
import json
import asyncio
class Kafka(BaseActivity):
def __init__(self, bootstrap_servers: str, polling_time: int,
group_id: str, logger: Logger, notification_handler: NotificationHandler):
self.polling_time = polling_time
self.bootstrap_servers = bootstrap_servers
self.group_id = group_id
self.consumers = {}
self._consumer_tasks = {}
BaseActivity.__init__(self, logger, notification_handler)
async def close(self):
"""Closes all consumer connections."""
self.info("Closing Kafka connectors...")
for _topic, consumer in self.consumers.items():
await consumer.stop()
async def __aenter__(self):
return self
async def __aexit__(self, exc_type, exc, tb):
await self.close()
async def create_consumer(self, topic: str):
consumer = AIOKafkaConsumer(
bootstrap_servers=self.bootstrap_servers,
auto_offset_reset="earliest",
enable_auto_commit=True,
group_id=f"{self.group_id}-{topic}",
value_deserializer=lambda x: json.loads(x.decode("utf-8"))
)
await consumer.start()
self.consumers[topic] = consumer
@activity.defn(name="load_from_kafka")
async def load_from_kafka(self, input_data: dict[str, Any]) -> dict[str, Any]:
"""
Loads data from a kafka topic. Polls the topic for a given time and returns the data.
Args:
input_data (dict[str, Any]): The data to load. Contains:
topic (str): The topic to load data from.
Returns:
dict[str, Any]: The data loaded from the topic.
"""
metadata = input_data['metadata']
self.debug(
f"Loading data from topic: {input_data['topic']}",
metadata=metadata
)
topic = input_data["topic"]
if topic not in self.consumers:
self.info(
f"Creating consumer for topic: {topic}",
metadata=metadata
)
await self.create_consumer(topic)
consumer = self.consumers[topic]
consumer.subscribe(topics=[topic])
message_values = []
messages = await consumer.getmany(timeout_ms=self.polling_time)
for tp, msgs in messages.items():
msg_topic = tp.topic
if msg_topic == topic:
for msg in msgs:
message_values.append(msg.value)
self.info(
f"Loaded {len(message_values)} messages from topic: {topic}",
metadata=metadata
)
self.debug(
f"Loaded data: {message_values}",
metadata=metadata
)
consumer.unsubscribe()
if not message_values:
return {}
return DataFrame(message_values).to_dict()

View File

@@ -1,71 +0,0 @@
from temporalio import workflow, activity
with workflow.unsafe.imports_passed_through():
import pandas as pd
from typing import List, Optional, Any
from datetime import datetime, timedelta
from pydruid.client import PyDruid
from pydruid.query import QueryBuilder
from sqlalchemy.engine import create_engine
from sqlalchemy import MetaData, Table, select, text
from sientia_do.temporal.activities.base import BaseActivity
from sientia_do.notifications.handlers import NotificationHandler
from sientia_do.temporal.utils.logger import Logger
class Druid(BaseActivity):
def __init__(self, host: str, port: int,
logger: Logger, notification_handler: NotificationHandler):
self.host = host
self.port = port
self.druid_engine = create_engine(
f'druid://{self.host}:{self.port}/druid/v2/sql/')
logger.info(
f"Druid client initialized with host: {self.host}, port: {self.port}")
BaseActivity.__init__(self, logger=logger,
notification_handler=notification_handler)
def shutdown(self):
self.client.close()
def __del__(self):
self.shutdown()
@activity.defn(name="load_latest_druid_data")
async def load_latest_druid_data(self, input_data: dict[str, Any]) -> dict[str, Any]:
"""
Loads the latest data from Druid.
"""
metadata = input_data['metadata']
datasource = f"raw_{input_data['schedule_name']}"
last_data_timestamp = input_data['last_data_timestamp']
last_data_timestamp = last_data_timestamp if last_data_timestamp is not None else '1970-01-01 00:00:00'
self.debug(
f"Loading data from Druid: {input_data}", metadata=metadata)
query = f'"__time" > TIMESTAMP \'{last_data_timestamp}\''
self.info(
f"Loading data from Druid: {datasource} with query: {query}"
)
places = Table(datasource, MetaData(), autoload_with=self.druid_engine)
stmt = select(places).where(text(query))
result = pd.read_sql(stmt, self.druid_engine)
result["inserted_at"] = pd.to_datetime(result["__time"]).dt.strftime(
"%Y-%m-%d %H:%M:%S.%f")
result.drop(columns=["__time"], inplace=True)
self.info(
f"Loaded {len(result)} rows from Druid"
)
self.debug(
f"Druid query result: {result}", metadata=metadata)
return result.to_dict()

View File

@@ -55,7 +55,7 @@ class Redis(RedisBase):
metadata=metadata
)
self.set(key, last_data_timestamp)
self.set(key, last_data_timestamp, ttl=None)
return last_data_timestamp

View File

@@ -14,10 +14,8 @@ with workflow.unsafe.imports_passed_through():
import asyncio
from scouter.utils.connectors_config import (
build_postgres_config,
build_kafka_config,
build_redis_config,
build_mongodb_config,
build_druid_config
build_mongodb_config
)
@@ -42,10 +40,8 @@ async def main():
logger=logger,
notification_handler=notification_handler,
postgres_config=build_postgres_config(),
kafka_config=build_kafka_config(),
redis_config=build_redis_config(),
mongodb_config=build_mongodb_config(),
druid_config=build_druid_config()
mongodb_config=build_mongodb_config()
)
logger.info('Starting Faker Activities...')
@@ -73,10 +69,8 @@ async def main():
workflows=[Scouter, CoreScouter],
activities=[
activities.load_latest_data,
activities.load_latest_druid_data,
activities.get_last_data_timestamp,
activities.put_last_data_timestamp,
activities.load_from_kafka,
activities.data_quality_gate,
activities.aggregate_data,
activities.group_and_hold_data,

View File

@@ -41,16 +41,6 @@ class Scouter:
}
}
# data = await workflow.execute_activity_method(
# Activities.load_from_kafka,
# {
# **metadata,
# 'topic': input_data['topic']
# },
# retry_policy=retry_policy,
# start_to_close_timeout=timedelta(seconds=60)
# )
last_data_timestamp = await workflow.execute_local_activity_method(
Activities.get_last_data_timestamp,
{
@@ -73,16 +63,8 @@ class Scouter:
retry_policy=retry_policy
)
# data = await workflow.execute_local_activity_method(
# Activities.load_latest_druid_data,
# {
# **metadata,
# 'schedule_name': input_data['schedule_name'],
# 'last_data_timestamp': last_data_timestamp
# },
# start_to_close_timeout=timedelta(seconds=60),
# retry_policy=retry_policy
# )
if data == {}:
return
await workflow.execute_activity_method(
Activities.put_last_data_timestamp,
@@ -96,9 +78,6 @@ class Scouter:
retry_policy=retry_policy
)
if data == {}:
return
input_data['data'] = data
input_data['metadata'] = metadata

View File

@@ -2,16 +2,16 @@ from unittest.mock import patch, MagicMock, ANY
from pytest import mark
from sientia_do.temporal.activities.postgres import Postgres
from scouter.activities.activities import Activities
from scouter.activities.mongodb import MongoDB
from scouter.activities.redis import Redis
from scouter.activities.kafka import Kafka
from scouter.activities.gates import Gates
@patch('scouter.activities.activities.MongoDB.__init__')
@patch('scouter.activities.activities.Postgres.__init__')
@patch('scouter.activities.activities.Redis.__init__')
@patch('scouter.activities.activities.Kafka.__init__')
@patch('scouter.activities.activities.Gates.__init__')
def test___init__(mock_gates_init, mock_kafka_init, mock_redis_init, mock_postgres_init):
def test___init__(mock_gates_init, mock_redis_init, mock_postgres_init, mock_mongodb_init):
postgres_config = {
'host': 'localhost',
@@ -30,10 +30,9 @@ def test___init__(mock_gates_init, mock_kafka_init, mock_redis_init, mock_postgr
'password': 'redis'
}
kafka_config = {
'bootstrap_servers': 'localhost:9092',
'polling_time': 1000,
'group_id': 'test-group'
mongodb_config = {
'connection_string': 'mongodb://localhost:27017',
'database_name': 'test_database'
}
logger = MagicMock()
@@ -42,7 +41,7 @@ def test___init__(mock_gates_init, mock_kafka_init, mock_redis_init, mock_postgr
activities = Activities(
postgres_config=postgres_config,
redis_config=redis_config,
kafka_config=kafka_config,
mongodb_config=mongodb_config,
logger=logger,
notification_handler=notification_handler
)
@@ -50,7 +49,7 @@ def test___init__(mock_gates_init, mock_kafka_init, mock_redis_init, mock_postgr
assert isinstance(activities, Activities)
assert isinstance(activities, Postgres)
assert isinstance(activities, Redis)
assert isinstance(activities, Kafka)
assert isinstance(activities, MongoDB)
assert isinstance(activities, Gates)
mock_postgres_init.assert_called_once_with(
@@ -76,11 +75,10 @@ def test___init__(mock_gates_init, mock_kafka_init, mock_redis_init, mock_postgr
notification_handler=notification_handler
)
mock_kafka_init.assert_called_once_with(
mock_mongodb_init.assert_called_once_with(
ANY,
bootstrap_servers=kafka_config['bootstrap_servers'],
polling_time=kafka_config['polling_time'],
group_id=kafka_config['group_id'],
connection_string=mongodb_config['connection_string'],
database_name=mongodb_config['database_name'],
logger=logger,
notification_handler=notification_handler
)
@@ -94,12 +92,12 @@ def test___init__(mock_gates_init, mock_kafka_init, mock_redis_init, mock_postgr
@patch('scouter.activities.activities.Postgres.__init__')
@patch('scouter.activities.activities.Redis.__init__')
@patch('scouter.activities.activities.Kafka.__init__')
@patch('scouter.activities.activities.Gates.__init__')
@patch('scouter.activities.activities.MongoDB.__init__')
@patch('scouter.activities.activities.Postgres.close')
@patch('scouter.activities.activities.Kafka.close')
def test_shutdown(mock_kafka_close, mock_postgres_close,
_mock_gates_init, _mock_redis_init, _mock_kafka_init, _mock_postgres_init):
@patch('scouter.activities.activities.MongoDB.shutdown')
def test_shutdown(mock_mongodb_close, mock_postgres_close, _mock_mongodb_init,
_mock_gates_init, _mock_redis_init, _mock_postgres_init):
postgres_config = {
'host': 'localhost',
'port': 5432,
@@ -117,10 +115,9 @@ def test_shutdown(mock_kafka_close, mock_postgres_close,
'password': 'redis'
}
kafka_config = {
'bootstrap_servers': 'localhost:9092',
'polling_time': 1000,
'group_id': 'test-group'
mongodb_config = {
'connection_string': 'mongodb://localhost:27017',
'database_name': 'test_database'
}
logger = MagicMock()
@@ -129,7 +126,7 @@ def test_shutdown(mock_kafka_close, mock_postgres_close,
activities = Activities(
postgres_config=postgres_config,
redis_config=redis_config,
kafka_config=kafka_config,
mongodb_config=mongodb_config,
logger=logger,
notification_handler=notification_handler
)
@@ -137,4 +134,4 @@ def test_shutdown(mock_kafka_close, mock_postgres_close,
activities.shutdown()
mock_postgres_close.assert_called_once()
mock_kafka_close.assert_called_once()
mock_mongodb_close.assert_called_once()

View File

@@ -290,8 +290,8 @@ async def test_aggregate_data(gates_fixture):
'value': None, 'timestamp': '2023-01-04'},
],
'model_tags': {
'name1': {'aggr_function': 'avg'},
'name2': {'aggr_function': 'max'},
'name1': {'aggr_func': 'avg'},
'name2': {'aggr_func': 'max'},
},
**metadata
}

View File

@@ -1,92 +0,0 @@
from unittest.mock import MagicMock, patch, ANY
from pytest import fixture, mark
from pandas import DataFrame
from scouter.activities.kafka import Kafka
@fixture
@patch("scouter.activities.kafka.KafkaConsumer")
def kafka(_kafka_consumer):
return Kafka(
bootstrap_servers="localhost:9092",
polling_time=1000,
group_id="test-group",
logger=MagicMock(),
notification_handler=MagicMock()
)
@patch("scouter.activities.kafka.KafkaConsumer")
def test___init__(kafka_consumer):
kafka = Kafka(
bootstrap_servers="localhost:9092",
polling_time=1000,
group_id="test-group",
logger=MagicMock(),
notification_handler=MagicMock()
)
assert kafka.polling_time == 1000
assert kafka.kafka_connector == kafka_consumer.return_value
kafka_consumer.assert_called_once_with(
bootstrap_servers="localhost:9092",
auto_offset_reset="earliest",
enable_auto_commit=True,
group_id="test-group",
value_deserializer=ANY
)
metadata = {
'metadata': {
'model_id': 'test_model_id',
'model_name': 'test_model',
'schedule_name': 'test_schedule',
'workflow_name': 'scouter'
}
}
@mark.asyncio
async def test_load_from_kafka(kafka):
input_data = {"topic": "test-topic", **metadata}
data = [
("test-topic", [
MagicMock(
value=f"test-value-{i}"
) for i in range(10)
])
]
kafka.kafka_connector.poll.return_value = MagicMock(
items=MagicMock(return_value=data)
)
expected = DataFrame([d.value for d in data[0][1]]).to_dict()
result = await kafka.load_from_kafka(input_data)
assert result == expected
kafka.kafka_connector.subscribe.assert_called_once_with(["test-topic"])
kafka.kafka_connector.poll.assert_called_once_with(timeout_ms=1000)
@mark.asyncio
async def test_load_from_kafka_empty(kafka):
input_data = {"topic": "test-topic", **metadata}
kafka.kafka_connector.poll.return_value = MagicMock(
items=MagicMock(return_value=[])
)
result = await kafka.load_from_kafka(input_data)
assert result == {}
kafka.kafka_connector.subscribe.assert_called_once_with(["test-topic"])
kafka.kafka_connector.poll.assert_called_once_with(timeout_ms=1000)

View File

@@ -87,7 +87,7 @@ async def test_group_and_hold_data_new_key(redis_activity):
# Verify set was called with correct arguments
redis_activity.set.assert_called_once()
args, kwargs = redis_activity.set.call_args
assert args[0] == 'test_pipeline_test_schedule'
assert args[0] == 'held_data_test_pipeline_test_schedule'
assert args[1] == {
'sensor1': 25.5,
'sensor2': 30.0,
@@ -139,7 +139,7 @@ async def test_group_and_hold_data_update_existing(redis_activity):
# Verify set was called with correct arguments
redis_activity.set.assert_called_once()
args, kwargs = redis_activity.set.call_args
assert args[0] == 'test_workflow_test_schedule'
assert args[0] == 'held_data_test_workflow_test_schedule'
assert args[1] == {
'sensor1': 25.5,
'sensor2': 28.0,

View File

@@ -16,6 +16,14 @@ async def test_core_scouter_workflow_success(mock_workflow, core_scouter):
'filtered_data', 'grouped_data', 'held_data']
await core_scouter.run(
input_data={
'metadata': {
'metadata': {
'model_id': 'test_model_id',
'model_name': 'test_model',
'schedule_name': 'test_schedule',
'workflow_name': 'test_workflow'
}
},
'workflow_name': 'test_workflow',
'schedule_name': 'test_schedule',
'model_name': 'test_model',
@@ -97,6 +105,14 @@ async def test_core_scouter_workflow_with_empty_data(mock_workflow, core_scouter
mock_workflow.execute_local_activity_method.return_value = {}
await core_scouter.run(
input_data={
'metadata': {
'metadata': {
'model_id': 'test_model_id',
'model_name': 'test_model',
'schedule_name': 'test_schedule',
'workflow_name': 'test_workflow'
}
},
'workflow_name': 'test_workflow',
'schedule_name': 'test_schedule',
'model_name': 'test_model',

View File

@@ -1,4 +1,4 @@
from unittest.mock import AsyncMock, patch, ANY
from unittest.mock import AsyncMock, patch, ANY, call
from pytest import fixture, mark
from scouter.workflow.scouter import Scouter
from scouter.activities.activities import Activities
@@ -13,7 +13,10 @@ def scouter():
@patch('scouter.workflow.scouter.workflow', new_callable=AsyncMock)
async def test_scouter_workflow(mock_workflow, scouter):
mock_workflow.execute_activity_method.return_value = 'test_data'
mock_workflow.execute_local_activity_method.side_effect = [
'test_last_data_timestamp',
'test_data'
]
await scouter.run(
input_data={
'topic': 'test_topic',
@@ -32,11 +35,41 @@ async def test_scouter_workflow(mock_workflow, scouter):
}
}
mock_workflow.execute_activity_method.assert_called_once_with(
Activities.load_from_kafka,
mock_workflow.execute_local_activity_method.assert_has_calls(
[
call(
Activities.get_last_data_timestamp,
{
**expected_metadata,
'topic': 'test_topic'
'workflow_name': 'scouter',
'schedule_name': 'test_schedule'
},
retry_policy=ANY,
start_to_close_timeout=ANY
)
])
mock_workflow.execute_local_activity_method.assert_has_calls(
[
call(
Activities.load_latest_data,
{
**expected_metadata,
'collection_name': "raw_test_schedule",
'last_data_timestamp': 'test_last_data_timestamp'
},
retry_policy=ANY,
start_to_close_timeout=ANY
)
])
mock_workflow.execute_activity_method.assert_called_once_with(
Activities.put_last_data_timestamp,
{
**expected_metadata,
'data': 'test_data',
'workflow_name': 'scouter',
'schedule_name': 'test_schedule'
},
retry_policy=ANY,
start_to_close_timeout=ANY
@@ -45,6 +78,7 @@ async def test_scouter_workflow(mock_workflow, scouter):
mock_workflow.execute_child_workflow.assert_called_once_with(
'core_scouter',
{
'metadata': expected_metadata,
'topic': 'test_topic',
'data': 'test_data',
'workflow_name': 'scouter',
@@ -58,7 +92,10 @@ async def test_scouter_workflow(mock_workflow, scouter):
@mark.asyncio
@patch('scouter.workflow.scouter.workflow', new_callable=AsyncMock)
async def test_scouter_workflow_empty(mock_workflow, scouter):
mock_workflow.execute_activity_method.return_value = {}
mock_workflow.execute_local_activity_method.side_effect = [
'test_last_data_timestamp',
{}
]
await scouter.run(
input_data={
'topic': 'test_topic',
@@ -77,14 +114,19 @@ async def test_scouter_workflow_empty(mock_workflow, scouter):
}
}
mock_workflow.execute_activity_method.assert_called_once_with(
Activities.load_from_kafka,
mock_workflow.execute_local_activity_method.assert_has_calls(
[
call(
Activities.load_latest_data,
{
**expected_metadata,
'topic': 'test_topic'
'collection_name': "raw_test_schedule",
'last_data_timestamp': 'test_last_data_timestamp'
},
retry_policy=ANY,
start_to_close_timeout=ANY
)
])
mock_workflow.execute_activity_method.assert_not_called()
mock_workflow.execute_child_workflow.assert_not_called()

View File

@@ -3,7 +3,7 @@
# Declare variables to be passed into your templates.
# This will set the replicaset count more information can be found here: https://kubernetes.io/docs/concepts/workloads/controllers/replicaset/
replicaCount: 3
replicaCount: 1
# This sets the container image more information can be found here: https://kubernetes.io/docs/concepts/containers/images/
image: