diff --git a/requirements.txt b/requirements.txt index a6a4b25..8c7edc9 100644 --- a/requirements.txt +++ b/requirements.txt @@ -5,7 +5,7 @@ asyncua redis aiokafka pymongo -git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.4.2 +git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.4.3 git+ssh://git@github.com/Aignosi/sientia-mlops-library.git@0.38.5 pydruid[pandas] prometheus-client \ No newline at end of file diff --git a/scouter/activities/mongodb.py b/scouter/activities/mongodb.py index e36c563..878b185 100644 --- a/scouter/activities/mongodb.py +++ b/scouter/activities/mongodb.py @@ -10,6 +10,7 @@ with workflow.unsafe.imports_passed_through(): from sientia_do.notifications.models import NotificationLevel from sientia_do.temporal.activities.base import BaseActivity from sientia_do.observability.logger import Logger + from sientia_do.temporal.constants import DATETIME_FORMAT_MS_WITH_TZ def clear_mongo_id(docs: list) -> list: @@ -99,7 +100,7 @@ class MongoDB(BaseActivity): else: data_filter = { "inserted_at": { - "$gt": datetime.strptime(last_data_timestamp, "%Y-%m-%d %H:%M:%S.%f") + "$gt": datetime.strptime(last_data_timestamp, DATETIME_FORMAT_MS_WITH_TZ) } } @@ -120,7 +121,7 @@ class MongoDB(BaseActivity): for item in data: item['inserted_at'] = item['inserted_at'].strftime( - "%Y-%m-%d %H:%M:%S.%f") + DATETIME_FORMAT_MS_WITH_TZ) self.info( f"Loaded {len(data)} documents from MongoDB", diff --git a/scouter/activities/redis.py b/scouter/activities/redis.py index fb3f556..284acd5 100644 --- a/scouter/activities/redis.py +++ b/scouter/activities/redis.py @@ -11,6 +11,7 @@ with workflow.unsafe.imports_passed_through(): from pandas import DataFrame from datetime import datetime from scouter import metrics + from sientia_do.temporal.constants import DATETIME_FORMAT class Redis(RedisBase): @@ -169,7 +170,7 @@ class Redis(RedisBase): (row['name'], value)) data_hold['timestamp'] = data['timestamp'].max() if not data.empty else \ - datetime.now().strftime("%Y-%m-%d %H:%M:%S") + data_hold['timestamp'] self.set(key, data_hold, ttl=retention_time) @@ -224,7 +225,7 @@ class Redis(RedisBase): data: The data used to collect the data. """ metadata = input_data['metadata'] - key = f"data_package_{input_data['workflow_name']}_{input_data['schedule_name']}_{datetime.now().strftime('%Y-%m-%d_%H-%M-%S')}" + key = f"data_package_{input_data['workflow_name']}_{input_data['schedule_name']}_{datetime.now().strftime(DATETIME_FORMAT)}" data = DataFrame(input_data['data']) held_data = DataFrame(input_data['held_data']) diff --git a/tests/activities/test_mongo.py b/tests/activities/test_mongo.py index 24ae1c4..8b8cd87 100644 --- a/tests/activities/test_mongo.py +++ b/tests/activities/test_mongo.py @@ -2,6 +2,7 @@ from datetime import datetime from unittest.mock import ANY, MagicMock, call, patch from pytest import fixture, mark from sientia_do.notifications.models import NotificationLevel +from sientia_do.temporal.constants import DATETIME_FORMAT_MS_WITH_TZ from scouter.activities.mongodb import MongoDB, clear_mongo_id @@ -95,7 +96,7 @@ async def test_load_latest_data_none_last_data_timestamp(mongodb_activity): 'name': 'test1', 'value': 1, 'inserted_at': datetime.strptime( - '2023-01-01 12:00:00.000000', '%Y-%m-%d %H:%M:%S.%f') + '2023-01-01 12:00:00.000000', DATETIME_FORMAT_MS_WITH_TZ) } ] @@ -137,7 +138,7 @@ async def test_load_latest_data_not_none_last_data_timestamp(mongodb_activity): 'name': 'test1', 'value': 1, 'inserted_at': datetime.strptime( - '2023-01-01 12:00:00.000000', '%Y-%m-%d %H:%M:%S.%f') + '2023-01-01 12:00:00.000000', DATETIME_FORMAT_MS_WITH_TZ) } ] @@ -154,7 +155,7 @@ async def test_load_latest_data_not_none_last_data_timestamp(mongodb_activity): { 'inserted_at': { '$gt': datetime.strptime( - '2023-01-01 12:00:00.000000', '%Y-%m-%d %H:%M:%S.%f') + '2023-01-01 12:00:00.000000', DATETIME_FORMAT_MS_WITH_TZ) } }, {"_id": 0}