diff --git a/requirements.txt b/requirements.txt index 3e098f3..c249857 100644 --- a/requirements.txt +++ b/requirements.txt @@ -7,3 +7,4 @@ aiokafka pymongo git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.2.0 git+ssh://git@github.com/Aignosi/sientia-mlops-library.git@0.38.1 +pydruid[pandas] \ No newline at end of file diff --git a/scouter/activities/activities.py b/scouter/activities/activities.py index ddeae2e..53f7f0f 100644 --- a/scouter/activities/activities.py +++ b/scouter/activities/activities.py @@ -8,10 +8,11 @@ with workflow.unsafe.imports_passed_through(): 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): +class Activities(Postgres, Redis, Kafka, Gates, MongoDB, Druid): """Activities class that combines multiple services with proper initialization.""" def __init__(self, @@ -19,6 +20,7 @@ class Activities(Postgres, Redis, Kafka, Gates, MongoDB): 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): @@ -73,6 +75,15 @@ class Activities(Postgres, Redis, Kafka, Gates, MongoDB): 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) diff --git a/scouter/activities/mongodb.py b/scouter/activities/mongodb.py index a02db81..2191d55 100644 --- a/scouter/activities/mongodb.py +++ b/scouter/activities/mongodb.py @@ -1,15 +1,15 @@ -from pandas import DataFrame from temporalio import workflow, activity with workflow.unsafe.imports_passed_through(): from typing import Any import traceback - from logging import Logger from datetime import datetime from pymongo import MongoClient + from pandas import DataFrame from sientia_do.notifications.handlers import NotificationHandler from sientia_do.notifications.models import NotificationLevel from sientia_do.temporal.activities.base import BaseActivity + from sientia_do.temporal.utils.logger import Logger def clear_mongo_id(docs: list) -> list: diff --git a/scouter/activities/pydruid.py b/scouter/activities/pydruid.py new file mode 100644 index 0000000..2dc2478 --- /dev/null +++ b/scouter/activities/pydruid.py @@ -0,0 +1,74 @@ +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 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, + endpoint: str = "druid/v2"): + + self.host = host + self.port = port + self.endpoint = endpoint + self.client = PyDruid( + f"http://{self.host}:{self.port}", {self.endpoint} + ) + 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 = datetime.strptime( + input_data['last_data_timestamp'], "%Y-%m-%d %H:%M:%S.%f") + + self.debug( + f"Loading data from Druid: {input_data}", metadata=metadata) + + end_time = datetime(9999, 12, 31, 23, 59, 59) + + interval = f"{last_data_timestamp.isoformat()}Z/{end_time.isoformat()}Z" + + builder = QueryBuilder() + + query = builder.scan( + { + "datasource": datasource, + "intervals": interval, + "columns": ["timestamp", "value", "tag"], + "limit": 10000, + } + ) + + result = query.export_pandas() + + self.info( + f"Loaded {len(result)} rows from Druid" + ) + + self.debug( + f"Druid query result: {result}", metadata=metadata) + + return result.to_dict(orient="records") diff --git a/scouter/utils/connectors_config.py b/scouter/utils/connectors_config.py index 64a7feb..2179fac 100644 --- a/scouter/utils/connectors_config.py +++ b/scouter/utils/connectors_config.py @@ -40,3 +40,10 @@ def build_mongodb_config(): 'connection_string': connection_string, 'database_name': getenv('MONGODB_DATABASE_NAME', 'sientia') } + + +def build_druid_config(): + return { + 'host': getenv('DRUID_HOST', 'localhost'), + 'port': int(getenv('DRUID_PORT', '8082')), + } diff --git a/scouter/worker/worker.py b/scouter/worker/worker.py index 74ad479..b99c05c 100644 --- a/scouter/worker/worker.py +++ b/scouter/worker/worker.py @@ -16,7 +16,8 @@ with workflow.unsafe.imports_passed_through(): build_postgres_config, build_kafka_config, build_redis_config, - build_mongodb_config + build_mongodb_config, + build_druid_config ) @@ -43,7 +44,8 @@ async def main(): postgres_config=build_postgres_config(), kafka_config=build_kafka_config(), redis_config=build_redis_config(), - mongodb_config=build_mongodb_config() + mongodb_config=build_mongodb_config(), + druid_config=build_druid_config() ) logger.info('Starting Faker Activities...') @@ -71,6 +73,7 @@ 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, diff --git a/scouter/workflow/scouter.py b/scouter/workflow/scouter.py index cf36c49..e9ca1f2 100644 --- a/scouter/workflow/scouter.py +++ b/scouter/workflow/scouter.py @@ -62,11 +62,22 @@ class Scouter: retry_policy=retry_policy ) + # data = await workflow.execute_local_activity_method( + # Activities.load_latest_data, + # { + # **metadata, + # 'collection_name': f"raw_{input_data['schedule_name']}", + # 'last_data_timestamp': last_data_timestamp + # }, + # start_to_close_timeout=timedelta(seconds=60), + # retry_policy=retry_policy + # ) + data = await workflow.execute_local_activity_method( - Activities.load_latest_data, + Activities.load_latest_druid_data, { **metadata, - 'collection_name': f"raw_{input_data['schedule_name']}", + 'schedule_name': input_data['schedule_name'], 'last_data_timestamp': last_data_timestamp }, start_to_close_timeout=timedelta(seconds=60),