usando apio con SQS en la aplicación del matraz
pero el apio recibe la misma tarea dos veces con la misma identificación de tarea al mismo tiempo ,
corriendo trabajador como este,
celery worker -A app.jobs.run -l info --pidfile=/var/run/celery/celery.pid --logfile=/var/log/celery/celery.log --time-limit=7200 --concurrency=8
aquí están los troncos de apio
[2019-11-29 08:07:35,464: INFO/MainProcess] Received task: app.jobs.booking.bookFlightTask[657985d5-c3a3-438d-a524-dbb129529443] [2019-11-29 08:07:35,465: INFO/MainProcess] Received task: app.jobs.booking.bookFlightTask[657985d5-c3a3-438d-a524-dbb129529443] [2019-11-29 08:07:35,471: WARNING/ForkPoolWorker-4] in booking funtion1 [2019-11-29 08:07:35,473: WARNING/ForkPoolWorker-3] in booking funtion1 [2019-11-29 08:07:35,537: WARNING/ForkPoolWorker-3] book_request_pp [2019-11-29 08:07:35,543: WARNING/ForkPoolWorker-4] book_request_ppmisma tarea recibida dos veces y ambas se ejecutan simultáneamente,
usando celery==4.4.0rc4 , boto3==1.9.232, kombu==4.6.6 con SQS en matraz de Python.
En SQS, el tiempo de espera de visibilidad predeterminado es de 30 minutos, y mi tarea es no tener ETA y no reconocer
mi tarea.py
from app import app as flask_app from app.jobs.run import capp from flask_sqlalchemy import SQLAlchemy db = SQLAlchemy(flask_app) class BookingTasks: def addBookingToTask(self): request_data = request.json print ('in addBookingToTask',request_data['request_id']) print (request_data) bookFlightTask.delay(request_data) return 'addBookingToTask added' @capp.task(max_retries=0) def bookFlightTask(request_data): task_id = capp.current_task.request.id try: print ('in booking funtion1') ----mi archivo de configuración, config.py
import os from urllib.parse import quote_plus aws_access_key = quote_plus(os.getenv('AWS_ACCESS_KEY')) aws_secret_key = quote_plus(os.getenv('AWS_SECRET_KEY')) broker_url = "sqs://{aws_access_key}:{aws_secret_key}@".format( aws_access_key=aws_access_key, aws_secret_key=aws_secret_key, ) imports = ('app.jobs.run',) ## Using the database to store task state and results. result_backend = 'db' + '+' + os.getenv('SQLALCHEMY_DATABASE_URI')y, por último, mi archivo de aplicación de apio, run.py
from __future__ import absolute_import, unicode_literals import os from celery import Celery from flask import Flask from app import app as flask_app import sqlalchemy capp = Celery() capp.config_from_object('app.jobs.config') # Optional configuration, see the capplication user guide. capp.conf.update( result_expires=3600, ) # SQS_QUEUE_NAME is like 'celery_test.fifo' , .fifo is required capp.conf.task_default_queue = os.getenv('FLIGHT_BOOKINNG_SQS_QUEUE_NAME') if __name__ == '__main__': capp.start()El tiempo de espera de visibilidad predeterminado de SQS es de 30 s. Debe actualizar el valor de configuración de apio: broker_transport_options={'visibility_timeout': 3600} .
Cuando el apio vaya a crear la cola, establecerá el tiempo de espera de visibilidad en 1 hora.
NOTA: si especifica task_default_queue y la cola ya se ha creado sin especificar broker_transport_options={'visibility_timeout': 3600} , celery no actualizará el tiempo de espera de visibilidad cuando se reinicie con broker_transport_options={'visibility_timeout': 3600} . Deberá eliminar la cola y hacer que el apio la vuelva a crear.