Created
September 26, 2012 08:47
-
-
Save reusee/3786868 to your computer and use it in GitHub Desktop.
golang实现的http get批量采集器
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| package main | |
| import ( | |
| "bufio" | |
| "flag" | |
| "fmt" | |
| "io/ioutil" | |
| "net" | |
| "net/http" | |
| "reus/utils" | |
| "strings" | |
| "time" | |
| ) | |
| var poolSize = flag.Int("j", 100, "number of concurrent threads") | |
| func collect(url string) []byte { | |
| resp, err := http.Get(url) | |
| if err != nil { | |
| return []byte("") | |
| } | |
| defer resp.Body.Close() | |
| body, err := ioutil.ReadAll(resp.Body) | |
| if err != nil { | |
| return []byte("") | |
| } | |
| return body | |
| } | |
| type Content struct { | |
| url string | |
| content []byte | |
| } | |
| func makeStorage() chan *Content { | |
| ch := make(chan *Content) | |
| go func() { | |
| for { | |
| content := <-ch | |
| fmt.Printf("received: %s, length %d\n", | |
| content.url, | |
| len(content.content)) | |
| } | |
| }() | |
| return ch | |
| } | |
| func workerStat(query chan chan *utils.WorkerStat) { | |
| ticker := time.NewTicker(time.Second * 1) | |
| ret := make(chan *utils.WorkerStat) | |
| delta := uint64(0) | |
| var stat *utils.WorkerStat | |
| for _ = range ticker.C { | |
| query <- ret | |
| stat = <-ret | |
| fmt.Printf("total: %d, done: %d, pending: %d, delta: %d\n", | |
| stat.Total, stat.Done, stat.Total-stat.Done, stat.Done-delta) | |
| delta = stat.Done | |
| } | |
| } | |
| func main() { | |
| flag.Parse() | |
| fmt.Printf("pool size: %d\n", *poolSize) | |
| workers, query := utils.WorkerPool(*poolSize) | |
| storage := makeStorage() | |
| // worker statistic | |
| go workerStat(query) | |
| // start socket interface | |
| ln, err := net.Listen("tcp", ":43210") | |
| if err != nil { | |
| panic("Listen error") | |
| } | |
| for { | |
| conn, err := ln.Accept() | |
| if err != nil { | |
| continue | |
| } | |
| go func(conn net.Conn) { | |
| defer conn.Close() | |
| reader := bufio.NewReader(conn) | |
| for { | |
| line, err := reader.ReadString('\n') | |
| if err != nil { | |
| return | |
| } | |
| line = strings.TrimSpace(line) | |
| fmt.Printf("add: %s\n", line) | |
| workers <- func() { | |
| content := collect(line) | |
| storage <- &Content{line, content} | |
| } | |
| } | |
| }(conn) | |
| } | |
| } |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment