基于OpenMP的C++ LZW数据压缩并行化问题求助
问题分析与解决方案
核心错误原因
你的代码出现访问违例和逻辑错误,根源有三个:
- 共享变量
current的数据竞争:current是所有线程共享的变量,循环中current = next操作没有同步保护,多个线程会同时读写current,导致next的计算完全混乱,甚至生成无效字符串引发内存访问错误。 unordered_map非线程安全:即使在临界区修改字典,临界区外的dict.find(next)是并发读操作。unordered_map在插入元素时可能触发rehash,此时所有迭代器失效,其他线程的find会触发未定义行为,包括访问违例。- 错误的并行模型:LZW是强顺序依赖的算法,每个字符的处理结果依赖前一个字符的
current状态,直接用parallel for遍历单个字符,完全违背了算法的核心逻辑,从根本上就不成立。
正确并行化思路
LZW无法直接并行单个字符的处理,必须采用分块并行+边界同步的策略:
- 将输入分成多个独立块,每个线程用局部字典处理块内内容,避免并发修改全局字典的问题。
- 收集每个块末尾的
current字符串,合并到全局字典,处理跨块的字符串匹配。 - 重新处理每个块的边界部分,确保使用统一的全局字典,保证压缩结果与顺序版本一致。
修复后的并行实现代码
#include <omp.h> #include <string> #include <unordered_map> #include <vector> #include <fstream> #include <iostream> #include <iterator> #include <algorithm> using namespace std; // 存储每个块的处理结果:压缩数据、块末尾的current字符串 struct BlockResult { vector<int> compressed; string trailing_current; }; // 单块处理逻辑,使用局部字典 BlockResult process_block(const string& block, const unordered_map<string, int>& global_dict, int& local_dict_size) { BlockResult result; unordered_map<string, int> dict = global_dict; int dict_size = local_dict_size; string current; for (char c : block) { string next = current + c; if (dict.count(next)) { current = next; } else { result.compressed.push_back(dict[current]); dict[next] = dict_size++; current = string(1, c); } } result.trailing_current = current; local_dict_size = dict_size; return result; } void lzw_parallel(const string& input, const string& output_file) { // 初始化全局基础字典(单字符) unordered_map<string, int> global_dict; for (int i = 0; i < 256; i++) { global_dict[string(1, i)] = i; } int global_dict_size = 256; int num_threads = omp_get_max_threads(); int block_size = input.length() / num_threads; vector<BlockResult> results(num_threads); double start_time = omp_get_wtime(); // 第一步:并行处理每个块,使用局部字典,记录边界状态 #pragma omp parallel for shared(input, global_dict, results, global_dict_size) for (int t = 0; t < num_threads; t++) { int start = t * block_size; int end = (t == num_threads - 1) ? input.length() : (t + 1) * block_size; string block(input.begin() + start, input.begin() + end); int local_dict_size = global_dict_size; results[t] = process_block(block, global_dict, local_dict_size); // 同步更新全局字典的最大尺寸 #pragma omp critical { if (local_dict_size > global_dict_size) { global_dict_size = local_dict_size; } } } // 第二步:合并跨块的字符串到全局字典 for (size_t t = 0; t < num_threads - 1; t++) { const string& prev_trailing = results[t].trailing_current; if (prev_trailing.empty()) continue; // 取当前块的第一个字符,拼接前一块的末尾字符串 char next_char = input[(t + 1) * block_size]; string next = prev_trailing + next_char; if (global_dict.find(next) == global_dict.end()) { global_dict[next] = global_dict_size++; } } // 第三步:重新处理每个块的边界,确保使用统一的全局字典 #pragma omp parallel for shared(input, global_dict, results, global_dict_size) for (int t = 1; t < num_threads; t++) { int start = t * block_size; int end = (t == num_threads - 1) ? input.length() : (t + 1) * block_size; const string& prev_trailing = results[t-1].trailing_current; if (prev_trailing.empty()) continue; unordered_map<string, int> dict = global_dict; int dict_size = global_dict_size; vector<int> new_compressed; string current = prev_trailing; char c = input[start]; string next = current + c; // 处理边界第一个字符 if (dict.count(next)) { current = next; } else { new_compressed.push_back(dict[current]); dict[next] = dict_size++; current = string(1, c); } // 处理块内剩余字符 for (int i = start + 1; i < end; i++) { c = input[i]; next = current + c; if (dict.count(next)) { current = next; } else { new_compressed.push_back(dict[current]); dict[next] = dict_size++; current = string(1, c); } } results[t].compressed = new_compressed; results[t].trailing_current = current; } // 合并所有块的压缩结果 vector<int> compressed; for (auto& res : results) { compressed.insert(compressed.end(), res.compressed.begin(), res.compressed.end()); } // 处理最后一个块的剩余current if (!results.back().trailing_current.empty()) { compressed.push_back(global_dict[results.back().trailing_current]); } double end_time = omp_get_wtime(); // 写入压缩文件 ofstream outfile(output_file); for (int i : compressed) { outfile << i << " "; } // 输出统计信息 cout << "Input length: " << input.length() << endl; cout << "Compressed length: " << compressed.size() * 12 / 8 << endl; printf("Work took %f seconds \n", end_time - start_time); }
额外优化建议
- 如果使用Intel TBB库,可以用
tbb::concurrent_unordered_map替代unordered_map,简化全局字典的并发访问,但仍需解决顺序依赖问题。 - 分块大小可以根据输入长度动态调整,避免过小的块导致并行开销大于收益。
- 若对压缩率要求不高,可以让每个线程使用独立的局部字典,最后将局部ID映射到全局ID,这种实现更简单,但会降低压缩率。
内容的提问来源于stack exchange,提问作者Dababy
相关产品推荐
相关产品推荐

