worker.py 1.8 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849
  1. import logging
  2. from .service import AiTranslateService, SectionTimeout, Message
  3. from .decode_dataclass import ns_to_dataclass
  4. from .utils import is_stopped
  5. logger = logging.getLogger(__name__)
  6. class TaskFailException(Exception):
  7. def __init__(self, message="task fail"):
  8. self.message = message
  9. super().__init__(self.message)
  10. def handle_message(redis, ch, method, id, content_type, body, api_url, customer_timeout):
  11. MaxRetry = 3
  12. try:
  13. logger.info("process message start (%s) messages", len(body.payload))
  14. consumer = AiTranslateService(
  15. redis, ch, method, api_url, customer_timeout)
  16. messages = ns_to_dataclass([body], Message)
  17. consumer.process_translate(id, messages[0])
  18. ch.basic_ack(delivery_tag=method.delivery_tag) # 确认消息
  19. except SectionTimeout as e:
  20. # 时间到了,活还没干完 NACK 并重新入队
  21. ch.basic_nack(delivery_tag=method.delivery_tag, requeue=True)
  22. except Exception as e:
  23. # retry
  24. retryKey = f'{redis[1]}/message/retry/{id}'
  25. retry = 0
  26. if redis[0].exists(retryKey):
  27. retry = redis[0].get(retryKey)
  28. if retry > MaxRetry:
  29. logger.error(f'超过最大重试次数[{MaxRetry}],任务失败')
  30. # NACK 丢弃或者进入死信队列
  31. ch.basic_nack(delivery_tag=method.delivery_tag,
  32. requeue=False)
  33. raise TaskFailException
  34. retry = retry+1
  35. redis[0].set(retryKey, retry)
  36. # NACK 并重新入队
  37. logger.warning(f'消息处理错误,重新压入队列 [{retry}/{MaxRetry}]')
  38. ch.basic_nack(delivery_tag=method.delivery_tag,
  39. requeue=True)
  40. logger.error(f"error: {e}")
  41. logger.exception("发生异常")
  42. finally:
  43. is_stopped()