C++三线程按时间戳合并排序两个CSV文件的实现及问题排查
多线程合并两个CSV文件并按时间戳排序的问题解决
我有两个CSV文件A、B,每行数据格式为「时间戳, 参数...」。希望用一个线程读取其中一个文件,另一个线程读取另一个文件,由第三个线程比较时间戳,最终构造一个包含两个文件中所有行内容、按时间戳排序的vector。但以下实现无法输出所有时间戳,求解决方法:
#include <iostream> #include <fstream> #include <string> #include <mutex> #include <condition_variable> #include <thread> #include <vector> #include <chrono> std::mutex mtx; std::condition_variable cv; long long timestamp1, timestamp2; std::vector<long long> timestamps; bool finished1 = false, finished2 = false; void thread2() { std::ifstream file2("a.csv"); std::string line; while (std::getline(file2, line)) { std::vector<std::string> values = split(line, ','); long long current_timestamp = std::stoll(values[4]); { std::unique_lock<std::mutex> lock(mtx); while (timestamp1 >= current_timestamp) { cv.wait(lock); } timestamp2 = current_timestamp; } cv.notify_one(); } { std::unique_lock<std::mutex> lock(mtx); finished2 = true; } cv.notify_one(); } void thread3() { while (!finished1 || !finished2) { std::unique_lock<std::mutex> lock(mtx); cv.wait(lock); if (finished1 && finished2) { break; } if (timestamp1 >= timestamp2) { timestamps.push_back(timestamp1); std::cout << timestamp1 <<"\n" << std::flush; } else { timestamps.push_back(timestamp2); std::cout << timestamp2 <<"\n" << std::flush; } } } #include <algorithm> int main() { std::thread t1(thread1); std::thread t2(thread2); std::thread t3(thread3); t1.join(); t2.join(); t3.join(); std::cout << std::is_sorted(timestamps.begin(),timestamps.end()); }
问题分析
现有代码存在多个关键缺陷,导致数据丢失或逻辑错误:
- 缺少
thread1实现:main函数启动了t1(thread1)但未定义该函数,直接导致编译失败。 - 共享变量设计错误:仅用单个
timestamp1和timestamp2存储当前读取的时间戳,会覆盖未处理的数据,直接丢失部分时间戳。 - 条件变量逻辑混乱:读取线程的等待条件
timestamp1 >= current_timestamp无合理依据,会导致线程阻塞或处理顺序完全错误。 - 处理线程逻辑缺陷:thread3每次唤醒后仅处理一个时间戳,但未标记已处理状态,且未处理其中一个文件读完后剩余的所有数据。
- 缺少
split函数实现:代码调用了split但未定义,编译无法通过。
修正后的完整实现
#include <iostream> #include <fstream> #include <string> #include <mutex> #include <condition_variable> #include <thread> #include <vector> #include <queue> #include <sstream> #include <algorithm> // 存储每行的时间戳和原始内容 struct LineData { long long timestamp; std::string line; }; std::mutex mtx; std::condition_variable cv; std::queue<LineData> queue1; // 文件A的缓存队列 std::queue<LineData> queue2; // 文件B的缓存队列 bool finished1 = false, finished2 = false; std::vector<std::string> result; // 最终排序后的所有行内容 // 字符串分割工具函数 std::vector<std::string> split(const std::string& s, char delimiter) { std::vector<std::string> tokens; std::string token; std::istringstream tokenStream(s); while (std::getline(tokenStream, token, delimiter)) { tokens.push_back(token); } return tokens; } // 读取文件A的线程 void thread1() { std::ifstream file("a.csv"); std::string line; while (std::getline(file, line)) { std::vector<std::string> values = split(line, ','); if (values.size() < 1) continue; // 跳过无效行 // 假设时间戳是第一列,若实际是第5列则改为values[4] long long ts = std::stoll(values[0]); { std::lock_guard<std::mutex> lock(mtx); queue1.push({ts, line}); } cv.notify_one(); } { std::lock_guard<std::mutex> lock(mtx); finished1 = true; } cv.notify_one(); } // 读取文件B的线程 void thread2() { std::ifstream file("b.csv"); std::string line; while (std::getline(file, line)) { std::vector<std::string> values = split(line, ','); if (values.size() < 1) continue; long long ts = std::stoll(values[0]); { std::lock_guard<std::mutex> lock(mtx); queue2.push({ts, line}); } cv.notify_one(); } { std::lock_guard<std::mutex> lock(mtx); finished2 = true; } cv.notify_one(); } // 合并排序的线程 void thread3() { while (true) { std::unique_lock<std::mutex> lock(mtx); // 等待有数据或者两个文件都读取完毕 cv.wait(lock, []{ return (!queue1.empty() || !queue2.empty() || finished1 || finished2); }); bool has_data = false; // 每次取时间戳较小的元素加入结果 if (!queue1.empty() && !queue2.empty()) { if (queue1.front().timestamp <= queue2.front().timestamp) { result.push_back(queue1.front().line); queue1.pop(); } else { result.push_back(queue2.front().line); queue2.pop(); } has_data = true; } else if (!queue1.empty()) { result.push_back(queue1.front().line); queue1.pop(); has_data = true; } else if (!queue2.empty()) { result.push_back(queue2.front().line); queue2.pop(); has_data = true; } // 两个文件都读完且队列清空,退出循环 if (finished1 && finished2 && queue1.empty() && queue2.empty()) { break; } if (has_data) { cv.notify_all(); } } } int main() { std::thread t1(thread1); std::thread t2(thread2); std::thread t3(thread3); t1.join(); t2.join(); t3.join(); // 验证结果是否按时间戳排序并输出 bool sorted = true; long long prev_ts = -1; for (const auto& line : result) { std::vector<std::string> values = split(line, ','); long long ts = std::stoll(values[0]); if (ts < prev_ts) { sorted = false; } prev_ts = ts; std::cout << line << std::endl; } std::cout << "\n结果是否已排序:" << (sorted ? "是" : "否") << std::endl; return 0; }
核心改进点
- 队列缓存数据:用
std::queue<LineData>存储每行的时间戳和原始内容,避免覆盖未处理的数据,确保所有行都能被处理。 - 修正同步逻辑:读取线程写入队列后通知处理线程;处理线程等待队列有数据或文件读完,每次取时间戳最小的元素加入结果。
- 补充缺失功能:实现了
split分割函数和完整的thread1读取逻辑。 - 处理剩余数据:当其中一个文件读完后,自动处理另一个队列中的所有剩余数据。
- 结果验证:输出所有行并验证排序状态。
内容的提问来源于stack exchange,提问作者Remí1998
相关产品推荐
相关产品推荐

