Страницы

Поиск по вопросам

Показаны сообщения с ярлыком celery. Показать все сообщения
Показаны сообщения с ярлыком celery. Показать все сообщения

пятница, 24 января 2020 г.

Import Error при запуске celery

#python #celery


У меня есть такая структура папок: 

- src
     -- \frontend
          --- \views
              ---- tasks.py
      -- \generator
          --- tasks.py


Во frontend/views/tasks я делаю: from generator import tasks.
При запуске селери из папки views питон пишет:  


  ImportError: No module named generator


В чем причина такой ошибки?
    


Ответы

Ответ 1



Родительская директория generator обязана быть в sys.path, чтобы from generator import tasks в этом случае работал. Достаточно из src директории запускать, чтобы путь в PYTHONPATH автоматически добавился: src$ python -m frontend.views.tasks Если хочется запускать из других директорий во время разработки, то можно создать setup.py для каждого пакета и установить их: $ pip install -e . Не стоит руками изменять sys.path в своём коде -- это ведёт к сюрпризам с неочевидным происхождением, например, см. Traps for the Unwary. Некоторые пакеты автоматически модифицируют sys.path, например, twisted использует _preamble.py, чтобы скрипты из bin директории могли без установки twisted пакет найти. Но подобная практика не поощряется, например, Pypy имел в прошлом похожий скрипт autopath.py, но сейчас он больше не используется -- он создаёт больше проблем чем решает. Пример проблем с импортом: Why python finds module instead of package if they have the same name?

пятница, 27 декабря 2019 г.

Celery накапливает периодические задания и выполняет пачкой

#python #django #celery


Celery контролируется by superviser. Есть файл для celery и для beat, но запускаются
3 процесса (worker дублируется). В логах всегда фигурирует Worker-1.
Есть выдержки из логов тут.

Проблема заключается в том, что таски, выполняющиеся каждые 2 минуты, запускаются
вовремя только первые 20-40 минут. Дальше все начинает идти не так. По прошествии очередных
2-х минут, задания перестают выполняться вовремя. Остаются только сообщение отправки

Scheduler: Sending due task contestsapp.tasks.apply_votes

без получения 

Received task: contestsapp.tasks.apply_votes[786b5cc6-5b08-40b0-9638-970a5ce6990f] 

Таким образом они накапливаются (неясно где, но celery процессы держат 90MB в памяти).
Через некоторое время воркер хватает все невыполненные задания и мгновенно (они примитивные)
выполняет. 

В итоге задания с периодом 2 мин выполняются каждые 10 минут по 5 раз (показатели
варьируются). Больше 5-ти одинаковых тасков не скапливается.

В остальном система работает. Отказов и исключений нет, брокер сбрасывался, база
синхронизировалась, PeriodicTask.objects.update(last_run_at=None) не помог, TZ везде
(даже ОС) стоит UTC.

PS
Иногда приходит пара десятков отложенных на дни-недели тасков. Такой "нагрузки" хрупкий
баланс двухминуток обычно не выдерживает, и проявляется эта проблема (но и без них
сценарий всегда один). 

Компоненты системы:


Rabbitmq - celery 3
Apache2 - wsgi - Django 1.9.5
Supervisor for celery and beat (but it runs one additional worker )


Файл настроек:

settings.py

CELERYBEAT_SCHEDULER = "djcelery.schedulers.DatabaseScheduler"
CELERY_TIMEZONE = 'UTC'
TIME_ZONE = 'UTC'


Пример задачи:

@periodic_task(ignore_result=True, run_every=crontab(minute='*/2'))
def apply_denial():
    print('apply_denial')
    denialDicts =  {id:cache.get(id) for id in cache.keys("denial:*") if  cache.ttl(id)<=290}
    for k, v in denialDicts.items():
        user = Profile.objects.get(id=k[(k.index('denial:') + 7):])
        user.denial = (list(set(v) - set(user.denial)) + user.denial)[:300]
        user.save()
        cache.delete(k)


Копал уже во все стороны...
    


Ответы

Ответ 1



А вы урл брокера задали??? Типа BROKER_URL = "amqp://user:pass@localhost/queue", где user и pass это логин пароль от раббита, а queue - очередь куда пуляются ваши сообщения P.S если так запускаете воркеры -A project worker --loglevel=INFO -A project worker --loglevel=INFO -A project beat --loglevel=INFO попробуйте заменить на: -A app_celery worker -l info -B -Q "наименование очереди" если никакую очередь явно не используете то -Q не нужно указывать

воскресенье, 8 декабря 2019 г.

Одна переменная на несколько процессов в Python

#python #многопоточность #websocket #tornado #celery


По совету @andreymal из моего предыдущего вопроса:  


  Если разными частями приложения являются разные процессы, то для этого надо организовывать
межпроцессное взаимодействие, чтобы один процесс работал с соединениями, а другие процессы
отправляли этому процессу сообщения. Можно для этого написать свой велосипед на сокетах,
можно для этого использовать готовые решения вроде RabbitMQ и Redis.  


Пытаюсь реализовать межпроцессное взаимодействие при помощи celery.  

Исходные данные и какую задачу нужно решить:
К серверу по вебсокетам должны подключаться клиенты (люди, использующие браузер)
и ожидать сообщений от сервера.
На сервере время от времени срабатывает процесс, генерирующий это самое сообщение.
Это сообщение (которое нужно отдать клиенту) можно попросить у процесса передать в
указанный скрипт.  

