Enhance worker configuration in worker.py by adding parameters for max concurrent tasks and cached workflows to improve performance and scalability.
117 lines
3.5 KiB
Python
117 lines
3.5 KiB
Python
from temporalio import workflow, client
|
|
from temporalio.worker import Worker
|
|
|
|
with workflow.unsafe.imports_passed_through():
|
|
import sys
|
|
import os
|
|
from sientia_do.notifications.handlers import NotificationHandler
|
|
from sientia_do.temporal.utils.logger import get_logger
|
|
from scouter.activities.activities import Activities
|
|
from scouter.workflow.scouter import Scouter
|
|
from scouter.workflow.sub_workflows.core_scouter import CoreScouter
|
|
from scouter.workflow.fake_data import FakeData
|
|
from scouter.activities.faker import Faker
|
|
import asyncio
|
|
from scouter.utils.connectors_config import (
|
|
build_postgres_config,
|
|
build_redis_config,
|
|
build_mongodb_config
|
|
)
|
|
|
|
|
|
async def main():
|
|
host = os.getenv('TEMPORAL_HOST', 'localhost:7233')
|
|
|
|
logger = get_logger(__name__)
|
|
|
|
logger.info('Starting Worker...')
|
|
|
|
logger.info('Starting Notification Handler...')
|
|
|
|
notification_handler = NotificationHandler(
|
|
servers=os.getenv('KAFKA_BOOTSTRAP_SERVERS', 'http://localhost:9092'),
|
|
logger=logger,
|
|
project_name=os.getenv('PROJECT_NAME', 'scouter')
|
|
)
|
|
|
|
logger.info('Starting Activities...')
|
|
|
|
activities = Activities(
|
|
logger=logger,
|
|
notification_handler=notification_handler,
|
|
postgres_config=build_postgres_config(),
|
|
redis_config=build_redis_config(),
|
|
mongodb_config=build_mongodb_config()
|
|
)
|
|
|
|
logger.info('Starting Faker Activities...')
|
|
|
|
faker_activities = Faker(
|
|
logger=logger,
|
|
notification_handler=notification_handler,
|
|
bootstrap_servers=os.getenv(
|
|
'KAFKA_BOOTSTRAP_SERVERS', 'localhost:9092')
|
|
)
|
|
|
|
logger.info('Starting Temporal Client...')
|
|
|
|
temporal_client = await client.Client.connect(
|
|
target_host=host,
|
|
namespace=os.getenv('TEMPORAL_NAMESPACE', 'default')
|
|
)
|
|
|
|
logger.info('Starting Workers...')
|
|
|
|
workers = [
|
|
Worker(
|
|
temporal_client,
|
|
task_queue='scouter-queue',
|
|
workflows=[Scouter, CoreScouter],
|
|
activities=[
|
|
activities.load_latest_data,
|
|
activities.get_last_data_timestamp,
|
|
activities.put_last_data_timestamp,
|
|
activities.data_quality_gate,
|
|
activities.aggregate_data,
|
|
activities.group_and_hold_data,
|
|
activities.export_data_to_postgres,
|
|
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,
|
|
),
|
|
Worker(
|
|
temporal_client,
|
|
task_queue='fake_data-queue',
|
|
workflows=[FakeData],
|
|
activities=[
|
|
faker_activities.generate_and_send_data,
|
|
]
|
|
)
|
|
]
|
|
|
|
handlers = []
|
|
for w in workers:
|
|
handlers.append(w.run())
|
|
|
|
logger.info('Workers started successfully')
|
|
|
|
try:
|
|
await asyncio.gather(*handlers)
|
|
|
|
except BaseException as e:
|
|
logger.error("An unhandled exception occurred: %s", e, exc_info=True)
|
|
finally:
|
|
if notification_handler:
|
|
notification_handler.shutdown()
|
|
if activities:
|
|
activities.shutdown()
|
|
# Exit with a non-zero status code to indicate failure to Kubernetes
|
|
sys.exit(1)
|
|
|
|
if __name__ == '__main__':
|
|
asyncio.run(main())
|