En el proyecto en el que estoy trabajando actualmente, no tengo permitido usar un ORM, así que hice el mío propio
Funciona muy bien, pero tengo problemas con Celery y su concurrencia. Durante un tiempo, lo configuré en 1 (usando --concurrency=1 ) pero estoy agregando nuevas tareas que requieren más tiempo para procesarse de lo que necesitan para ejecutarse con celery beat, lo que provoca una gran acumulación de tareas.
Cuando configuro la concurrencia de apio en> 1, esto es lo que sucede (pastebin porque es grande):
¿Alguna idea de cómo podría implementar algún tipo de bloqueo/espera en los otros procesos para que los diferentes trabajadores no se crucen entre sí?
Editar: aquí es donde configuro mi instancia de PyMySQL y cómo se manejan la apertura y el cierre
PyMSQL no permite que los subprocesos compartan la misma conexión (el módulo se puede compartir, pero los subprocesos no pueden compartir una conexión). Su clase Model está reutilizando la misma conexión en todas partes.
Entonces, cuando diferentes trabajadores llaman a los modelos para hacer consultas, están usando el mismo objeto de conexión, lo que genera conflictos.
Asegúrese de que sus objetos de conexión sean locales de subprocesos. En lugar de tener un atributo de clase db , considere un método que recuperará un objeto de conexión local de subproceso, en lugar de reutilizar uno potencialmente creado en un subproceso diferente.
Por ejemplo, cree su conexión en la tarea .
En este momento, está utilizando una conexión global en todas partes para cada modelo.
# Connect to the database connection = pymysql.connect(**database_config) class Model(object): """ Base Model class, all other Models will inherit from this """ db = connection Para evitar esto, puede crear la base de datos en el método __init__ en su lugar...
class Model(object): """ Base Model class, all other Models will inherit from this """ def __init__(self, *args, **kwargs): self.db = pymysql.connect(**database_config)Sin embargo, esto puede no ser eficiente/práctico porque cada instancia del objeto db creará una sesión.
Para mejorar esto, podría usar un enfoque usando threading.local para mantener las conexiones locales a los hilos.
class Model(object): """ Base Model class, all other Models will inherit from this """ _conn = threading.local() @property def db(self): if not hasattr(self._conn, 'db'): self._conn.db = pymysql.connect(**database_config) return self._conn.dbTenga en cuenta que una solución local de subprocesos funciona asumiendo que está utilizando un modelo de concurrencia de subprocesos. Tenga en cuenta también que el apio utiliza múltiples procesos (prefork) de forma predeterminada. Esto puede o no ser un problema. Si es un problema, es posible que pueda solucionarlo si cambia los trabajadores para que usen eventlet en su lugar.