Что для этого сделано:  


Обработчик сокет-соединений на Python tornado;
.py скрипт, которому на stdin подается сообщение от процесса, который генерирует
само сообщение;
функция-таск для celery, которая должна отправить сообщение в браузер.


То есть, процесс генерирует сообщение -> вызывает скрипт, передавая ему сообщение
на stdin -> скрипт вызывает celery таск -> таск берет массив соединений и рассылает
им сообщение.

Суть проблемы:
Сокет сервер складывает все соединения в массив, но этот массив оказывается пустым
для celery таска.   

Собственно, вопрос: как завести один массив на несколько процессов, или где я налажал,
и как это делать правильно? 

UPD добавил упрощенный код для понимания того, как и что есть.
Сокет-сервер на tornado: 

class SocketHandler(tornado.websocket.WebSocketHandler):
    connections = []

    def open(self):
        self.__class__.connections.append(self)

    @classmethod
    def get_connections(cls):
        return cls.connections


и celery таск, который должен отправлять сообщения клиентам: 

from handlers import SocketHandler

@celery.task
def send_msg(msg):
    for conn in SocketHandler.get_connections():
        conn.write_message(msg)


Но список полученный путем SocketHandler.get_connections() в файле с celery тасками
оказывается пустым.
Запускаю все это в разных терминалах так:   

(venv)$ python app.py  # это торнадо апп, импортирующий приведенный handlers
(venv)$ celery -A tasks worker

    


Ответы

Ответ 1



Из описания проблемы и документации Celery я так понимаю, что таски выполняются в отдельных воркерах, которые может запускать сам Celery, но суть в том, что таск должен выполняться в одном процессе с сокет-сервером, тогда массив с сокетами ему будет виден. Пробежавшсь по диагонали по документации Celery, я подобного не нашёл; если это на нём и правда невозможно, то, видимо, стоит взять что-то попроще. Вообще я представляю это как-то так (несколько кривенько, но, думаю, суть должна быть понятна): def queue_thread(): while thread_should_work(): msg = redis.blpop("websocket_queue", timeout=5) if not msg: continue data = pickle.loads(msg[1]) send_message_to_all_sockets(data) Threading.thread(target=queue_thread).start()

пятница, 12 октября 2018 г.

Одна переменная на несколько процессов в Python

По совету @andreymal из моего предыдущего вопроса:
Если разными частями приложения являются разные процессы, то для этого надо организовывать межпроцессное взаимодействие, чтобы один процесс работал с соединениями, а другие процессы отправляли этому процессу сообщения. Можно для этого написать свой велосипед на сокетах, можно для этого использовать готовые решения вроде RabbitMQ и Redis.
Пытаюсь реализовать межпроцессное взаимодействие при помощи celery.
Исходные данные и какую задачу нужно решить: К серверу по вебсокетам должны подключаться клиенты (люди, использующие браузер) и ожидать сообщений от сервера. На сервере время от времени срабатывает процесс, генерирующий это самое сообщение. Это сообщение (которое нужно отдать клиенту) можно попросить у процесса передать в указанный скрипт.
Что для этого сделано:
Обработчик сокет-соединений на Python tornado; .py скрипт, которому на stdin подается сообщение от процесса, который генерирует само сообщение; функция-таск для celery, которая должна отправить сообщение в браузер.
То есть, процесс генерирует сообщение -> вызывает скрипт, передавая ему сообщение на stdin -> скрипт вызывает celery таск -> таск берет массив соединений и рассылает им сообщение.
Суть проблемы: Сокет сервер складывает все соединения в массив, но этот массив оказывается пустым для celery таска.
Собственно, вопрос: как завести один массив на несколько процессов, или где я налажал, и как это делать правильно?
UPD добавил упрощенный код для понимания того, как и что есть. Сокет-сервер на tornado:
class SocketHandler(tornado.websocket.WebSocketHandler): connections = []
def open(self): self.__class__.connections.append(self)
@classmethod def get_connections(cls): return cls.connections
и celery таск, который должен отправлять сообщения клиентам:
from handlers import SocketHandler
@celery.task def send_msg(msg): for conn in SocketHandler.get_connections(): conn.write_message(msg)
Но список полученный путем SocketHandler.get_connections() в файле с celery тасками оказывается пустым. Запускаю все это в разных терминалах так:
(venv)$ python app.py # это торнадо апп, импортирующий приведенный handlers (venv)$ celery -A tasks worker


Ответ

Из описания проблемы и документации Celery я так понимаю, что таски выполняются в отдельных воркерах, которые может запускать сам Celery, но суть в том, что таск должен выполняться в одном процессе с сокет-сервером, тогда массив с сокетами ему будет виден. Пробежавшсь по диагонали по документации Celery, я подобного не нашёл; если это на нём и правда невозможно, то, видимо, стоит взять что-то попроще.
Вообще я представляю это как-то так (несколько кривенько, но, думаю, суть должна быть понятна):
def queue_thread(): while thread_should_work(): msg = redis.blpop("websocket_queue", timeout=5) if not msg: continue data = pickle.loads(msg[1]) send_message_to_all_sockets(data)
Threading.thread(target=queue_thread).start()