HR-ATS-Portal/backend/taskiq_management/broker_setup.py

60 lines
1.9 KiB
Python

"""Taskiq broker — Redis Streams + smart retry + DLQ.
Worker: taskiq worker taskiq_management.broker_setup:broker inbox.tasks inbox.sync_tasks cron_schdule.tasks taskiq_management.tasks job.candidate.bank_tasks
Scheduler: taskiq scheduler taskiq_management.broker_setup:scheduler inbox.sync_tasks cron_schdule.tasks
"""
from __future__ import annotations
import os
from dotenv import load_dotenv
from taskiq import TaskiqScheduler
from taskiq.middlewares import SmartRetryMiddleware
from taskiq.schedule_sources import LabelScheduleSource
from taskiq_redis import (
ListRedisScheduleSource,
RedisAsyncResultBackend,
RedisStreamBroker,
)
from taskiq_management.middleware import DeadLetterMiddleware
load_dotenv()
REDIS_URL=os.getenv("REDIS_URL","redis://localhost:6379/0")
QUEUE_NAME=os.getenv("TASKIQ_QUEUE_NAME","inbox")
# 2 retries after first failure → max_retries=3
MAX_RETRIES=int(os.getenv("TASKIQ_MAX_RETRIES","3"))
RETRY_DELAY=float(os.getenv("TASKIQ_RETRY_DELAY","5"))
result_backend=RedisAsyncResultBackend(redis_url=REDIS_URL)
schedule_source=ListRedisScheduleSource(url=REDIS_URL,prefix="taskiq:schedule")
broker=(
RedisStreamBroker(
url=REDIS_URL,
queue_name=QUEUE_NAME,
consumer_group_name=os.getenv("TASKIQ_CONSUMER_GROUP","taskiq"),
idle_timeout=int(os.getenv("TASKIQ_IDLE_TIMEOUT_MS","600000")),
)
.with_result_backend(result_backend)
.with_middlewares(
DeadLetterMiddleware(redis_url=REDIS_URL),
SmartRetryMiddleware(
default_retry_count=MAX_RETRIES,
default_retry_label=True,
default_delay=RETRY_DELAY,
use_jitter=True,
use_delay_exponent=True,
max_delay_exponent=float(os.getenv("TASKIQ_MAX_DELAY","120")),
schedule_source=schedule_source,
),
)
)
scheduler=TaskiqScheduler(
broker=broker,
sources=[schedule_source,LabelScheduleSource(broker)],
)