【发布时间】:2020-05-16 03:51:37
【问题描述】:
此函数从消息容器中删除过期消息
void Broker::removeExpiredMessages(){
while(true){
std::this_thread::sleep_for(std::chrono::milliseconds(1000));
messageMut.lock();
std::cout<<Broker::getMessages().size()<<std::endl;
for(auto& i : Broker::getMessages()){
if(i.second.getHeader().expireAfter <= 0){
std::cout<<"del"<<std::endl;
Broker::getMessages().erase(i.first);
}
else
i.second.getHeader().setExpireAfter(i.second.getHeader().getExpireAfter()-1);
}
messageMut.unlock();
}
}
Message.hpp
class Message{
public:
struct Header{
std::string time;
int expireAfter;
std::string topicName;
int getExpireAfter() const {
return expireAfter;
}
void setExpireAfter(int expireAfter) {
this->expireAfter = expireAfter;
}
const std::string& getTime() const {
return time;
}
void setTime(const std::string &time) {
this->time = time;
}
const std::string& getTopicName() const {
return topicName;
}
void setTopicName(const std::string &topicName) {
this->topicName = topicName;
}
};
Message(){}
void setPayload(std::string _payload){
this->payload = _payload;
}
std::string getPayload()const{
return this->payload;
}
void setHeader(const std::string& _time, const int _expireAfter, const std::string _topicName){
this->header.setExpireAfter(_expireAfter);
this->header.setTime(_time);
this->header.setTopicName(_topicName);
}
Header getHeader()const{
return this->header;
}
private:
Header header;
std::string payload;
};
Broker.hpp
class Broker{
public:
static bool isSubscriberLoopRunning;
Broker(){std::cout<<"Broker()"<<std::endl;}
static std::queue<const char*>& getMessageQueue() {
return messageQueue;
}
static std::unordered_map<const char*, Message>& getMessages() {
return messages;
}
static std::vector<Subscriber>& getSubscriberList() {
return subscriberList;
}
static void pushNewSubscriber(Subscriber&);
static void receiveMessageFromPublisher(const Message&);
static void sendToSubscriber();
static void removeSubscriber(const Subscriber&);
static void removeExpiredMessages();
private:
static std::unordered_map<const char*,Message> messages;
static std::queue<const char*> messageQueue;
static std::vector<Subscriber> subscriberList;
};
我正在尝试制作消息传递系统,并且我有一个代理作为订阅者和发布者共享消息的中间件,消息存储在unordered_set<const char*,Message> messages。代理的一个功能是它会自动删除过期的消息,为此我使函数在与其父线程分离的单独线程中运行。
我的问题是它没有删除过期时间已到的消息,因为即使我正在这样做,expireAfter 变量也没有递减。我不明白为什么它不更新字段 expireAfter 的值。
【问题讨论】:
-
请提取minimal reproducible example 并将其作为您问题的一部分提供。注意它既完整又最小,你的既不是最小也不完整。作为新用户,也可以使用tour 并阅读How to Ask。另外,一旦你运行了这个,考虑在 codereview.stackexchange.com 上提交审查,我已经可以在那里发现一些禁忌。
-
我立即跳出来的是,您正在从容器中删除一个元素,同时迭代该容器。这通常是一个禁忌,但根据en.cppreference.com/w/cpp/container/unordered_map,
std::unordered_map::erase操作只会使您刚刚删除的迭代器无效。所以我想这很好。只是不和谐。 -
std::this_thread::sleep_for(std::chrono::milliseconds(1000))看起来很可疑。每当我看到明确的睡眠(或类似的)时,我的第一个想法就是“这里有一个错误”。这有什么意义? -
哦,等一下,您可能会在删除
i后立即访问它。没有bueno。 -
@UlrichEckhardt 好的,我会记住这一点
标签: c++ multithreading message-queue