rabbitmq-c初探 (二)

2014-11-24 00:04:26 · 作者: · 浏览: 15
channel");

{

amqp_queue_declare_ok_t *r = amqp_queue_declare(conn, 1, AMQP_EMPTY_BYTES/*amqp_empty_bytes*/, 0, 0, 0, 1,

AMQP_EMPTY_TABLE/*amqp_empty_table*/);

die_on_amqp_error(amqp_get_rpc_reply(conn), "Declaring queue");

queuename = amqp_bytes_malloc_dup(r->queue);

if (queuename.bytes == NULL) {

fprintf(stderr, "Out of memory while copying queue name");

return 1;

}

}

amqp_queue_bind(conn, 1, queuename, amqp_cstring_bytes(exchange), amqp_cstring_bytes(bindingkey),

AMQP_EMPTY_TABLE/*amqp_empty_table*/);

die_on_amqp_error(amqp_get_rpc_reply(conn), "Binding queue");

amqp_basic_consume(conn, 1, queuename, AMQP_EMPTY_BYTES/*amqp_empty_bytes*/, 0, 1, 0, AMQP_EMPTY_TABLE/*amqp_empty_table*/);

die_on_amqp_error(amqp_get_rpc_reply(conn), "Consuming");

run(conn);

die_on_amqp_error(amqp_channel_close(conn, 1, AMQP_REPLY_SUCCESS), "Closing channel");

die_on_amqp_error(amqp_connection_close(conn, AMQP_REPLY_SUCCESS), "Closing connection");

die_on_error(amqp_destroy_connection(conn), "Ending connection");

return 0;

}


2 producer
2.1 producer.pro的内容
SOURCES=utils.cpp amqp_producer.cpp platform_utils.cpp
HEADERS=utils.h
VPATH+=/usr/include
CONFIG+=release
TARGET=producer
LIBS += -lrabbitmq
2.2 amqp_producer.cpp代码
这里的代码来自于rabbitmq-c-v0.3.0 具体查看https://github.com/alanxz/rabbitmq-c/blob/rabbitmq-c-v0.3.0/examples/amqp_producer.c。(对于几个特殊的宏引用作了调整)

#include

#include

#include

#include

#include

#include

#include "utils.h"

#define SUMMARY_EVERY_US 1000000


static void send_batch(amqp_connection_state_t conn,

char const *queue_name,

int rate_limit,

int message_count)

{

uint64_t start_time = now_microseconds();

int i;

int sent = 0;

int previous_sent = 0;

uint64_t previous_report_time = start_time;

uint64_t next_summary_time = start_time + SUMMARY_EVERY_US;

char message[256];

amqp_bytes_t message_bytes;


for (i = 0; i < (int)sizeof(message); i++) {

message[i] = i & 0xff;

}

message_bytes.len = sizeof(message);

message_bytes.bytes = message;

for (i = 0; i < message_count; i++) {

uint64_t now = now_microseconds();

die_on_error(amqp_basic_publish(conn,1,amqp_cstring_bytes("amq.direct"),amqp_cstring_bytes(queue_name),

0,0,NULL,message_bytes),"Publishing");

sent++;

if (now > next_summary_time) {

int countOverInterval = sent - previous_sent;

double intervalRate = countOverInterval / ((now - previous_report_time) / 1000000.0);

printf("%d ms: Sent %d - %d since last report (%d Hz)\n",(int)(now - start_time) / 1000, sent,

countOverInterval, (int) intervalRate);

previous_sent = sent;

previous_report_time = now;

next_summary_time += SUMMARY_EVERY_US;

}

while (((i * 1000000.0) / (now - start_time)) > rate_limit) {

microsleep(2000);

now = now_microseconds();

}

}

{

uint64_t stop_time = now_microseconds();

int total_delta = stop_time - start_time;

printf("PRODUCER - Message count: %d\n", message_count);

printf("Total time, milliseconds: %d\n", total_delta / 1000);

printf("Overall messages-per-second: %g\n", (message_count / (total_delta / 1000000.0)));

}

}


int main(int argc, char const * const *argv) {

char const *hostname;

int port;

int rate_limit;

int message_count;

int sockfd