如何使用cilk_reducer无竞争写入文件及实现identity、reduce函数
使用OpenCilk的cilk_reducer实现无数据竞争的文件写入
核心思路
OpenCilk的reducer通过为每个工作线程创建独立的视图(view),避免直接共享文件句柄或直接写文件,而是让线程先将输出内容暂存到各自的视图缓冲区,最后通过reduce函数合并所有线程的内容,再一次性写入文件,从根源上消除数据竞争。
步骤1:定义自定义Reducer的视图结构
首先定义每个线程独立持有的视图,用缓冲区暂存该线程的输出内容:
#include <cilk/cilk.h> #include <cilk/reducer.h> #include <stdio.h> #include <stdlib.h> #include <string.h> #include <stdarg.h> // 每个线程的视图:独立的输出缓冲区 typedef struct { char* buffer; size_t length; // 当前已写入内容长度 size_t capacity; // 缓冲区总容量 } FileWriteView;
步骤2:实现identity函数
identity函数负责初始化每个线程的视图,创建一个空的可扩展缓冲区:
void identity(void* view) { FileWriteView* v = (FileWriteView*)view; v->capacity = 256; // 初始缓冲区大小 v->buffer = (char*)malloc(v->capacity); v->length = 0; v->buffer[v->length] = '\0'; // 确保字符串终止符 }
这个函数会为每个新启动的工作线程生成一个干净的初始视图,保证线程间的缓冲区完全隔离。
步骤3:实现reduce函数
reduce函数负责合并两个视图的内容,将右侧视图的缓冲区内容追加到左侧,之后释放右侧的缓冲区避免内存泄漏:
void reduce(void* left, void* right) { FileWriteView* l = (FileWriteView*)left; FileWriteView* r = (FileWriteView*)right; // 左侧缓冲区不足时扩容 if (l->length + r->length + 1 > l->capacity) { l->capacity = l->length + r->length + 256; // 预留额外空间 l->buffer = (char*)realloc(l->buffer, l->capacity); } // 追加右侧内容到左侧 strncat(l->buffer, r->buffer, r->length); l->length += r->length; // 清理右侧视图资源 free(r->buffer); r->buffer = NULL; r->length = 0; r->capacity = 0; }
如果需要保证写入顺序(比如按并行迭代的顺序输出),可以在视图中增加序号字段,合并时按序号排序后再拼接内容,否则默认是任意顺序的合并。
步骤4:注册并使用自定义Reducer
用CILK_CUSTOM_REDUCER宏注册我们的reducer类型,再封装辅助函数简化写入操作,最后在并行区域中使用:
// 注册自定义reducer类型 CILK_CUSTOM_REDUCER(FileWriteReducer, FileWriteView, identity, reduce); // 辅助函数:向reducer视图写入格式化内容 void reducer_printf(CILK_REDUCER(FileWriteReducer)* r, const char* format, ...) { va_list args; va_start(args, format); FileWriteView* v = CILK_REDUCER_VIEW(r); char temp_buf[1024]; int written = vsnprintf(temp_buf, sizeof(temp_buf), format, args); // 扩容视图缓冲区 if (v->length + written + 1 > v->capacity) { v->capacity = v->length + written + 256; v->buffer = (char*)realloc(v->buffer, v->capacity); } // 将临时内容追加到视图缓冲区 strncat(v->buffer, temp_buf, written); v->length += written; va_end(args); } int main() { // 初始化自定义reducer CILK_REDUCER(FileWriteReducer) reducer; CILK_REDUCER_INIT(reducer); // 并行任务:每个线程向自己的视图写入内容 cilk_for (int i = 0; i < 10; ++i) { reducer_printf(&reducer, "线程%d:并行迭代%d的输出内容\n", cilk_spawn_id(), i); } // 合并所有视图后,一次性写入文件 FileWriteView* final_view = CILK_REDUCER_VIEW(&reducer); FILE* fp = fopen("output.txt", "w"); if (fp != NULL) { fwrite(final_view->buffer, 1, final_view->length, fp); fclose(fp); } // 清理reducer资源 CILK_REDUCER_DESTROY(reducer); return 0; }
关键注意事项
- 禁止并行直接操作文件:所有线程的写入都先缓存到独立视图,最后统一写入,彻底避免文件IO的竞争。
- 内存管理:严格在
identity和reduce中管理缓冲区的分配与释放,防止内存泄漏。 - 顺序控制:如果需要严格的输出顺序,需在视图中添加序号标记,合并时按序号排序后再拼接内容。
内容的提问来源于stack exchange,提问作者Mohan S
相关产品推荐
相关产品推荐

