PostgreSQL broker for taskiq
This commit is contained in:
+16
@@ -39,6 +39,22 @@ broker = ResilientListQueueBroker(env.redis.url).with_result_backend(
|
||||
{% endif %}{% if use_dishka %}
|
||||
setup_dishka(container, broker)
|
||||
{% endif %}
|
||||
{%- elif taskiq_broker == 'postgres' -%}
|
||||
{% if use_dishka %}from dishka.integrations.taskiq import setup_dishka
|
||||
{% endif %}{% if 'scheduler' in backend_services %}from taskiq import TaskiqScheduler
|
||||
from taskiq.schedule_sources import LabelScheduleSource
|
||||
{% endif %}from taskiq_pg import AsyncpgBroker, AsyncpgResultBackend
|
||||
|
||||
{% if use_dishka %}from dependencies.container import container
|
||||
{% endif %}from utils.env import env
|
||||
|
||||
broker = AsyncpgBroker(env.db.connection_url).with_result_backend(
|
||||
AsyncpgResultBackend(env.db.connection_url)
|
||||
)
|
||||
{% if 'scheduler' in backend_services %}scheduler = TaskiqScheduler(broker, sources=[LabelScheduleSource(broker)])
|
||||
{% endif %}{% if use_dishka %}
|
||||
setup_dishka(container, broker)
|
||||
{% endif %}
|
||||
{%- else -%}
|
||||
{% if use_dishka %}from dishka.integrations.taskiq import setup_dishka
|
||||
{% endif %}from taskiq import InMemoryBroker{% if 'scheduler' in backend_services %}, TaskiqScheduler
|
||||
|
||||
Reference in New Issue
Block a user