Worker.js 4.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156
  1. const db = require('../DataBase/db')
  2. const path = require('path')
  3. const Logger = require('../../lib/Logger')
  4. const mq = require('.')
  5. const { mq: mqName } = require('./mqPrefix')
  6. class Worker {
  7. constructor() {
  8. this.logger = new Logger(
  9. path.join(__dirname, '../logs/Worker.log'),
  10. 'INFO'
  11. )
  12. this.handlers = {}
  13. this.running = false
  14. // 队列名
  15. this.taskQueue = mqName('task_queue')
  16. this.resultQueue = mqName('task_result_queue')
  17. // channel 名称(避免和别的模块冲突)
  18. this.channelName = 'worker_channel'
  19. }
  20. /**
  21. * 注册任务处理器
  22. */
  23. register(type, handler) {
  24. this.handlers[type] = handler
  25. this.logger.info(`注册处理器: ${type}`)
  26. }
  27. /**
  28. * 启动 Worker
  29. */
  30. async start() {
  31. if (this.running) return
  32. this.running = true
  33. this.logger.info('Worker 启动中...')
  34. try {
  35. const channel = await mq.getChannel(this.channelName)
  36. channel.on('close', () => {
  37. if (!this.running) return
  38. this.logger.warn('Worker channel 已关闭,准备重启消费')
  39. this.running = false
  40. setTimeout(() => {
  41. this.start().catch((e) => {
  42. this.logger.error('重启 Worker 失败: ' + (e?.stack || e))
  43. })
  44. }, 1000)
  45. })
  46. // 控制并发(重要)
  47. await channel.prefetch(5)
  48. // 确保队列存在
  49. await channel.assertQueue(this.taskQueue, { durable: true })
  50. await channel.assertQueue(this.resultQueue, { durable: true })
  51. // 开始消费
  52. await channel.consume(
  53. this.taskQueue,
  54. async (msg) => {
  55. if (!msg) return
  56. let content
  57. try {
  58. content = JSON.parse(msg.content.toString())
  59. } catch (err) {
  60. this.logger.error('消息解析失败: ' + err.message)
  61. channel.ack(msg)
  62. return
  63. }
  64. const { id, type, data } = content
  65. this.logger.info(`收到任务: ${id} 类型: ${type}`)
  66. const handler = this.handlers[type]
  67. if (!handler) {
  68. this.logger.error(`未找到处理器: ${type}`)
  69. channel.ack(msg)
  70. return
  71. }
  72. try {
  73. const result = await handler(data, {
  74. db,
  75. logger: this.logger
  76. })
  77. this.logger.info(`任务完成: ${id}`)
  78. await this.sendResult(channel, {
  79. id,
  80. success: true,
  81. result
  82. })
  83. channel.ack(msg)
  84. } catch (err) {
  85. this.logger.error(`任务失败: ${id} - ${err.stack}`)
  86. await this.sendResult(channel, {
  87. id,
  88. success: false,
  89. error: err.message
  90. })
  91. // 简单策略:失败直接 ack(避免死循环)
  92. channel.ack(msg)
  93. }
  94. },
  95. {
  96. noAck: false
  97. }
  98. )
  99. this.logger.info('Worker 启动成功')
  100. } catch (err) {
  101. this.logger.error('Worker 启动失败: ' + err.stack)
  102. }
  103. }
  104. /**
  105. * 发送结果
  106. */
  107. async sendResult(channel, data) {
  108. try {
  109. await mq.sendToQueueSafe(
  110. this.channelName,
  111. this.resultQueue,
  112. Buffer.from(JSON.stringify(data)),
  113. { persistent: true, contentType: 'application/json' }
  114. )
  115. } catch (err) {
  116. this.logger.error('结果发送失败: ' + err.message)
  117. }
  118. }
  119. /**
  120. * 停止 Worker
  121. */
  122. async stop() {
  123. this.running = false
  124. await mq.close()
  125. this.logger.info('Worker 已停止')
  126. }
  127. }
  128. module.exports = Worker