【问题标题】:Webcrawler in GoGo 中的网络爬虫
【发布时间】:2015-04-07 12:37:47
【问题描述】:

我正在尝试在 Go 中构建一个网络爬虫,我想在其中指定并发工作人员的最大数量。只要队列中有要探索的链接,它们都会工作。当队列的元素少于workers时,workers应该向下喊,但如果发现更多链接,则恢复。

我试过的代码是

const max_workers = 6
// simulating links with int
func crawl(wg *sync.WaitGroup, queue chan int) {
    for element := range queue {   
        wg.Done() // why is defer here causing a deadlock?
        fmt.Println("adding 2 new elements ")
        if element%2 == 0 {
            wg.Add(2)
            queue <- (element*100 + 11)
            queue <- (element*100 + 33)
        }

    }
}

func main() {
    var wg sync.WaitGroup
    queue := make(chan int, 10)
    queue <- 0
    queue <- 1
    queue <- 2
    queue <- 3
    var min int
    if (len(queue) < max_workers) {
        min = len(queue)
    } else {
        min = max_workers
    }
    for i := 0; i < min; i++ {
        wg.Add(1)
        go crawl(&wg, queue)
    }
    wg.Wait()
    close(queue)
}

Link to playground

这似乎可行,但有一个问题:当我开始时,我必须用多个元素填充队列。我希望它从(单个)种子页面(在我的示例中为queue &lt;- 0)开始,然后动态地增大/缩小工作池。

我的问题是:

  • 如何获取行为?

  • 为什么 defer wg.Done() 会导致死锁? wg.Done()函数实际完成的时候正常吗?我认为如果没有defer,goroutine 不会等待另一部分完成(在解析 HTML 的实际工作示例中这可能需要更长的时间)。

【问题讨论】:

标签: go web-crawler


【解决方案1】:

如果您使用自己喜欢的网络搜索“Go web crawler”(或“golang web crawler”) 你会发现很多例子,包括: Go Tour Exercise: Web Crawler。 Go 中也有一些关于并发的讨论,涵盖了这类事情。

在 Go 中执行此操作的“标准”方式根本不需要涉及等待组。 要回答您的一个问题,使用defer 排队的事情只有在函数返回时才会运行。你有一个长时间运行的函数,所以不要在这样的循环中使用defer

“标准”方式是在他们自己的 goroutine 中启动任意数量的工人。 他们都从同一个频道读取“作业”,如果/当无事可做时会阻塞。 完全完成后,该通道将关闭,它们都将退出。

在爬虫之类的情况下,工作人员会发现更多“工作”要做,并希望将它们排入队列。 你不希望他们写回同一个频道,因为它会有一些有限的缓冲(或者没有!),你最终会阻止所有试图加入更多工作的工人!

一个简单的解决方案是使用单独的频道 (例如,每个工人都有in &lt;-chan Job, out chan&lt;- Job) 以及读取这些请求的单个队列/过滤器 goroutine, 将它们附加到一个切片上,它要么让其增长任意大,要么对其进行一些全局限制, 并且还从切片的头部馈送另一个通道 (即从一个通道读取并写入另一个通道的简单 for-select 循环)。 该代码通常还负责跟踪已经完成的工作 (例如访问的 URL 地图)并丢弃传入的重复请求。

队列 goroutine 可能看起来像这样(这里的参数名称过于冗长):

type Job string

func queue(toWorkers chan<- Job, fromWorkers <-chan Job) {
    var list []Job
    done := make(map[Job]bool)
    for {
        var send chan<- Job
        var item Job
        if len(list) > 0 {
            send = toWorkers
            item = list[0]
        }
        select {
        case send <- item:
            // We sent an item, remove it
            list = list[1:]
        case thing := <-fromWorkers:
            // Got a new thing
            if !done[thing] {
                list = append(list, thing)
                done[thing] = true
            }
        }
    }
}

在这个简单的例子中,有几件事被掩盖了。 比如终止。如果“工作”是一些更大的结构,您想使用 chan *Job[]*Job 代替。 在这种情况下,您还需要将地图类型更改为您从作业中提取的一些键 (例如,Job.URL 也许) 并且您想在 list = list[1:] 之前执行 list[0] = nil 以摆脱对 *Job 指针的引用,并让垃圾收集器更早地处理它。

编辑:关于干净终止的一些说明。

有几种方法可以干净地终止上述代码。可以使用等待组,但添加/完成调用的放置需要仔细完成,并且您可能需要另一个 goroutine 来执行等待(然后关闭其中一个通道以开始关闭)。工作人员不应该关闭他们的输出通道,因为有多个工作人员并且您不能多次关闭一个通道;队列 goroutine 无法告诉何时关闭它与工人的通道,而不知道工人何时完成。

