mirror of
https://github.com/fastapi-practices/fastapi-best-architecture.git
synced 2026-09-21 13:12:24 +00:00
Update the celery configuration and tasks (#458)
* Update the celery configuration and tasks * fix message notifications
This commit is contained in:
+30
-30
@@ -10,7 +10,7 @@ __all__ = ['celery_app']
|
||||
|
||||
|
||||
def init_celery() -> celery.Celery:
|
||||
"""创建 celery 应用"""
|
||||
"""初始化 celery 应用"""
|
||||
|
||||
# TODO: Update this work if celery version >= 6.0.0
|
||||
# https://github.com/fastapi-practices/fastapi_best_architecture/issues/321
|
||||
@@ -18,54 +18,54 @@ def init_celery() -> celery.Celery:
|
||||
celery.app.trace.build_tracer = celery_aio_pool.build_async_tracer
|
||||
celery.app.trace.reset_worker_optimizations()
|
||||
|
||||
app = celery.Celery(
|
||||
'fba_celery',
|
||||
broker_connection_retry_on_startup=True,
|
||||
worker_pool=celery_aio_pool.pool.AsyncIOPool,
|
||||
trace=celery_aio_pool.build_async_tracer,
|
||||
)
|
||||
# Celery Schedule Tasks
|
||||
# https://docs.celeryq.dev/en/stable/userguide/periodic-tasks.html
|
||||
beat_schedule = task_settings.CELERY_SCHEDULE
|
||||
|
||||
# Celery Config
|
||||
# https://docs.celeryq.dev/en/stable/userguide/configuration.html
|
||||
_redis_broker = (
|
||||
f'redis://:{settings.REDIS_PASSWORD}@{settings.REDIS_HOST}:'
|
||||
f'{settings.REDIS_PORT}/{task_settings.CELERY_BROKER_REDIS_DATABASE}'
|
||||
broker_url = (
|
||||
(
|
||||
f'redis://:{settings.REDIS_PASSWORD}@{settings.REDIS_HOST}:'
|
||||
f'{settings.REDIS_PORT}/{task_settings.CELERY_BROKER_REDIS_DATABASE}'
|
||||
)
|
||||
if task_settings.CELERY_BROKER == 'redis'
|
||||
else (
|
||||
f'amqp://{task_settings.RABBITMQ_USERNAME}:{task_settings.RABBITMQ_PASSWORD}@'
|
||||
f'{task_settings.RABBITMQ_HOST}:{task_settings.RABBITMQ_PORT}'
|
||||
)
|
||||
)
|
||||
_amqp_broker = (
|
||||
f'amqp://{task_settings.RABBITMQ_USERNAME}:{task_settings.RABBITMQ_PASSWORD}@'
|
||||
f'{task_settings.RABBITMQ_HOST}:{task_settings.RABBITMQ_PORT}'
|
||||
)
|
||||
_result_backend = (
|
||||
result_backend = (
|
||||
f'redis://:{settings.REDIS_PASSWORD}@{settings.REDIS_HOST}:'
|
||||
f'{settings.REDIS_PORT}/{task_settings.CELERY_BACKEND_REDIS_DATABASE}'
|
||||
)
|
||||
_result_backend_transport_options = {
|
||||
'global_keyprefix': f'{task_settings.CELERY_BACKEND_REDIS_PREFIX}_',
|
||||
result_backend_transport_options = {
|
||||
'global_keyprefix': f'{task_settings.CELERY_BACKEND_REDIS_PREFIX}',
|
||||
'retry_policy': {
|
||||
'timeout': task_settings.CELERY_BACKEND_REDIS_TIMEOUT,
|
||||
},
|
||||
}
|
||||
|
||||
# Celery Schedule Tasks
|
||||
# https://docs.celeryq.dev/en/stable/userguide/periodic-tasks.html
|
||||
_beat_schedule = task_settings.CELERY_SCHEDULE
|
||||
|
||||
# Update celery settings
|
||||
app.conf.update(
|
||||
broker_url=_redis_broker if task_settings.CELERY_BROKER == 'redis' else _amqp_broker,
|
||||
result_backend=_result_backend,
|
||||
result_backend_transport_options=_result_backend_transport_options,
|
||||
timezone=settings.DATETIME_TIMEZONE,
|
||||
app = celery.Celery(
|
||||
'fba_celery',
|
||||
enable_utc=False,
|
||||
timezone=settings.DATETIME_TIMEZONE,
|
||||
beat_schedule=beat_schedule,
|
||||
broker_url=broker_url,
|
||||
broker_connection_retry_on_startup=True,
|
||||
result_backend=result_backend,
|
||||
result_backend_transport_options=result_backend_transport_options,
|
||||
task_cls='app.task.celery_task.base:TaskBase',
|
||||
task_track_started=True,
|
||||
beat_schedule=_beat_schedule,
|
||||
# TODO: Update this work if celery version >= 6.0.0
|
||||
worker_pool=celery_aio_pool.pool.AsyncIOPool,
|
||||
)
|
||||
|
||||
# Load task modules
|
||||
app.autodiscover_tasks(task_settings.CELERY_TASKS_PACKAGES)
|
||||
app.autodiscover_tasks(task_settings.CELERY_TASK_PACKAGES)
|
||||
|
||||
return app
|
||||
|
||||
|
||||
# 创建 celery 实例
|
||||
celery_app = init_celery()
|
||||
celery_app: celery.Celery = init_celery()
|
||||
|
||||
Reference in New Issue
Block a user