SIENTIAPDE-1199

Enhance logging in Gates, MongoDB, and Redis activities by replacing debug statements with info level logs, improving observability of data processing steps. Update worker configuration to adjust concurrency settings and enable autoscaling for task polling.
This commit is contained in:
vitor-aignosi
2025-08-20 12:39:43 -03:00
parent 3bbc2c4993
commit 2a76f2e3c6
4 changed files with 45 additions and 18 deletions

View File

@@ -89,8 +89,8 @@ class Gates(BaseActivity):
# Convert input data to DataFrame
df = DataFrame(input_data['data'])
self.debug(
f"Aggregating time series data: {df.to_string()}",
self.info(
f"Aggregating time series data for {len(df)} rows",
metadata=metadata
)
@@ -147,10 +147,16 @@ class Gates(BaseActivity):
}
result_df = DataFrame(list(result.values()))
self.debug(
f"Aggregated data:\n {result_df.to_string()}",
self.info(
f"Aggregated data has {len(result_df)} rows",
metadata=metadata
)
self.debug(
f"Aggregated data: {result_df.to_string()}",
metadata=metadata
)
return result_df.to_dict()
except Exception as e:
@@ -195,8 +201,8 @@ class Gates(BaseActivity):
data = DataFrame(input_data['data'])
model_tags = input_data['model_tags']
self.debug(
f"Applying quality gate to data: {data.to_string()}",
self.info(
f"Applying quality gate to data to {len(data)} rows",
metadata=metadata
)
@@ -249,8 +255,8 @@ class Gates(BaseActivity):
if policy == "DISCARD":
data = data[~data.index.isin(filtered_data.index)]
self.debug(
"Data quality gate applied",
self.info(
f"Data quality gate applied, final data has {len(data)} rows",
metadata=metadata
)
@@ -265,8 +271,18 @@ class Gates(BaseActivity):
"""
metadata = input_data['metadata']
self.info(
f"Writing metrics for {metadata['model_name']}",
metadata=metadata
)
metrics.LABORIOUS_DATA_WRITTEN_COUNT.labels(
pod_id=self.pod_id,
model_name=metadata['model_name'],
pipeline_name=metadata['workflow_name']
).inc()
self.info(
f"Metrics written for {metadata['model_name']}",
metadata=metadata
)

View File

@@ -87,7 +87,7 @@ class MongoDB(BaseActivity):
collection_name = input_data['collection_name']
last_data_timestamp = input_data['last_data_timestamp']
self.debug(
self.info(
f"Loading data from MongoDB: {input_data}",
metadata=metadata
)

View File

@@ -29,6 +29,8 @@ class Redis(RedisBase):
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}")
try:
data_hold = self.get(key)
except Exception as e:
@@ -42,7 +44,7 @@ class Redis(RedisBase):
)
raise e
self.debug(
self.info(
f"Last collected timestamp: {data_hold}",
metadata=metadata
)
@@ -60,6 +62,8 @@ class Redis(RedisBase):
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}")
data = DataFrame(input_data['data'])
if data.empty:
@@ -70,7 +74,7 @@ class Redis(RedisBase):
last_data_timestamp = data['inserted_at'].max()
self.debug(
self.info(
f"Last collected timestamp to insert: {last_data_timestamp}",
metadata=metadata
)
@@ -116,6 +120,8 @@ class Redis(RedisBase):
key = f"held_data_{input_data['workflow_name']}_{input_data['schedule_name']}"
self.info(f"Getting held data for {key}")
try:
data_hold = self.get(key)
except Exception as e:
@@ -137,6 +143,8 @@ class Redis(RedisBase):
)
return data_hold
self.info(f"Grouping and holding data for {len(data)} rows")
try:
# Remove possibly removed tags
@@ -195,8 +203,10 @@ class Redis(RedisBase):
)
raise e
self.info(f"Data held and melted has {len(data_hold_melted)} rows")
self.debug(
f"Data grouped and held successfully:\n {data_hold_melted.to_string()}",
f"Data held and melted:\n {data_hold_melted.to_string()}",
metadata=metadata
)

View File

@@ -1,5 +1,5 @@
from temporalio import workflow, client
from temporalio.worker import Worker
from temporalio.worker import Worker, PollerBehaviorAutoscaling
from temporalio.runtime import Runtime, TelemetryConfig, PrometheusConfig
with workflow.unsafe.imports_passed_through():
@@ -99,11 +99,12 @@ async def main():
activities.write_metrics,
activities.store_data_package,
],
max_concurrent_workflow_tasks=100,
max_concurrent_activities=100,
max_concurrent_local_activities=100,
max_concurrent_workflow_task_polls=100,
max_cached_workflows=50,
max_concurrent_workflow_tasks=50,
max_concurrent_activities=50,
max_concurrent_local_activities=50,
max_cached_workflows=200,
workflow_task_poller_behavior=PollerBehaviorAutoscaling(),
activity_task_poller_behavior=PollerBehaviorAutoscaling()
),
Worker(
temporal_client,