Skip to content

Commit dadb30c

Browse files
authored
[M2-11007] Added error logger to worker and scheduler (#2098)
* Added error logger to worker and scheduler * cqf
1 parent 3dd180e commit dadb30c

1 file changed

Lines changed: 9 additions & 2 deletions

File tree

src/broker.py

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@
22

33
import structlog
44
import taskiq_fastapi
5-
from taskiq import AsyncBroker, InMemoryBroker, TaskiqMessage, TaskiqMiddleware
5+
from taskiq import AsyncBroker, InMemoryBroker, TaskiqMessage, TaskiqMiddleware, TaskiqResult
66
from taskiq.formatters.json_formatter import JSONFormatter
77
from taskiq_aio_pika import AioPikaBroker
88
from taskiq_redis import RedisAsyncResultBackend
@@ -34,11 +34,18 @@ def pre_execute(
3434
return message
3535

3636

37+
class ErrorLoggerMiddleware(TaskiqMiddleware):
38+
"""Custom error logging middleware so Datadog receives errors"""
39+
40+
async def on_error(self, message: TaskiqMessage, result: TaskiqResult[Any], exception: BaseException) -> None:
41+
logger.error(f"Task {message.task_name} failed! ", exc_info=exception)
42+
43+
3744
if settings.env == "testing" or settings.env == "local":
3845
logger.info("Starting in memory broker")
3946
broker = InMemoryBroker().with_formatter(JSONFormatter())
4047

41-
middlewares = [StructlogMiddleware()]
48+
middlewares = [StructlogMiddleware(), ErrorLoggerMiddleware()]
4249
broker.add_middlewares(*middlewares)
4350

4451
taskiq_fastapi.init(broker, "main:app")

0 commit comments

Comments
 (0)