diff --git a/scouter/activities/activities.py b/scouter/activities/activities.py index 53f7f0f..59990b4 100644 --- a/scouter/activities/activities.py +++ b/scouter/activities/activities.py @@ -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) diff --git a/scouter/activities/kafka.py b/scouter/activities/kafka.py deleted file mode 100644 index ea36655..0000000 --- a/scouter/activities/kafka.py +++ /dev/null @@ -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() diff --git a/scouter/activities/pydruid.py b/scouter/activities/pydruid.py deleted file mode 100644 index f19e4bb..0000000 --- a/scouter/activities/pydruid.py +++ /dev/null @@ -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() diff --git a/scouter/activities/redis.py b/scouter/activities/redis.py index b4ad730..d9f54e4 100644 --- a/scouter/activities/redis.py +++ b/scouter/activities/redis.py @@ -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 diff --git a/scouter/worker/worker.py b/scouter/worker/worker.py index b99c05c..d0e4fdd 100644 --- a/scouter/worker/worker.py +++ b/scouter/worker/worker.py @@ -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, diff --git a/scouter/workflow/scouter.py b/scouter/workflow/scouter.py index eedebc8..ffcb625 100644 --- a/scouter/workflow/scouter.py +++ b/scouter/workflow/scouter.py @@ -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 diff --git a/tests/activities/test_activities.py b/tests/activities/test_activities.py index 2150dca..2740d18 100644 --- a/tests/activities/test_activities.py +++ b/tests/activities/test_activities.py @@ -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() diff --git a/tests/activities/test_gates.py b/tests/activities/test_gates.py index 245b176..ba87f5a 100644 --- a/tests/activities/test_gates.py +++ b/tests/activities/test_gates.py @@ -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 } diff --git a/tests/activities/test_kafka.py b/tests/activities/test_kafka.py deleted file mode 100644 index 681c9dd..0000000 --- a/tests/activities/test_kafka.py +++ /dev/null @@ -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) diff --git a/tests/activities/test_redis.py b/tests/activities/test_redis.py index a266c71..f4d8977 100644 --- a/tests/activities/test_redis.py +++ b/tests/activities/test_redis.py @@ -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, diff --git a/tests/workflow/sub_workflows/test_core_scouter.py b/tests/workflow/sub_workflows/test_core_scouter.py index 905b864..c92a175 100644 --- a/tests/workflow/sub_workflows/test_core_scouter.py +++ b/tests/workflow/sub_workflows/test_core_scouter.py @@ -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', diff --git a/tests/workflow/test_scouter.py b/tests/workflow/test_scouter.py index 9727293..8d1b3a3 100644 --- a/tests/workflow/test_scouter.py +++ b/tests/workflow/test_scouter.py @@ -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_local_activity_method.assert_has_calls( + [ + call( + Activities.get_last_data_timestamp, + { + **expected_metadata, + '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.load_from_kafka, + Activities.put_last_data_timestamp, { **expected_metadata, - 'topic': 'test_topic' + '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, - { - **expected_metadata, - 'topic': 'test_topic' - }, - 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_not_called() mock_workflow.execute_child_workflow.assert_not_called() diff --git a/values.yaml b/values.yaml index 3bf6fa4..ded30e4 100644 --- a/values.yaml +++ b/values.yaml @@ -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: