【问题标题】:rabbitMQ return queue is empty when i trying to get message by using outgoing AMQP from zato当我尝试使用来自 zato 的传出 AMQP 获取消息时,rabbitMQ 返回队列为空
【发布时间】:2018-11-05 15:39:00
【问题描述】:

我有一个由 ESB (zato) 调用的服务,该服务的作用是通过 AMQP 传出在 rabbitMQ 中发布消息,但是当我咨询 rabbitMQ 并获取消息时,答案是队列为空。这是服务在佐藤

from zato.server.service import Service

class HelloService(Service):
    def handle(self):

        # Request parameters
        msg = 'Hello AMQP broker!'
        out_name = 'My CRM connection'
        exchange = 'My exchange'
        routing_key = ''
        properties = {'app_id': 'ESB'}
        headers = {'X-Foo': 'bar'}

        # Send a message to the broker
        self.outgoing.amqp.send(msg, out_name, exchange, routing_key,
            properties, headers)

【问题讨论】:

  • 你有哪种exchange?你是通过什么方式将队列绑定到交易所的?
  • 我是 rabbitMQ 的新手,exchange 的类型是 declad 主题,我不知道将队列绑定到 exchange 的方式

标签: python rabbitmq esb zato


【解决方案1】:

从 zato 服务的兔子队列消费的完整工作示例如下:

兔子

  1. 创建交流
  2. 创建队列
  3. 将队列绑定到交换器
  4. 在 zato 中创建连接定义
  5. 在 zato 中创建传出 AMQP 连接定义
  6. 编写 zato 服务以发布或使用

前三个步骤可以通过多种方式完成,这是一个简单的 Python 脚本,您可以使用它(只需安装 kombu,然后单击):

import click
import os
import sys
import settings
from kombu import Connection, Exchange, Queue


BROKER_URL = 'amqp://{user}:{password}@{server}:{port}/{vhost}'.format(user=settings.RABBIT_USER,
                                                                       password=settings.RABBIT_PASS,
                                                                       server=settings.RABBIT_SERVER,
                                                                       port=settings.RABBIT_PORT,
                                                                       vhost=settings.RABBIT_VHOST)


@click.command()
@click.option('--remove/--no-remove', default=False, help='Remove current Queues/Exchanges.')
@click.option('--create/--no-create', default=False, help='Create needed Queues/Exchanges')
def job(remove, create):
    exchanges = {'dead_letter': Exchange(name=settings.DEAD_LETTER_EXCHANGE,
                                         type=settings.DEAD_LETTER_EXCHANGE_TYPE,
                                         durable=settings.DEAD_LETTER_EXCHANGE_DURABLE),
                 'results': Exchange(name=settings.RESULTS_EXCHANGE_NAME,
                                     type=settings.RESULTS_EXCHANGE_TYPE,
                                     durable=settings.RESULTS_EXCHANGE_DURABLE)}

    queues = {'dead_letter': Queue(name=settings.DEAD_LETTER_QUEUE,
                                   exchange=exchanges['dead_letter'],
                                   routing_key=settings.DEAD_LETTER_ROUTING,
                                   durable=settings.DEAD_LETTER_EXCHANGE_DURABLE),
              'results': Queue(name=settings.RESULTS_QUEUE_NAME,
                               exchange=exchanges['results'],
                               routing_key=settings.RESULTS_QUEUE_ROUTING,
                               durable=settings.RESULTS_EXCHANGE_DURABLE),
              'task': Queue(name=settings.TASK_QUEUE_NAME,
                            exchange=exchanges['results'],
                            routing_key=settings.TASK_ROUTING_KEY,
                            queue_arguments={
                                "x-message-ttl": settings.TASK_QUEUE_TTL,
                                "x-dead-letter-exchange": settings.DEAD_LETTER_EXCHANGE,
                                "x-dead-letter-routing-key": settings.DEAD_LETTER_ROUTING})}

    print 'using broker: {}'.format(BROKER_URL)

    with Connection(BROKER_URL) as conn:
        channel = conn.channel()
        if remove:
            # remove exchanges
            for (key, exchange) in exchanges.items():
                print 'removing exchange: {}'.format(exchange.name)
                bound_exchange = exchange(channel)
                bound_exchange.delete()

            # remove queues
            for (key, queue) in queues.items():
                print 'removing queue {} '.format(queues[key].name)
                bound_queue = queues[key](channel)
                bound_queue.delete()

        if create:
            # create exchanges
            for (key, exchange) in exchanges.items():
                print 'creating exchange: {}'.format(exchange.name)
                bound_exchange = exchange(channel)
                bound_exchange.declare()

            # add queues
            for (key, queue) in queues.items():
                # if key in exchanges:
                print 'binding queue {} to exchange {} with routing key {}'.format(queue.name,
                                                                                   queue.exchange.name,
                                                                                   queue.routing_key)
                bound_queue = queue(channel)
                bound_queue.declare()


if __name__ == '__main__':
    job()

还有设置文件:

# rabbit stuff
RABBIT_SERVER = 'localhost'
RABBIT_USER = 'guest'
RABBIT_PASS = 'guest'
RABBIT_PORT = 5672
RABBIT_VHOST = '/'

# default task queue
TASK_EXCHANGE_NAME = 'test.service.request'
TASK_EXCHANGE_TYPE = 'direct'
TASK_EXCHANGE_DURABLE = True
TASK_QUEUE_NAME = 'test.service.request'
TASK_ROUTING_KEY = 'request'
TASK_QUEUE_TTL = 604800000

# dead letter settings
DEAD_LETTER_EXCHANGE = 'test.service.deadletter'
DEAD_LETTER_EXCHANGE_TYPE = 'direct'
DEAD_LETTER_EXCHANGE_DURABLE = True
DEAD_LETTER_QUEUE = 'test.service.deadletter'
DEAD_LETTER_ROUTING = 'deadletter'

# results settings
RESULTS_EXCHANGE_NAME = 'test.service.results'
RESULTS_EXCHANGE_TYPE = 'direct'
RESULTS_EXCHANGE_DURABLE = True
RESULTS_QUEUE_NAME = 'test.service.results'
RESULTS_QUEUE_ROUTING = 'results'

现在让我们使用 python 2.7 在新的 virtualenv 上创建运行上述脚本的队列:

$ virtualenv rabbit_test
New python executable in /home/ivan/rabbit_test/bin/python
Installing setuptools, pip, wheel...done.

$ source /home/ivan/rabbit_test/bin/activate

$ pip install kombu
Collecting kombu
...
$ pip install click
Collecting click
...

复制上面的脚本

$ mkdir ~/rabbit_test/app
$ vi ~/rabbit_test/app/create_queues.py
$ vi ~/rabbit_test/app/settings.py

并运行 create_queues.py。

$ cd ~/rabbit_test/app
$ python create_queues.py --create
using broker: amqp://guest:guest@localhost:5672//
creating exchange: test.service.results
creating exchange: test.service.deadletter
binding queue test.service.request to exchange test.service.results with routing key request
binding queue test.service.results to exchange test.service.results with routing key results
binding queue test.service.deadletter to exchange test.service.deadletter with routing key deadletter

您可以使用cli 工具或management plugin 验证交换和队列是否在rabbit 上:

$ rabbitmqadmin list exchanges
+-------------------------+---------+
|          name           |  type   |
+-------------------------+---------+
| test.service.deadletter | direct  |
| test.service.results    | direct  |
+-------------------------+---------+

$ rabbitmqadmin list queues
+-------------------------+----------+
|          name           | messages |
+-------------------------+----------+
| test.service.deadletter | 0        |
| test.service.request    | 0        |
| test.service.results    | 0        |
+-------------------------+----------+

$ rabbitmqadmin list bindings
+-------------------------+-------------------------+-------------------------+
|         source          |       destination       |       routing_key       |
+-------------------------+-------------------------+-------------------------+
|                         | test.service.deadletter | test.service.deadletter |
|                         | test.service.request    | test.service.request    |
|                         | test.service.results    | test.service.results    |
| test.service.deadletter | test.service.deadletter | deadletter              |
| test.service.results    | test.service.request    | request                 |
| test.service.results    | test.service.results    | results                 |
+-------------------------+-------------------------+-------------------------+

