C++簡單實現(xiàn)消息隊列的示例代碼
簡介
消息隊列是一種應用間的通訊方式,消息發(fā)送后可以立即放回,由消息系統(tǒng)來確保消息的可靠傳遞。消息發(fā)布者只需要將消息發(fā)布到消息隊列中,而不需要管誰來取。消息使用者只管從消息隊列中取消息而不管誰發(fā)布的。這樣發(fā)布者和使用者都不同知道對方的存在。
消息隊列普遍使用在生產(chǎn)者和消費者模型中。
- 優(yōu)點
- 應用解耦: 應用之間不用那么多的同步調(diào)用,發(fā)消息到消息隊列就行,消費者可以自己消費,消費生產(chǎn)者不用管了,降低應用之間的耦合。
- 降低延時: 應用之間用同步調(diào)用,需要等待對方響應,等待時間比較長,用消息之后,發(fā)送消息到消息隊列就行,應用就可以返回了,對客戶來講降低了應用延時。
- 削峰填谷:請求比較多的時候,應用處理不過來,會丟棄請求;請求比較少時,應用不飽和。
請求比較多時,把請求放到消息隊列,消費者按特定處理速度來處理,請求少時,也讓應用有事情可以做;能做到忙時不丟請求,閑時不閑置應用資源。
主流的消息隊列:Kafka、ActiveMQ、RabbitMQ、RocketMQ
下面使用C++實現(xiàn)一個簡單的消息隊列。
具體實現(xiàn)
- 消息隊列類
里面主要包含一個數(shù)組和一個隊列,都保存消息。
應用從數(shù)組中拿消息處理,當數(shù)組滿的時候,消息保存到隊列中,當數(shù)組中消息處理完,從隊列中取消息,再處理。
數(shù)組實現(xiàn)的是一個環(huán)形數(shù)組,記錄下消息寫游標和讀游標,超過數(shù)組大小,對數(shù)組大小取余。
//消息長度和消息
typedef std::pair<size_t, char*> msgPair;
#define MsgQueueSize 102400
class zMsgQueue
{
//消息對,first表示是否存放消息
typedef std::pair<bool, msgPair> msgQueue;
public:
zMsgQueue();
~zMsgQueue();
void* msgMalloc(const size_t len);
void msgFree(void* p);
//獲得一個消息
msgPair* get();
//放入一個消息
bool put(const void* msg, size_t msgLen);
//將隊列中的消息放到消息數(shù)組中
bool putMsgQueue2Arr();
//刪除一個消息
void erase();
bool empty();
bool msgArrEmpty();
private:
void clear();
// 保存正在處理的消息
msgQueue msgArr_[MsgQueueSize];
// 保存等待處理的消息
std::queue<msgPair> msgQueue_;
//消息寫游標
size_t queueWrite_;
//消息讀游標
size_t queueRead_;
};
實現(xiàn):
zMsgQueue::zMsgQueue()
{
bzero(msgArr_, sizeof(msgArr_));
queueWrite_ = 0;
queueRead_ = 0;
}
zMsgQueue::~zMsgQueue()
{
clear();
}
void* zMsgQueue::msgMalloc(const size_t len)
{
char* p = (char*)malloc(len + 1);
return (void*)(p + 1);
}
void zMsgQueue::msgFree(void* p)
{
free((char*)p - 1);
}
//獲得一個消息
msgPair* zMsgQueue::get()
{
if(queueRead_ >= MsgQueueSize)
return NULL;
if(msgArrEmpty())
putMsgQueue2Arr();
msgPair* ret = NULL;
if(msgArr_[queueRead_].first)
ret = &msgArr_[queueRead_].second;
return ret;
}
//放入一個消息
bool zMsgQueue::put(const void* msg, size_t msgLen)
{
char* buf = (char*)msgMalloc(msgLen);
if(buf)
{
bcopy(msg, buf, msgLen);
//先將隊列中的消息放到數(shù)組中
//數(shù)組中還有位置直接放到數(shù)組中
//沒有位置放到隊列中
if(!putMsgQueue2Arr() && !msgArr_[queueWrite_].first)
{
msgArr_[queueWrite_].first = true;
msgArr_[queueWrite_].second.first = msgLen;
msgArr_[queueWrite_].second.second = buf;
queueWrite_++;
queueWrite_ %= MsgQueueSize;
}
else
{
msgQueue_.push(std::make_pair(msgLen, buf));
}
return true;
}
return false;
}
//將隊列中的消息放到消息數(shù)組中
bool zMsgQueue::putMsgQueue2Arr()
{
bool isLeft = false;
while(!msgQueue_.empty())
{
if(!msgArr_[queueWrite_].first)
{
msgArr_[queueWrite_].first = true;
msgArr_[queueWrite_].second = msgQueue_.front();
queueWrite_++;
queueWrite_ %= MsgQueueSize;
msgQueue_.pop();
}
else
{
isLeft = true;
break;
}
}
return isLeft;
}
//刪除一個消息
void zMsgQueue::erase()
{
if(!msgArr_[queueRead_].first)
return;
msgFree(msgArr_[queueRead_].second.second);
msgArr_[queueRead_].second.second = NULL;
msgArr_[queueRead_].second.first = 0;
msgArr_[queueRead_].first = false;
queueRead_++;
queueRead_ %= MsgQueueSize;
}
void zMsgQueue::clear()
{
//隊列中還有消息
while(putMsgQueue2Arr())
{
//數(shù)組中還有消息
while(get())
{
erase();
}
}
//數(shù)組中還有消息
while(get())
{
erase();
}
}
bool zMsgQueue::empty()
{
if(putMsgQueue2Arr()) return false;
return msgArrEmpty();
}
bool zMsgQueue::msgArrEmpty()
{
if(queueRead_ == queueWrite_ && !msgArr_[queueRead_].first)
{
return true;
}
return false;
}
- 消息隊列的封裝
對消息隊列的封裝主要是為了對消息進行解析和處理。
消息解析和處理函數(shù)定義成了虛函數(shù),當需要使用消息隊列并處理消息時,只需要繼承消息隊列,然后重寫虛函數(shù),進行對應處理即可。
類中還使用到了讀寫鎖,當多線程的情況下,消息隊列是一個臨界資源,線程共享,需要進行上鎖。單線程的情況下不需要加鎖。
//T表示使用的消息隊列
//msgT表示消息的類型,有的需要消息頭,消息正文等,需要解析,這里是直接使用
template<class T=zMsgQueue, class msgT=char>
class messageQueue : public rwLocker
{
public:
messageQueue()
{}
~messageQueue()
{}
bool putMsg(const msgT* msg, const size_t msgLen)
{
rwLocker::wlock();
msgQueue_.put(msg, msgLen);
rwLocker::unlock();
return true;
}
//解析消息,處理消息
virtual bool msgParse(const msgT* msg, const size_t msgLen) = 0;
//獲取消息,解析消息,處理消息
bool doCmd()
{
rwLocker::wlock();
msgPair* msg = msgQueue_.get();
while(msg)
{
msgParse(msg->second, msg->first);
msgQueue_.erase();
msg = msgQueue_.get();
}
rwLocker::unlock();
return true;
}
bool empty()
{
return msgQueue_.empty();
}
private:
T msgQueue_;
};
- 讀寫鎖的封裝
讀寫鎖:可以多個線程進行讀,只能一個線程進行寫。寫時獨享資源,讀時共享資源。寫鎖的優(yōu)先級高。 - 為什么讀寫鎖需要讀鎖?
為了防止其他線程請求寫鎖。一個線程請求了讀鎖,其他線程在請求寫鎖會阻塞,但是請求讀鎖不會阻塞。一個線程請求了寫鎖,其他線程請求讀鎖和寫鎖都會阻塞。
#include <pthread.h>
class rwLock
{
public:
rwLock()
{
pthread_rwlock_init(&rwlc_, NULL);
}
~rwLock()
{
pthread_rwlock_destroy(&rwlc_);
}
void rlock()
{
pthread_rwlock_rdlock(&rwlc_);
}
void wlock()
{
pthread_rwlock_wrlock(&rwlc_);
}
void unlock()
{
pthread_rwlock_unlock(&rwlc_);
}
private:
pthread_rwlock_t rwlc_;
};
class rwLocker
{
public:
void rlock()
{
rwlc_.rlock();
}
void wlock()
{
rwlc_.wlock();
}
void unlock()
{
rwlc_.unlock();
}
private:
rwLock rwlc_;
};
- Makefile:
# ini1=main.cpp # in2=messageQueue.cpp out=main cc=g++ std=-std=c++11 -lpthread #$(out):$(in1) $(in2) $(out): main.cpp messageQueue.cpp rwlock.h $(cc) $^ -o $@ $(std) .PHONY:clean clean: rm -rf $(out)
- 代碼測試
實現(xiàn)一個類繼承消息隊列,重寫消息處理函數(shù)。
定義對象,調(diào)用doCmd函數(shù)即可。
#include "messageQueue.h"
class test : public messageQueue<>
{
bool msgParse(const char* msg, const size_t msgLen)
{
std::cout << msgLen << ":" << msg << std::endl;
return true;
}
};
int main()
{
//模擬客戶端發(fā)送消息
char buf[256] = "hello world!";
test t;
//消息隊列放消息
t.putMsg(buf, strlen(buf));
//處理消息
t.doCmd();
return 0;
}

到此這篇關(guān)于C++簡單實現(xiàn)消息隊列的示例代碼的文章就介紹到這了,更多相關(guān)C++ 消息隊列內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
C數(shù)據(jù)結(jié)構(gòu)中串簡單實例
這篇文章主要介紹了C數(shù)據(jù)結(jié)構(gòu)中串簡單實例的相關(guān)資料,需要的朋友可以參考下2017-06-06
Vs2022環(huán)境下安裝低版本.net framework的實現(xiàn)步驟
本文主要介紹了Vs2022環(huán)境下安裝低版本.net framework的實現(xiàn)步驟,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧2022-04-04
C語言編程gcc如何生成靜態(tài)庫.a和動態(tài)庫.so示例詳解
本文主要敘述了gcc如何生成靜態(tài)庫(.a)和動態(tài)庫(.so),幫助我們更好的進行嵌入式編程。因為有些時候,涉及安全,所以可能會提供靜態(tài)庫或動態(tài)庫供我們使用2021-10-10

