Golang爬虫并发写入JSONL文件异常的解决方法
解决Golang爬虫并发写入JSONL文件的异常问题
你遇到的空行、内容连续写入异常,核心原因是多个goroutine并发操作同一个文件句柄导致的竞态条件。os.File的Write方法并非并发安全,当多个goroutine同时调用Write时,字节流会被交错写入,出现内容重叠、空行等问题。另外你的代码里kcdLoop3和Main函数中都直接操作文件,两处并发写入进一步加剧了冲突。
解决方案:单一写入goroutine + 通道控制
最可靠的修复方式是用专属goroutine负责文件写入,其他爬取逻辑的goroutine只把数据发送到通道,由写入goroutine统一处理,彻底避免并发写入冲突。同时完善错误处理,避免资源泄漏。
修改后的完整代码
package next import ( "encoding/json" "fmt" "io/fs" "net/http" "net/url" "os" "strconv" "sync" ) type IApi struct { Tree []ITree `json:"tree"` } type ITree struct { CodeTitle string `json:"code_title"` Degree string `json:"degree"` EngText string `json:"eng_text"` GubunCode string `json:"gubun_code"` KorText string `json:"kor_text"` LastNodeAt string `json:"last_node_at"` LevelNo int `json:"level_no"` LevelTitle string `json:"level_title"` } func Main() { KCD_numbers := []int{5, 6, 7, 8} c := make(chan IApi) c1 := make(chan IApi) c2 := make(chan IApi) c3 := make(chan IApi) // 新增写入专用通道,传递需要持久化的ITree对象 writeChan := make(chan ITree) var wg sync.WaitGroup // 启动独立写入goroutine,唯一操作文件句柄 f, err := os.OpenFile("./koicd_code_list.json", os.O_CREATE|os.O_RDWR|os.O_TRUNC, fs.FileMode(0644)) if err != nil { fmt.Printf("打开文件失败: %s\n", err) os.Exit(1) } defer f.Close() go func() { for item := range writeChan { data, err := json.Marshal(item) if err != nil { fmt.Printf("JSON序列化失败: %s\n", err) continue } // 写入JSON数据+换行,保证每行一个有效JSON if _, err := f.Write(append(data, '\n')); err != nil { fmt.Printf("写入文件失败: %s\n", err) } } }() // 爬取KCD一级数据 for _, KCD_NO := range KCD_numbers { wg.Add(1) go func(no int) { defer wg.Done() kcdLoop(no, c) }(KCD_NO) } // 处理各级数据 for _, KCD_NO := range KCD_numbers { api_url := fmt.Sprintf("https://www.koicd.kr/kcd/0%s/kcd.tree.json", strconv.Itoa(KCD_NO)) k1_list := (<-c).Tree for _, k1 := range k1_list { wg.Add(1) go func(node ITree) { defer wg.Done() kcdLoop1(node, c1, api_url) }(k1) } for i := 0; i < len(k1_list); i++ { k2_list := (<-c1).Tree for _, k2 := range k2_list { wg.Add(1) go func(node ITree) { defer wg.Done() kcdLoop2(node, c2, api_url) }(k2) } for j := 0; j < len(k2_list); j++ { k3_list := (<-c2).Tree for _, k3 := range k3_list { // 将k3数据发送到写入通道,不再直接操作文件 writeChan <- k3 wg.Add(1) go func(node ITree) { defer wg.Done() kcdLoop3(node, c3, api_url) }(k3) } for k := 0; k < len(k3_list); k++ { k4_list := (<-c3).Tree for _, k4 := range k4_list { // 将k4数据发送到写入通道 writeChan <- k4 } } } } } // 等待所有爬取goroutine完成,再关闭写入通道 wg.Wait() close(writeChan) } func kcdLoop(KCD_NO int, c chan IApi) { api_url := fmt.Sprintf("https://www.koicd.kr/kcd/0%s/kcd.tree.json", strconv.Itoa(KCD_NO)) resp, herr := http.PostForm(api_url, url.Values{}) if herr != nil { fmt.Printf("请求失败: %s\n", herr) os.Exit(1) } defer resp.Body.Close() // 释放响应体资源 var res IApi if err := json.NewDecoder(resp.Body).Decode(&res); err != nil { fmt.Printf("JSON解析失败: %s\n", err) return } c <- res } func kcdLoop1(k1 ITree, c1 chan IApi, api_url string) { k1_resp, k1_herr := http.PostForm(api_url, url.Values{"level_no": {strconv.Itoa(k1.LevelNo)}, "level_title": {k1.LevelTitle}}) if k1_herr != nil { fmt.Printf("请求失败: %s\n", k1_herr) os.Exit(1) } defer k1_resp.Body.Close() var k1_res IApi if err := json.NewDecoder(k1_resp.Body).Decode(&k1_res); err != nil { fmt.Printf("JSON解析失败: %s\n", err) return } c1 <- k1_res } func kcdLoop2(k2 ITree, c2 chan IApi, api_url string) { k2_resp, k2_herr := http.PostForm(api_url, url.Values{"level_no": {strconv.Itoa(k2.LevelNo)}, "level_title": {k2.LevelTitle}}) if k2_herr != nil { fmt.Printf("请求失败: %s\n", k2_herr) os.Exit(1) } defer k2_resp.Body.Close() var k2_res IApi if err := json.NewDecoder(k2_resp.Body).Decode(&k2_res); err != nil { fmt.Printf("JSON解析失败: %s\n", err) return } c2 <- k2_res } func kcdLoop3(k3 ITree, c3 chan IApi, api_url string) { k3_resp, k3_herr := http.PostForm(api_url, url.Values{"level_no": {strconv.Itoa(k3.LevelNo)}, "level_title": {k3.LevelTitle}}) if k3_herr != nil { fmt.Printf("请求失败: %s\n", k3_herr) os.Exit(1) } defer k3_resp.Body.Close() var k3_res IApi if err := json.NewDecoder(k3_resp.Body).Decode(&k3_res); err != nil { fmt.Printf("JSON解析失败: %s\n", err) return } c3 <- k3_res }
关键修改说明
- 新增写入通道与专属goroutine:所有需要持久化的数据通过
writeChan传递,由单一goroutine负责文件写入,彻底避免并发冲突。 - 移除直接文件操作:删除
kcdLoop3和Main中直接调用f.Write的代码,统一改为向通道发送数据。 - 完善错误与资源处理:添加文件打开、JSON序列化、响应体关闭的错误检查与资源释放,避免内存泄漏。
- 使用WaitGroup同步:通过
sync.WaitGroup等待所有爬取goroutine完成,再关闭写入通道,确保所有数据都被持久化。
内容的提问来源于stack exchange,提问作者Hyeonkyu Han
相关产品推荐
相关产品推荐

