You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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;
}

核心改进点

  1. 队列缓存数据:用std::queue<LineData>存储每行的时间戳和原始内容,避免覆盖未处理的数据,确保所有行都能被处理。
  2. 修正同步逻辑:读取线程写入队列后通知处理线程;处理线程等待队列有数据或文件读完,每次取时间戳最小的元素加入结果。
  3. 补充缺失功能:实现了split分割函数和完整的thread1读取逻辑。
  4. 处理剩余数据:当其中一个文件读完后,自动处理另一个队列中的所有剩余数据。
  5. 结果验证:输出所有行并验证排序状态。

内容的提问来源于stack exchange,提问作者Remí1998

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.01 23:42:51