🐛 fix(worker.py): add exception handling and retry mechanism to build_vertex task to handle SoftTimeLimitExceeded error and retry the task up to 3 times with a countdown of 2 seconds
✨ feat(worker.py): add import statement for SoftTimeLimitExceeded exception from celery.exceptions module to handle SoftTimeLimitExceeded error in build_vertex task
This commit is contained in:
parent
2dd43d3f99
commit
82e6b38ca4
1 changed files with 11 additions and 4 deletions
|
|
@ -1,6 +1,7 @@
|
||||||
from langflow.core.celery_app import celery_app
|
from langflow.core.celery_app import celery_app
|
||||||
from typing import Any, Dict, Optional
|
from typing import Any, Dict, Optional
|
||||||
from typing import TYPE_CHECKING
|
from typing import TYPE_CHECKING
|
||||||
|
from celery.exceptions import SoftTimeLimitExceeded
|
||||||
|
|
||||||
if TYPE_CHECKING:
|
if TYPE_CHECKING:
|
||||||
from langflow.graph.vertex.base import Vertex
|
from langflow.graph.vertex.base import Vertex
|
||||||
|
|
@ -11,13 +12,19 @@ def test_celery(word: str) -> str:
|
||||||
return f"test task return {word}"
|
return f"test task return {word}"
|
||||||
|
|
||||||
|
|
||||||
@celery_app.task
|
@celery_app.task(bind=True, soft_time_limit=30, max_retries=3)
|
||||||
def build_vertex(vertex: "Vertex") -> "Vertex":
|
def build_vertex(self, vertex: "Vertex") -> "Vertex":
|
||||||
"""
|
"""
|
||||||
Build a vertex
|
Build a vertex
|
||||||
"""
|
"""
|
||||||
vertex.build()
|
try:
|
||||||
return vertex
|
vertex.task_id = self.request.id
|
||||||
|
vertex.build()
|
||||||
|
return vertex
|
||||||
|
except SoftTimeLimitExceeded as e:
|
||||||
|
raise self.retry(
|
||||||
|
exc=SoftTimeLimitExceeded("Task took too long"), countdown=2
|
||||||
|
) from e
|
||||||
|
|
||||||
|
|
||||||
@celery_app.task(acks_late=True)
|
@celery_app.task(acks_late=True)
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue