【问题标题】:Celery (Django) + RabbitMQ + nodejs server exchange dataCelery (Django) + RabbitMQ + nodejs 服务器交换数据
【发布时间】:2016-10-12 11:40:31
【问题描述】:

我有一个使用CeleryRabbitMQ 经纪人的Django 项目。现在我想从 NodeJS 服务器调用 django (celery) 任务。

在我的NodeJS 中,我使用的是amqplib。允许我向RabbitMQ发送任务:

amqp.connect('amqp://localhost', function(err, conn) {
  conn.createChannel(function(err, ch) {
    var q = 'celery';

    ch.assertQueue(q, {durable: true});
    ch.sendToQueue(q, new Buffer('What should I write here?'));
  });
});

我的问题是 celery 使用什么格式?我应该写什么给 Buffer 来调用 celery worker?

例如,在我的CELERY_ROUTES(django 设置)中,我有blabla.tasks.add

CELERY_ROUTES = {
   ...
   'blabla.tasks.add': 'high-priority',
}

如何调用这个blabla.tasks.add函数?

我尝试了很多方法,但芹菜工人给了我错误:Received and deleted unknown message. Wrong destination?!?

【问题讨论】:

    标签: node.js django rabbitmq celery


    【解决方案1】:

    我在这里找到了解决方案http://docs.celeryproject.org/en/latest/internals/protocol.html

    示例消息格式为:

    { "id": "4cc7438e-afd4-4f8f-a2f3-f46567e7ca77", "task": "celery.task.PingTask", "args": [], "kwargs": {}, "retries": 0, "eta": "2009-11-17T12:30:56.527191" }

    所以代码应该是:

    amqp.connect('amqp://localhost', function(err, conn) {
      conn.createChannel(function(err, ch) {
        var q = 'celery';
    
        ch.assertQueue(q, {durable: true});
        ch.sendToQueue(q, new Buffer('{"id": "this-is-soo-unique-id", "task": "blabla.tasks.add", "args": [1, 2], "kwargs": {}, "retries": 0}'), {
          contentType: 'application/json',
          contentEncoding: 'utf-8',
        });
      });
    });
    

    【讨论】:

      猜你喜欢
      • 2015-02-19
      • 2011-07-24
      • 2011-07-18
      • 2011-07-17
      • 2011-11-09
      • 2021-11-20
      • 1970-01-01
      • 2014-01-18
      • 2017-10-24
      相关资源
      最近更新 更多