现在 zato 部分(步骤 4,5 和 6)可以使用公共 api 或 webadmin 完成,我将向您展示如何使用公共 api 来完成,但通过 UI 更容易做到因为这只会做很少的次数。

创建 AMQP 连接定义 doc

$ curl -X POST -H "Authorization: Basic cHViYXBpOjEyMw==" -d '{
    "cluster_id": 1,
    "name": "SO_Test",
    "host": "127.0.0.1",
    "port": "5672",
    "vhost": "/",
    "username": "guest",
    "frame_max": 131072,
    "heartbeat": 10
}' "http://localhost:11223/zato/json/zato.definition.amqp.create"

{
  "zato_env": {
    "details": "",
    "result": "ZATO_OK",
    "cid": "K04DWBPMYF8A7768C7N482E75YM3"
  },
  "zato_definition_amqp_create_response": {
    "id": 2,
    "name": "SO_Test"
  }
}

为我们的 AMQP 连接设置密码 doc

$ curl -X POST -H "Authorization: Basic cHViYXBpOjEyMw=="  -d '{
    "id": 2,
    "password1": "guest",
    "password2": "guest"
}' "http://localhost:11223/zato/json/zato.definition.amqp.change-password"

{
  "zato_env": {
    "details": "",
    "result": "ZATO_OK",
    "cid": "K07K9YY21XZAX4QKWJB3ZFXN2ZFT"
  }
}

创建传出 AMQP 连接定义 doc

curl -X POST -H "Authorization: Basic cHViYXBpOjEyMw==" -d '{
    "cluster_id": 1,
    "name": "SO Test",
    "is_active": true,
    "def_id": 2,
    "delivery_mode": 1,
    "priority": 6,
    "content_type": "application/json",
    "content_encoding": "utf-8",
    "expiration": 30000
}' "http://localhost:11223/zato/json/zato.outgoing.amqp.create"

{
  "zato_outgoing_amqp_create_response": {
    "id": 1,
    "name": "SO Test"
  },
  "zato_env": {
    "details": "",
    "result": "ZATO_OK",
    "cid": "K05F2CR954BFNBP14KGTM26V47PC"
  }
}

最后是要发送消息的服务

from zato.server.service import Service

class HelloService(Service):
    def handle(self):
        # Request parameters 
        msg = 'Hello AMQP broker!'
        out_name = 'SO Test'
        exchange = 'test.service.results'
        routing_key = 'request'
        properties = {'app_id': 'ESB', 'user_id': 'guest'}
        headers = {'X-Foo': 'bar'}

        # Send a message to the broker
        info = self.outgoing.amqp.send(msg, out_name, exchange, routing_key,
            properties, headers)
        self.logger.info(info)

如果您要使用属性 user_id,它必须与连接 user_id 匹配,否则请求将失败。

另外请注意,我在这里创建了一个死信交换,如果消息仍在test.service.request 队列中,则消息将在 30 秒后发送到这里

最后一步是测试

为了验证消息是否传递到我们的队列,我们​​可以创建一个 http/soap 通道或直接调用服务,我正在使用公共 api 进行后者。

curl -X POST -H "Authorization: Basic cHViYXBpOjEyMw==" -d '{
   "name": "test.hello-service",
   "data_format": "json"
}' "http://localhost:11223/zato/json/zato.service.invoke"

{
  "zato_env": {
    "details": "",
    "result": "ZATO_OK",
    "cid": "K050J64QQ8FXASXHKVCAQNC4JC4N"
  },
  "zato_service_invoke_response": {
    "response": ""
  }
}

然后我们检查队列中我们刚刚发送的消息:

$ rabbitmqadmin get queue=test.service.request requeue=true
+-------------+----------------------+---------------+--------------------+---------------+------------------+-------------+
| routing_key |       exchange       | message_count |      payload       | payload_bytes | payload_encoding | redelivered |
+-------------+----------------------+---------------+--------------------+---------------+------------------+-------------+
| request     | test.service.results | 0             | Hello AMQP broker! | 18            | string           | False       |
+-------------+----------------------+---------------+--------------------+---------------+------------------+-------------+

记得检查rabbit和zato服务器日志,以防仍有问题。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-09-16
    • 2013-05-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多