SIENTIAPDE-1110
Integrate Druid activity into the Activities class, adding support for Druid configuration and initialization. Update worker and connectors configuration to accommodate Druid, enhancing data processing capabilities.
This commit is contained in:
@@ -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]
|
||||
@@ -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)
|
||||
|
||||
@@ -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:
|
||||
|
||||
74
scouter/activities/pydruid.py
Normal file
74
scouter/activities/pydruid.py
Normal file
@@ -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")
|
||||
@@ -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')),
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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),
|
||||
|
||||
Reference in New Issue
Block a user