【发布时间】:2011-03-20 16:02:55
【问题描述】:
我编写了一个用于解析日志文件的多线程应用程序。它基于来自http://drdobbs.com/cpp/184401518?pgno=5 的互斥缓冲区示例。
这个想法是有一个缓冲区类,它具有将项目放入缓冲区并从中取出项目的功能。使用条件处理读取和写入线程的同步。当缓冲区未满时,新项目被写入缓冲区,当它不为空时,正在从中读取项目。否则线程将等待。
该示例使用固定数量的项目来处理,因此我将读取线程更改为在有文件输入时运行,而处理线程在有输入或缓冲区中有项目时运行。
我的问题是,如果我使用 1 个读取和 1 个处理线程,一切正常且稳定。当我添加另一个处理线程时,性能会有很大提升,即使经过 10.000 次测试运行它仍然很稳定。
现在,当我添加另一个处理器线程(1 个读取,3 个处理)时,程序似乎会定期(但不是每次)挂起(死锁?)并且等待缓冲区填充或变空。
为什么 2 个线程在做同样的事情时同步稳定,而其中 3 个线程崩溃了?
我是 C++ 线程的新手,所以也许你们中任何更有经验的编码人员知道什么会导致这种行为?
这是我的代码:
缓冲类:
#include "StringBuffer.h"
void StringBuffer::put(string str)
{
scoped_lock lock(mutex);
if (full == BUF_SIZE)
{
{
//boost::mutex::scoped_lock lock(io_mutex);
//std::cout << "Buffer is full. Waiting..." << std::endl;
}
while (full == BUF_SIZE)
cond.wait(lock);
}
str_buffer[p] = str;
p = (p+1) % BUF_SIZE;
++full;
cond.notify_one();
}
string StringBuffer::get()
{
scoped_lock lk(mutex);
if (full == 0)
{
{
//boost::mutex::scoped_lock lock(io_mutex);
//std::cout << "Buffer is empty. Waiting..." << std::endl;
}
while (full == 0)
cond.wait(lk);
}
string test = str_buffer[c];
c = (c+1) % BUF_SIZE;
--full;
cond.notify_one();
return test;
}
这是主要的:
Parser p;
StringBuffer buf;
Report report;
string transfer;
ifstream input;
vector <boost::regex> regs;
int proc_count = 0;
int push_count = 0;
bool pusher_done = false;
// Show filter configuration and init report by dimensioning counter vectors
void setup_report() {
for (int k = 0; k < p.filters(); k++) {
std::cout << "SID(NUM):" << k << " Name(TXT):\"" << p.name_at(k) << "\"" << " Filter(REG):\"" << p.filter_at(k) << "\"" << endl;
regs.push_back(boost::regex(p.filter_at(k)));
report.hits_filters.push_back(0);
report.names.push_back(p.name_at(k));
report.filters.push_back(p.filter_at(k));
}
}
// Read strings from sourcefiles and put them into buffer
void pusher() {
// as long as another string could be red, ...
while (input) {
// put it into buffer
buf.put(transfer);
// and get another string from source file
getline(input, transfer);
push_count++;
}
pusher_done = true;
}
// Get strings from buffer and check RegEx filters. Pass matches to report
void processor()
{
while (!pusher_done || buf.get_rest()) {
string n = buf.get();
for (unsigned sid = 0; sid < regs.size(); sid++) {
if (boost::regex_search(n, regs[sid])) report.report_hit(sid);
}
boost::mutex::scoped_lock lk(buf.count_mutex);
{
proc_count++;
}
}
}
int main(int argc, const char* argv[], char* envp[])
{
if (argc == 3)
{
// first add sourcefile from argv[1] filepath, ...
p.addSource(argv[1]);
std::cout << "Source File: *** Ok\n";
// then read configuration from argv[2] filepath, ...
p.readPipes(envp, argv[2]);
std::cout << "Configuration: *** Ok\n\n";
// and setup the Report Object.
setup_report();
// For all sourcefiles that have been parsed, ...
for (int i = 0; i < p.sources(); i++) {
input.close();
input.clear();
// open the sourcefile in a filestream.
input.open(p.source_at(i).c_str());
// check if file exist, otherwise throw error and exit
if (!input)
{
std::cout << "\nError! File not found: " << p.source_at(i);
exit(1);
}
// get start time
std::cout << "\n- started: ";
ptime start(second_clock::local_time());
cout << start << endl;
// read a first string into transfer to get the loops going
getline(input, transfer);
// create threads and pass a reference to functions
boost::thread push1(&pusher);
boost::thread proc1(&processor);
boost::thread proc2(&processor);
// start all the threads and wait for them to complete.
push1.join();
proc1.join();
proc2.join();
// calculate and output runtime and lines per second
ptime end(second_clock::local_time());
time_duration runtime = end - start;
std::cout << "- finished: " << ptime(second_clock::local_time()) << endl;
cout << "- processed lines: " << push_count << endl;
cout << "- runtime: " << to_simple_string(runtime) << endl;
float processed = push_count;
float lines_per_second = processed/runtime.total_seconds();
cout << "- lines per second: " << lines_per_second << endl;
// write report to file
report.create_filereport(); // after all threads finished write reported data to file
cout << "\nReport saved as: ./report.log\n\nBye!" << endl;
}
}
else std::cout << "Usage: \"./Speed-Extract [source][config]\"\n\n";
return 0;
}
编辑 1:
非常感谢您的帮助。通过在输出中添加一些计数器和线程 ID,我发现了问题所在:
我注意到有几个线程可能仍在等待缓冲区填充。
我的处理线程在有尚未读取的新源字符串或缓冲区不为空时运行。这不好。
假设我有 2 个线程等待缓冲区填充。一旦读者读到新行(可能是最后几行),就会有 6 个其他线程尝试获取此行并锁定项目,因此 2 个等待线程可能甚至没有机会尝试解锁它.
一旦他们检查到一行被另一个线程占用,他们就会继续等待。 Reading 线程在到达 eof 然后停止时不会通知它们。两个等待的线程永远等待。
我的 Reading 函数还必须通知所有线程它已到达 eof,因此只有当缓冲区为空且文件不是 EOF 时,线程才应保持等待。
【问题讨论】:
-
我们需要查看导致问题的代码,尤其是锁定问题。
-
现在我们需要所有可以编译和运行代码的东西。所以标题和适当的数据。
-
将其归结为一个完整的实体,我们可以编译和重现行为。
标签: c++ multithreading boost synchronization mutex