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

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
}

关键修改说明

  1. 新增写入通道与专属goroutine:所有需要持久化的数据通过writeChan传递,由单一goroutine负责文件写入,彻底避免并发冲突。
  2. 移除直接文件操作:删除kcdLoop3和Main中直接调用f.Write的代码,统一改为向通道发送数据。
  3. 完善错误与资源处理:添加文件打开、JSON序列化、响应体关闭的错误检查与资源释放,避免内存泄漏。
  4. 使用WaitGroup同步:通过sync.WaitGroup等待所有爬取goroutine完成,再关闭写入通道,确保所有数据都被持久化。

内容的提问来源于stack exchange,提问作者Hyeonkyu Han

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 23:15:35