from celery import Celery from kombu import Exchange, Queue from labelaudit import config from labelaudit.connectors import sentry from labelaudit.logic import prepare_asset_search sentry # pep8 app = Celery('label_audit_app') app.config_from_object(dict( BROKER_URL=config.BROKER_URL, CELERY_DEFAULT_QUEUE=config.LABEL_QUEUE, CELERY_ACCEPT_CONTENT=['pickle'], CELERY_QUEUES=(Queue( config.LABEL_QUEUE, Exchange(config.LABEL_QUEUE), routing_key='labelaudit.tasks.label_tasks.#'),), BROKER_TRANSPORT_OPTIONS=dict( queue_name_prefix=config.CELERY_QUEUE_PREFIX))) @app.task( name='labelaudit.tasks.label_tasks.process_audit', queue=config.LABEL_QUEUE) def process_audit(audit_id): """Fetches label audit from the label queue and prepares for asset search Args: youtube_audit_id: id of the audit to be processed """ prepare_asset_search.handle_releases(audit_id)