过去,当我使用与上述非常相似的代码时,我在“队列”goroutine 中使用了本地“杰出”计数器(这避免了对互斥体的任何需求或等待组所具有的任何同步开销)。将作业发送给工作人员时,未完成作业的计数会增加。当工人说它已经完成时,它会再次减少。我的代码碰巧有另一个渠道(我的“队列”除了要排队的更多节点外,还收集结果)。它在自己的频道上可能更干净,但可以使用现有频道上的特殊值(例如 nil 作业指针)。无论如何,有了这样一个计数器,本地列表上的现有长度检查只需要在列表为空并且该终止时看到没有任何未完成的事情;只需关闭通往工人的频道并返回即可。

例如:

    if len(list) > 0 {
        send = toWorkers
        item = list[0]
    } else if outstandingJobs == 0 {
        close(toWorkers)
        return
    }

【讨论】:

  • 谢谢,这个答案真的很有帮助!我读过很多关于爬虫的东西,但我不得不承认并发对我来说仍然很难。此外,当阅读其他人的代码时,一切似乎都更容易了。这就是为什么我决定用我自己的代码询问
  • 能否请您详细说明一下垃圾收集部分?我不明白为什么 list[0] = nil 应该有所帮助。
  • @meto 只有当list[]*Job 时才需要它。 list = list[1:] 重新切片以“删除”条目,但切片指向的底层数组不会改变。直到append 最终 命中len(list) == cap(list) 并分配一个新的底层数组(仅复制旧切片中剩余的项目),旧数组才变为未引用,进而它已成为任何指针未引用。 list[0] = nil 立即删除该引用,并且一旦它被发送到的工作人员也停止引用它,垃圾收集器就可以在不等待追加的情况下收集它。
  • 好的,如果可以的话,最后一个问题(但如果您愿意,我可以打开另一个答案)。决定何时停止并非易事。仅检查通道是否为空可能还不够(其中一名工作人员可能正忙于解析 HTML)。有没有办法检查所有工人是否都在等待?在这种情况下,可以通过bool 频道发送信号。
【解决方案2】:

我利用 Go 的互斥 (Mutex) 功能编写了一个解决方案。

当它在并发上运行时,一次只限制一个实例访问 url 映射可能很重要。我相信我按照下面的方式实现了它。请随意尝试一下。非常感谢您的反馈,因为我也会向您的 cmets 学习。

package main

import (
    "fmt"
    "sync"
)

type Fetcher interface {
    // Fetch returns the body of URL and
    // a slice of URLs found on that page.
    Fetch(url string) (body string, urls []string, err error)
}




// ! SafeUrlBook helps restrict only one instance access the central url map at a time. So that no redundant crawling should occur.
type SafeUrlBook struct {
    book map[string]bool
    mux  sync.Mutex
    }

func (sub *SafeUrlBook) doesThisExist(url string) bool {
    sub.mux.Lock()
    _ , key_exists := sub.book[url]
    defer sub.mux.Unlock()
    
    if key_exists {
    return true
    }  else { 
    sub.book[url] = true
    return false 
    }  
}
// End SafeUrlBook


// Crawl uses fetcher to recursively crawl
// pages starting with url, to a maximum of depth.
// Note that now I use safeBook (SafeUrlBook) to keep track of which url has been visited by a crawler.
func Crawl(url string, depth int, fetcher Fetcher, safeBook SafeUrlBook) {
    if depth <= 0 {
        return
    }
    
    
    exist := safeBook.doesThisExist(url)
    if exist { fmt.Println("Skip", url) ; return }
    
    
    body, urls, err := fetcher.Fetch(url)
    if err != nil {
        fmt.Println(err)
        return
    }
    fmt.Printf("found: %s %q\n", url, body)
    for _, u := range urls {
        Crawl(u, depth-1, fetcher, safeBook)
    }
    return
}

func main() {
    safeBook := SafeUrlBook{book: make(map[string]bool)}
    Crawl("https://golang.org/", 4, fetcher, safeBook)
}

// fakeFetcher is Fetcher that returns canned results.
type fakeFetcher map[string]*fakeResult

type fakeResult struct {
    body string
    urls []string
}

func (f fakeFetcher) Fetch(url string) (string, []string, error) {
    if res, ok := f[url]; ok {
        return res.body, res.urls, nil
    }
    return "", nil, fmt.Errorf("not found: %s", url)
}

// fetcher is a populated fakeFetcher.
var fetcher = fakeFetcher{
    "https://golang.org/": &fakeResult{
        "The Go Programming Language",
        []string{
            "https://golang.org/pkg/",
            "https://golang.org/cmd/",
        },
    },
    "https://golang.org/pkg/": &fakeResult{
        "Packages",
        []string{
            "https://golang.org/",
            "https://golang.org/cmd/",
            "https://golang.org/pkg/fmt/",
            "https://golang.org/pkg/os/",
        },
    },
    "https://golang.org/pkg/fmt/": &fakeResult{
        "Package fmt",
        []string{
            "https://golang.org/",
            "https://golang.org/pkg/",
        },
    },
    "https://golang.org/pkg/os/": &fakeResult{
        "Package os",
        []string{
            "https://golang.org/",
            "https://golang.org/pkg/",
        },
    },
}

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2011-12-11
    • 1970-01-01
    • 1970-01-01
    • 2018-08-08
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多