diff --git a/requirements.txt b/requirements.txt index 0f19531..1dc88da 100644 --- a/requirements.txt +++ b/requirements.txt @@ -5,6 +5,6 @@ asyncua redis aiokafka pymongo -git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.5.2 +git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.5.3 pydruid[pandas] prometheus-client \ No newline at end of file diff --git a/scouter/activities/redis.py b/scouter/activities/redis.py index 6e2b885..968c320 100644 --- a/scouter/activities/redis.py +++ b/scouter/activities/redis.py @@ -94,10 +94,10 @@ class Redis(SientiaMonitoring): metadata = input_data['metadata'] key = f'last_data_timestamp:{input_data["workflow_name"]}:{input_data["schedule_name"]}' - self.info(f'Getting last data timestamp for {key}') + self.info(f'Getting last data timestamp for {key}', metadata=metadata) try: - data_hold = await self.redis_repository.get(key) + data_hold = await self.redis_repository.get(key, metadata=metadata) except Exception as e: await self.send_notification_async( metadata=metadata, @@ -142,7 +142,7 @@ class Redis(SientiaMonitoring): metadata = input_data['metadata'] key = f'last_data_timestamp:{input_data["workflow_name"]}:{input_data["schedule_name"]}' - self.info(f'Putting last data timestamp for {key}') + self.info(f'Putting last data timestamp for {key}', metadata=metadata) data = DataFrame(input_data['data']) @@ -155,7 +155,7 @@ class Redis(SientiaMonitoring): self.info(f'Last collected timestamp to insert: {last_data_timestamp}', metadata=metadata) try: - await self.redis_repository.set(key, last_data_timestamp, ttl=60 * 60 * 5) + await self.redis_repository.set(key, last_data_timestamp, ttl=60 * 60 * 5, metadata=metadata) except Exception as e: await self.send_notification_async( metadata=metadata, @@ -205,10 +205,10 @@ class Redis(SientiaMonitoring): key = f'held_data_{input_data["workflow_name"]}_{input_data["schedule_name"]}' - self.info(f'Getting held data for {key}') + self.info(f'Getting held data for {key}', metadata=metadata) try: - data_hold = await self.redis_repository.get(key) + data_hold = await self.redis_repository.get(key, metadata=metadata) except Exception as e: await self.send_notification_async( metadata=metadata, @@ -252,7 +252,7 @@ class Redis(SientiaMonitoring): data['timestamp'].max() if not data.empty else data_hold['timestamp'] ) - await self.redis_repository.set(key, data_hold, ttl=retention_time) + await self.redis_repository.set(key, data_hold, ttl=retention_time, metadata=metadata) data_hold_df = DataFrame(data_hold, index=[0]) data_hold_melted = data_hold_df.melt( @@ -298,7 +298,7 @@ class Redis(SientiaMonitoring): cache = {'data': data.to_dict(), 'held_data': held_data.to_dict()} try: - await self.redis_repository.set(key, cache, ttl=120) + await self.redis_repository.set(key, cache, ttl=120, metadata=metadata) except Exception as e: await self.send_notification_async( metadata=metadata,