You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

NodeJS调用Celery任务报错:InvalidTaskError参数格式异常

解决Celery接收NodeJS任务时的InvalidTaskError问题

错误根源

日志中的AttributeError: 'str' object has no attribute 'get'表明Celery未将消息体解析为JSON对象,而是当作字符串处理,核心问题是AMQP消息属性格式错误,同时任务全名不匹配。

修复方案

1. 修正NodeJS的消息发送属性

amqplib的sendToQueue选项参数格式错误,需将content-type、correlation_id等直接作为顶级属性,而非嵌套在properties中;同时修正任务全名为Celery实际注册的名称。

2. 修正后的NodeJS任务函数

const amqp = require('amqplib');

async function delayScrapNewMovies() {
  const connection = await amqp.connect('amqps://');
  const channel = await connection.createChannel();
  const queueName = 'celery';
  
  // Celery任务全名:应用名.函数名
  const taskName = 'scraper.scrapNewMovies';
  const taskId = '123';

  // 正确的AMQP发送选项
  const sendOptions = {
    correlationId: taskId,
    contentType: 'application/json',
    contentEncoding: 'utf-8',
    headers: {
      task: taskName,
      id: taskId,
    }
  };

  // 符合Celery协议的消息体
  const message = {
    task: taskName,
    args: [],
    kwargs: {
      "url": "https://example.com/movies",
      "limit": 10
    },
    id: taskId,
    retries: 0,
    eta: null,
    expires: null
  };

  await channel.assertQueue(queueName);
  channel.sendToQueue(queueName, Buffer.from(JSON.stringify(message)), sendOptions);
  
  // 延迟关闭连接,确保消息发送完成
  setTimeout(() => connection.close(), 500);
}
  
module.exports = {delayScrapNewMovies};

3. 确认Celery Worker启动命令

启动Worker时指定正确的应用模块,确保任务被注册:

celery -A scraper worker --loglevel=info

核心要点

  • content-type必须设为application/json,否则Celery不会自动解析JSON消息体。
  • 任务名称必须与Celery中注册的全名一致(scraper.scrapNewMovies,而非app.celery.scrapNewMovies)。
  • 发送消息后需关闭AMQP连接,避免资源泄漏。

内容的提问来源于stack exchange,提问作者hsn

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.28 22:20:41