Започнете со пребарување: Совладување на конкурентно преземање податоци со издржливост
Продукцискиот crawler не е јамка околу http.Get. Тој е издржлива работна редица поврзана со ненадежна мрежа: URL-адресите се множат, хостовите возвраќаат, базите на податоци стануваат привремено недостапни, а успешен одговор може да пристигне мигови пред машината да изгуби напојување.
Овој туторијал гради ограничен Go crawler со глобални и ограничувања на истовременост по авторитет, откажување, повторни обиди за привремени грешки, дедупликација на URL-адреси, потврдени закупи, обновување на истечени закупи и издржлива SQLite фронта. Неговиот модел на испорака е најмалку еднаш: по пад, URL-адреса под закуп може повторно да се преземе. Контролираните несакани ефекти затоа користат идемпотентни операции со базата на податоци.
Предуслови и структура на проектот
Потребни ви се Go 1.22 или понов, излезен HTTP пристап и овластување за индексирање на избраните цели. Имплементацијата користи golang.org/x/net/html за парсирање HTML и modernc.org/sqlite, чист Go SQLite драјвер.
mkdir durable-crawler
cd durable-crawler
go mod init example.com/durable-crawler
go get golang.org/x/[email protected]
go get modernc.org/[email protected]
durable-crawler/
go.mod
go.sum
main.go
main_test.go
Архитектура: закупи пред HTTP
SQLite фронтата е и редицата и издржливата контролна точка. Секоја URL-адреса е pending, leased или done. Примарниот клуч на URL-адресата обезбедува издржлива дедупликација по намерно мала политика за каноникализација: мали букви за шемите и хостовите, изоставени стандардни порти, празни патеки претворени во / и отстранети фрагменти.
Преземањето е една атомска изјава за база на податоци што се потврдува пред работникот да чека капацитет или да изврши мрежен I/O. Секое успешно преземање добива нов случаен токен. Обновувањето, ослободувањето и завршувањето го бараат токму тој токен. Ако истечен закуп се преземе повторно, застарен работник не може да ја продолжи или заврши новата резервација.
Heartbeat го продолжува закупот за време на чекања на семафор, повторни обиди и преземања. Важно е дека heartbeat има сопствен сигнал за запирање. Завршувањето на задача го запира и приклучува heartbeat пред откажувањето на контекстот на задачата, така што намерното откажување не може да претвори обновување во тек во лажен резултат за загубен закуп.
Инфраструктурните неуспеси на базата на податоци се неуспеси на процесот, а не сигнали за празна редица. Грешките при преземање, броење на фронтата, heartbeat и завршување се пренесуваат до main, ги откажуваат другите работници и произведуваат излез различен од нула. Тоа му овозможува на управувач на услуги што користи Restart=on-failure да се опорави наместо тивко да прифати незавршени редови.
Оддалечениот GET сè уште може да се повтори по пад. Локална трансакција не може да направи оддалечен сервер идемпотентен, затоа индексирајте само крајни точки каде што повторените барања со безбедни методи се прифатливи.
Бројот на работници ја ограничува глобалната истовременост. Семафор со клуч според канонички авторитет го ограничува секој хост. Пренасочувањата се ограничени на истиот авторитет како претходното барање, така што секое барање во синџир на пренасочувања останува покриено од веќе задржаниот семафор. Ова ограничување на ниво на апликација исто така ги ограничува HTTP/2 тековите; само ограничувањата на конекции не го прават тоа.
Имплементирајте го crawler-от
Создадете main.go. Операциите со базата на податоци добиваат одделни рокови од две секунди, под закупот од две минути и буџетот од 25 секунди за исклучување при распоредување. HTTP клиентот независно го ограничува воспоставувањето конекција, TLS преговорите, заглавјата на одговорот и целото барање.
package main
import (
"bytes"
"context"
"crypto/rand"
"database/sql"
"encoding/hex"
"errors"
"expvar"
"flag"
"fmt"
"io"
"log/slog"
mrand "math/rand"
"net"
"net/http"
"net/url"
"os"
"os/signal"
"path/filepath"
"strconv"
"strings"
"sync"
"syscall"
"time"
"golang.org/x/net/html"
_ "modernc.org/sqlite"
)
var (
completed = expvar.NewInt("crawler_completed_total")
retries = expvar.NewInt("crawler_retries_total")
failures = expvar.NewInt("crawler_failures_total")
recovered = expvar.NewInt("crawler_expired_leases_total")
errCrossAuthority = errors.New("cross-authority redirect rejected")
errRedirectLimit = errors.New("redirect limit exceeded")
errLeaseLost = errors.New("lease ownership lost")
)
const (
dbTimeout = 2 * time.Second
leaseTime = 2 * time.Minute
)
type stringList []string
func (s *stringList) String() string { return strings.Join(*s, ",") }
func (s *stringList) Set(v string) error {
*s = append(*s, v)
return nil
}
type Frontier struct{ db *sql.DB }
func openFrontier(path string) (*Frontier, error) {
if err := os.MkdirAll(filepath.Dir(path), 0700); err != nil {
return nil, err
}
db, err := sql.Open("sqlite", path)
if err != nil {
return nil, err
}
db.SetMaxOpenConns(1)
db.SetMaxIdleConns(1)
for _, statement := range []string{
"PRAGMA journal_mode=WAL",
"PRAGMA synchronous=FULL",
"PRAGMA busy_timeout=2000",
`CREATE TABLE IF NOT EXISTS frontier (
url TEXT PRIMARY KEY,
status TEXT NOT NULL
CHECK(status IN ('pending','leased','done')),
lease_token TEXT,
lease_until INTEGER NOT NULL DEFAULT 0
)`,
`CREATE INDEX IF NOT EXISTS frontier_status_lease
ON frontier(status, lease_until)`,
} {
if _, err := db.Exec(statement); err != nil {
db.Close()
return nil, err
}
}
return &Frontier{db: db}, nil
}
func dbContext(parent context.Context) (context.Context, context.CancelFunc) {
return context.WithTimeout(parent, dbTimeout)
}
func claimToken() (string, error) {
var b [16]byte
if _, err := rand.Read(b[:]); err != nil {
return "", err
}
return hex.EncodeToString(b[:]), nil
}
func requireOwned(result sql.Result, err error) error {
if err != nil {
return err
}
n, err := result.RowsAffected()
if err != nil {
return err
}
if n != 1 {
return errLeaseLost
}
return nil
}
func (f *Frontier) enqueue(ctx context.Context, raw string) error {
ctx, cancel := dbContext(ctx)
defer cancel()
_, err := f.db.ExecContext(ctx,
`INSERT OR IGNORE INTO frontier(url,status)
VALUES(?, 'pending')`, raw)
return err
}
func (f *Frontier) recoverExpired(ctx context.Context) error {
ctx, cancel := dbContext(ctx)
defer cancel()
result, err := f.db.ExecContext(ctx, `UPDATE frontier
SET status='pending', lease_token=NULL, lease_until=0
WHERE status='leased' AND lease_until < ?`,
time.Now().UnixMilli())
if err == nil {
n, _ := result.RowsAffected()
recovered.Add(n)
}
return err
}
func (f *Frontier) claim(ctx context.Context) (string, string, bool, error) {
token, err := claimToken()
if err != nil {
return "", "", false, err
}
ctx, cancel := dbContext(ctx)
defer cancel()
now := time.Now().UnixMilli()
until := time.Now().Add(leaseTime).UnixMilli()
var raw string
err = f.db.QueryRowContext(ctx, `UPDATE frontier
SET status='leased', lease_token=?, lease_until=?
WHERE url = (
SELECT url FROM frontier
WHERE status='pending'
OR (status='leased' AND lease_until < ?)
ORDER BY url LIMIT 1
)
RETURNING url`, token, until, now).Scan(&raw)
if errors.Is(err, sql.ErrNoRows) {
return "", "", false, nil
}
if err != nil {
return "", "", false, err
}
return raw, token, true, nil
}
func (f *Frontier) renew(ctx context.Context, raw, token string) error {
ctx, cancel := dbContext(ctx)
defer cancel()
result, err := f.db.ExecContext(ctx, `UPDATE frontier
SET lease_until=?
WHERE url=? AND status='leased' AND lease_token=?`,
time.Now().Add(leaseTime).UnixMilli(), raw, token)
return requireOwned(result, err)
}
func (f *Frontier) release(raw, token string) error {
ctx, cancel := context.WithTimeout(context.Background(), dbTimeout)
defer cancel()
result, err := f.db.ExecContext(ctx, `UPDATE frontier
SET status='pending', lease_token=NULL, lease_until=0
WHERE url=? AND status='leased' AND lease_token=?`,
raw, token)
return requireOwned(result, err)
}
func (f *Frontier) finish(
ctx context.Context, raw, token string, links []string,
) error {
ctx, cancel := dbContext(ctx)
defer cancel()
tx, err := f.db.BeginTx(ctx, nil)
if err != nil {
return err
}
defer tx.Rollback()
for _, link := range links {
if _, err := tx.ExecContext(ctx,
`INSERT OR IGNORE INTO frontier(url,status)
VALUES(?, 'pending')`, link); err != nil {
return err
}
}
result, err := tx.ExecContext(ctx, `UPDATE frontier
SET status='done', lease_token=NULL, lease_until=0
WHERE url=? AND status='leased' AND lease_token=?`,
raw, token)
if err := requireOwned(result, err); err != nil {
return err
}
return tx.Commit()
}
func (f *Frontier) remaining(ctx context.Context) (int, error) {
ctx, cancel := dbContext(ctx)
defer cancel()
var n int
err := f.db.QueryRowContext(ctx,
`SELECT count(*) FROM frontier WHERE status!='done'`).Scan(&n)
return n, err
}
type hostLimits struct {
mu sync.Mutex
size int
hosts map[string]chan struct{}
}
func (h *hostLimits) forHost(authority string) chan struct{} {
h.mu.Lock()
defer h.mu.Unlock()
if ch := h.hosts[authority]; ch != nil {
return ch
}
ch := make(chan struct{}, h.size)
h.hosts[authority] = ch
return ch
}
type Crawler struct {
frontier *Frontier
client *http.Client
allowed map[string]struct{}
limits *hostLimits
logger *slog.Logger
maxBody int64
retryMax int
retryBase time.Duration
leaseEvery time.Duration
}
func canonical(raw string, base *url.URL) (string, *url.URL, error) {
u, err := url.Parse(raw)
if err != nil {
return "", nil, err
}
if base != nil {
u = base.ResolveReference(u)
}
u.Scheme = strings.ToLower(u.Scheme)
if u.Scheme != "http" && u.Scheme != "https" {
return "", nil, errors.New("unsupported URL scheme")
}
if u.User != nil || u.Hostname() == "" {
return "", nil, errors.New("credentials or empty host rejected")
}
host := strings.ToLower(u.Hostname())
port := u.Port()
if (u.Scheme == "http" && port == "80") ||
(u.Scheme == "https" && port == "443") {
port = ""
}
switch {
case port != "":
u.Host = net.JoinHostPort(host, port)
case strings.Contains(host, ":"):
u.Host = "[" + host + "]"
default:
u.Host = host
}
u.Fragment = ""
if u.Path == "" {
u.Path = "/"
}
return u.String(), u, nil
}
func redirectPolicy(req *http.Request, via []*http.Request) error {
if len(via) >= 10 {
return errRedirectLimit
}
if len(via) == 0 {
return nil
}
_, current, err := canonical(req.URL.String(), nil)
if err != nil {
return err
}
_, previous, err := canonical(via[len(via)-1].URL.String(), nil)
if err != nil {
return err
}
if current.Host != previous.Host {
return fmt.Errorf("%w: %s to %s",
errCrossAuthority, previous.Host, current.Host)
}
return nil
}
func extractLinks(body []byte, base *url.URL) []string {
doc, err := html.Parse(bytes.NewReader(body))
if err != nil {
return nil
}
var links []string
var walk func(*html.Node)
walk = func(n *html.Node) {
if n.Type == html.ElementNode && n.Data == "a" {
for _, a := range n.Attr {
if strings.EqualFold(a.Key, "href") {
if raw, _, err := canonical(a.Val, base); err == nil {
links = append(links, raw)
}
break
}
}
}
for child := n.FirstChild; child != nil; child = child.NextSibling {
walk(child)
}
}
walk(doc)
return links
}
func retryableStatus(status int) bool {
switch status {
case 408, 425, 429, 500, 502, 503, 504:
return true
}
return false
}
func retryAfter(value string) time.Duration {
if seconds, err := strconv.Atoi(strings.TrimSpace(value)); err == nil && seconds >= 0 {
return min(time.Duration(seconds)*time.Second, 30*time.Second)
}
if when, err := http.ParseTime(value); err == nil {
return min(max(time.Until(when), 0), 30*time.Second)
}
return 0
}
func (c *Crawler) fetch(
ctx context.Context, raw string,
) ([]string, bool, time.Duration, error) {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, raw, nil)
if err != nil {
return nil, false, 0, err
}
req.Header.Set("User-Agent", "durable-crawler/1.0")
resp, err := c.client.Do(req)
if err != nil {
if ctx.Err() != nil {
return nil, false, 0, ctx.Err()
}
if errors.Is(err, errCrossAuthority) ||
errors.Is(err, errRedirectLimit) {
return nil, false, 0, err
}
return nil, true, 0, err
}
defer resp.Body.Close()
if retryableStatus(resp.StatusCode) {
_, _ = io.Copy(io.Discard, io.LimitReader(resp.Body, 4096))
return nil, true, retryAfter(resp.Header.Get("Retry-After")),
fmt.Errorf("retryable HTTP status %d", resp.StatusCode)
}
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
return nil, false, 0,
fmt.Errorf("terminal HTTP status %d", resp.StatusCode)
}
if !strings.HasPrefix(
strings.ToLower(resp.Header.Get("Content-Type")), "text/html",
) {
return nil, false, 0, nil
}
body, err := io.ReadAll(io.LimitReader(resp.Body, c.maxBody+1))
if err != nil {
return nil, true, 0, err
}
if int64(len(body)) > c.maxBody {
return nil, false, 0,
fmt.Errorf("body exceeds %d bytes", c.maxBody)
}
return extractLinks(body, resp.Request.URL), false, 0, nil
}
func (c *Crawler) crawl(ctx context.Context, raw string) ([]string, error) {
for attempt := 0; ; attempt++ {
links, retry, delay, err := c.fetch(ctx, raw)
if err == nil || !retry || attempt >= c.retryMax {
return links, err
}
retries.Add(1)
if delay == 0 {
shift := min(attempt, 4)
base := c.retryBase * time.Duration(1<<shift)
delay = base +
time.Duration(mrand.Int63n(int64(base/2)+1))
}
c.logger.Warn("request retry",
"url", raw, "attempt", attempt+1,
"delay", delay, "error", err)
timer := time.NewTimer(delay)
select {
case <-ctx.Done():
timer.Stop()
return nil, ctx.Err()
case <-timer.C:
}
}
}
func (c *Crawler) maintainLease(
cancelJob context.CancelFunc, raw, token string,
) func() error {
stop := make(chan struct{})
done := make(chan error, 1)
interval := c.leaseEvery
if interval <= 0 {
interval = 30 * time.Second
}
go func() {
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-stop:
done <- nil
return
case <-ticker.C:
// Renewal is independent of deliberate job cancellation.
if err := c.frontier.renew(
context.Background(), raw, token,
); err != nil {
done <- err
cancelJob()
return
}
}
}
}()
return func() error {
close(stop)
return <-done
}
}
func reportFatal(
fatal chan<- error, cancel context.CancelFunc, err error,
) {
select {
case fatal <- err:
default:
}
cancel()
}
func (c *Crawler) reservationProblem(
raw, phase string, err error,
fatal chan<- error, cancel context.CancelFunc,
) bool {
if errors.Is(err, errLeaseLost) {
c.logger.Warn("lease ownership lost",
"url", raw, "phase", phase)
return false
}
wrapped := fmt.Errorf("%s %s: %w", phase, raw, err)
c.logger.Error("database operation failed",
"url", raw, "phase", phase, "error", err)
reportFatal(fatal, cancel, wrapped)
return true
}
func (c *Crawler) worker(
ctx context.Context, id int, fatal chan<- error,
cancelRun context.CancelFunc, wg *sync.WaitGroup,
) {
defer wg.Done()
for {
raw, token, ok, err := c.frontier.claim(ctx)
if err != nil {
if ctx.Err() == nil {
reportFatal(fatal, cancelRun,
fmt.Errorf("worker %d claim: %w", id, err))
}
return
}
if !ok {
n, err := c.frontier.remaining(ctx)
if err != nil {
if ctx.Err() == nil {
reportFatal(fatal, cancelRun,
fmt.Errorf("worker %d count frontier: %w", id, err))
}
return
}
if n == 0 {
return
}
timer := time.NewTimer(100 * time.Millisecond)
select {
case <-ctx.Done():
timer.Stop()
return
case <-timer.C:
continue
}
}
// Shutdown may have arrived while the blocking claim completed.
if ctx.Err() != nil {
if err := c.frontier.release(raw, token); err != nil {
if c.reservationProblem(
raw, "shutdown release", err, fatal, cancelRun,
) {
return
}
}
return
}
jobCtx, cancelJob := context.WithCancel(ctx)
stopLease := c.maintainLease(cancelJob, raw, token)
_, parsed, parseErr := canonical(raw, nil)
_, allowed := c.allowed[parsedHost(parsed)]
if parseErr != nil || !allowed {
leaseErr := stopLease()
cancelJob()
if leaseErr != nil {
if c.reservationProblem(
raw, "lease renewal", leaseErr, fatal, cancelRun,
) {
return
}
continue
}
if err := c.frontier.finish(
context.Background(), raw, token, nil,
); err != nil {
if c.reservationProblem(
raw, "completion", err, fatal, cancelRun,
) {
return
}
continue
}
completed.Add(1)
continue
}
limit := c.limits.forHost(parsed.Host)
select {
case limit <- struct{}{}:
case <-jobCtx.Done():
leaseErr := stopLease()
cancelJob()
if leaseErr != nil {
if c.reservationProblem(
raw, "lease renewal", leaseErr, fatal, cancelRun,
) {
return
}
}
if ctx.Err() != nil {
if err := c.frontier.release(raw, token); err != nil {
if c.reservationProblem(
raw, "shutdown release", err, fatal, cancelRun,
) {
return
}
}
return
}
continue
}
links, crawlErr := c.crawl(jobCtx, raw)
<-limit
// Join the heartbeat before canceling jobCtx. An in-flight
// renewal can finish normally instead of seeing context.Canceled.
leaseErr := stopLease()
cancelJob()
if leaseErr != nil {
if c.reservationProblem(
raw, "lease renewal", leaseErr, fatal, cancelRun,
) {
return
}
continue
}
if ctx.Err() != nil {
if err := c.frontier.release(raw, token); err != nil {
if c.reservationProblem(
raw, "shutdown release", err, fatal, cancelRun,
) {
return
}
}
return
}
var accepted []string
if crawlErr != nil {
failures.Add(1)
c.logger.Warn("crawl failed",
"worker", id, "url", raw, "error", crawlErr)
} else {
for _, link := range links {
canon, u, err := canonical(link, nil)
if err == nil {
if _, ok := c.allowed[u.Host]; ok {
accepted = append(accepted, canon)
}
}
}
}
if err := c.frontier.finish(
context.Background(), raw, token, accepted,
); err != nil {
if c.reservationProblem(
raw, "completion", err, fatal, cancelRun,
) {
return
}
continue
}
completed.Add(1)
}
}
func parsedHost(u *url.URL) string {
if u == nil {
return ""
}
return u.Host
}
func main() {
var seeds stringList
var workers, perHost, retryMax int
var dbPath, metricsAddr string
flag.Var(&seeds, "seed",
"seed URL; repeat to approve more authorities")
flag.IntVar(&workers, "workers", 16, "global worker limit")
flag.IntVar(&perHost, "per-host", 2,
"requests per authority")
flag.IntVar(&retryMax, "retries", 3,
"retries after the first attempt")
flag.StringVar(&dbPath, "db", "./state/frontier.db",
"SQLite frontier path")
flag.StringVar(&metricsAddr, "metrics", "127.0.0.1:9090",
"expvar address")
flag.Parse()
logger := slog.New(slog.NewJSONHandler(os.Stdout, nil))
if len(seeds) == 0 || workers < 1 ||
perHost < 1 || retryMax < 0 {
logger.Error("invalid configuration")
os.Exit(2)
}
frontier, err := openFrontier(dbPath)
if err != nil {
logger.Error("open frontier", "error", err)
os.Exit(1)
}
allowed := make(map[string]struct{})
for _, seed := range seeds {
raw, u, err := canonical(seed, nil)
if err != nil {
logger.Error("invalid seed",
"seed", seed, "error", err)
os.Exit(2)
}
allowed[u.Host] = struct{}{}
if err := frontier.enqueue(context.Background(), raw); err != nil {
logger.Error("enqueue seed", "error", err)
os.Exit(1)
}
}
if err := frontier.recoverExpired(context.Background()); err != nil {
logger.Error("recover leases", "error", err)
os.Exit(1)
}
transport := &http.Transport{
Proxy: http.ProxyFromEnvironment,
DialContext: (&net.Dialer{
Timeout: 3 * time.Second,
KeepAlive: 30 * time.Second,
}).DialContext,
ForceAttemptHTTP2: true,
MaxIdleConns: 100,
MaxIdleConnsPerHost: perHost,
MaxConnsPerHost: perHost,
IdleConnTimeout: 60 * time.Second,
TLSHandshakeTimeout: 5 * time.Second,
ResponseHeaderTimeout: 8 * time.Second,
}
client := &http.Client{
Transport: transport,
Timeout: 15 * time.Second,
CheckRedirect: redirectPolicy,
}
crawler := &Crawler{
frontier: frontier,
client: client,
allowed: allowed,
limits: &hostLimits{
size: perHost, hosts: make(map[string]chan struct{}),
},
logger: logger,
maxBody: 4 << 20,
retryMax: retryMax,
retryBase: 500 * time.Millisecond,
leaseEvery: 30 * time.Second,
}
signalCtx, stop := signal.NotifyContext(
context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop()
runCtx, cancelRun := context.WithCancel(signalCtx)
defer cancelRun()
fatal := make(chan error, 1)
mux := http.NewServeMux()
mux.Handle("/debug/vars", expvar.Handler())
metricsServer := &http.Server{
Addr: metricsAddr,
Handler: mux,
ReadHeaderTimeout: 2 * time.Second,
}
go func() {
if err := metricsServer.ListenAndServe();
!errors.Is(err, http.ErrServerClosed) {
reportFatal(fatal, cancelRun,
fmt.Errorf("metrics server: %w", err))
}
}()
var wg sync.WaitGroup
for id := 0; id < workers; id++ {
wg.Add(1)
go crawler.worker(runCtx, id, fatal, cancelRun, &wg)
}
wg.Wait()
var runErr error
select {
case runErr = <-fatal:
default:
}
checkpointCtx, checkpointCancel := context.WithTimeout(
context.Background(), dbTimeout)
if _, err := frontier.db.ExecContext(
checkpointCtx, "PRAGMA wal_checkpoint(TRUNCATE)",
); err != nil {
logger.Error("final WAL checkpoint", "error", err)
if runErr == nil {
runErr = fmt.Errorf("final WAL checkpoint: %w", err)
}
}
checkpointCancel()
shutdownCtx, shutdownCancel := context.WithTimeout(
context.Background(), 2*time.Second)
if err := metricsServer.Shutdown(shutdownCtx); err != nil && runErr == nil {
runErr = fmt.Errorf("metrics shutdown: %w", err)
}
shutdownCancel()
transport.CloseIdleConnections()
if err := frontier.db.Close(); err != nil && runErr == nil {
runErr = fmt.Errorf("close frontier: %w", err)
}
if runErr != nil {
logger.Error("crawler stopped with infrastructure failure",
"error", runErr)
os.Exit(1)
}
}
Тестирајте оградување, ограничувања, пренасочувања и пренос на неуспеси
Создадете main_test.go. Овие тестови покриваат повторни обиди за привремени грешки, потиснување на дупликатни врски, истовременост на ниво на апликација, отфрлање пренасочувања, обработка на истечени закупи, оградување со застарени токени, чисто исклучување на heartbeat и пренос на грешки од базата на податоци.
package main
import (
"context"
"fmt"
"io"
"log/slog"
"net/http"
"net/http/httptest"
"path/filepath"
"sync"
"sync/atomic"
"testing"
"time"
)
func runWorkers(ctx context.Context, c *Crawler, count int) error {
runCtx, cancel := context.WithCancel(ctx)
defer cancel()
fatal := make(chan error, 1)
var wg sync.WaitGroup
for id := 0; id < count; id++ {
wg.Add(1)
go c.worker(runCtx, id, fatal, cancel, &wg)
}
wg.Wait()
select {
case err := <-fatal:
return err
default:
return nil
}
}
func TestCrawler(t *testing.T) {
var current, maximum, flaky, expiredHits, targetHits atomic.Int32
target := httptest.NewServer(http.HandlerFunc(
func(w http.ResponseWriter, r *http.Request) {
targetHits.Add(1)
io.WriteString(w, "must not be reached")
}))
defer target.Close()
server := httptest.NewServer(http.HandlerFunc(
func(w http.ResponseWriter, r *http.Request) {
n := current.Add(1)
defer current.Add(-1)
for {
old := maximum.Load()
if n <= old || maximum.CompareAndSwap(old, n) {
break
}
}
w.Header().Set("Content-Type", "text/html")
switch r.URL.Path {
case "/":
fmt.Fprintf(w,
`<a href="/a">A</a>`+
`<a href="/a">duplicate</a>`+
`<a href="/flaky">F</a>`+
`<a href="/jump">J</a>`)
case "/flaky":
if flaky.Add(1) <= 2 {
w.WriteHeader(http.StatusServiceUnavailable)
return
}
io.WriteString(w, "recovered")
case "/jump":
http.Redirect(w, r, target.URL, http.StatusFound)
case "/expired":
expiredHits.Add(1)
io.WriteString(w, "reclaimed")
default:
time.Sleep(10 * time.Millisecond)
io.WriteString(w, "done")
}
}))
defer server.Close()
frontier, err := openFrontier(
filepath.Join(t.TempDir(), "frontier.db"))
if err != nil {
t.Fatal(err)
}
defer frontier.db.Close()
root, parsed, err := canonical(server.URL, nil)
if err != nil {
t.Fatal(err)
}
if err := frontier.enqueue(context.Background(), root); err != nil {
t.Fatal(err)
}
client := server.Client()
client.CheckRedirect = redirectPolicy
crawler := &Crawler{
frontier: frontier,
client: client,
allowed: map[string]struct{}{parsed.Host: {}},
limits: &hostLimits{
size: 1, hosts: make(map[string]chan struct{}),
},
logger: slog.New(slog.NewTextHandler(io.Discard, nil)),
maxBody: 1 << 20,
retryMax: 3,
retryBase: time.Millisecond,
leaseEvery: 20 * time.Millisecond,
}
ctx, cancel := context.WithTimeout(
context.Background(), 3*time.Second)
defer cancel()
if err := runWorkers(ctx, crawler, 4); err != nil {
t.Fatal(err)
}
if maximum.Load() > 1 {
t.Fatalf("per-host limit violated: %d", maximum.Load())
}
if flaky.Load() != 3 {
t.Fatalf("expected three flaky requests, got %d", flaky.Load())
}
if targetHits.Load() != 0 {
t.Fatal("cross-authority redirect reached target")
}
var total, done int
if err := frontier.db.QueryRow(
"SELECT count(*) FROM frontier").Scan(&total); err != nil {
t.Fatal(err)
}
if err := frontier.db.QueryRow(
"SELECT count(*) FROM frontier WHERE status='done'",
).Scan(&done); err != nil {
t.Fatal(err)
}
if total != 4 || done != 4 {
t.Fatalf("unexpected frontier: total=%d done=%d", total, done)
}
expiredURL := server.URL + "/expired"
if _, err := frontier.db.Exec(`INSERT INTO frontier
(url,status,lease_token,lease_until)
VALUES(?, 'leased', 'dead-token', 0)`,
expiredURL); err != nil {
t.Fatal(err)
}
if err := frontier.recoverExpired(context.Background()); err != nil {
t.Fatal(err)
}
if err := runWorkers(ctx, crawler, 1); err != nil {
t.Fatal(err)
}
var status string
if err := frontier.db.QueryRow(
"SELECT status FROM frontier WHERE url=?", expiredURL,
).Scan(&status); err != nil {
t.Fatal(err)
}
if status != "done" || expiredHits.Load() != 1 {
t.Fatalf("reclaimed job not processed: status=%s hits=%d",
status, expiredHits.Load())
}
}
func TestClaimTokenFencesStaleWorker(t *testing.T) {
frontier, err := openFrontier(
filepath.Join(t.TempDir(), "frontier.db"))
if err != nil {
t.Fatal(err)
}
defer frontier.db.Close()
const job = "https://example.test/"
if err := frontier.enqueue(context.Background(), job); err != nil {
t.Fatal(err)
}
raw, stale, ok, err := frontier.claim(context.Background())
if err != nil || !ok {
t.Fatalf("first claim: ok=%v err=%v", ok, err)
}
if _, err := frontier.db.Exec(
"UPDATE frontier SET lease_until=0 WHERE url=?", raw,
); err != nil {
t.Fatal(err)
}
raw, current, ok, err := frontier.claim(context.Background())
if err != nil || !ok {
t.Fatalf("reclaim: ok=%v err=%v", ok, err)
}
if stale == current {
t.Fatal("reclaimed job reused its claim token")
}
if err := frontier.renew(context.Background(), raw, stale); err == nil {
t.Fatal("stale renewal succeeded")
}
if err := frontier.release(raw, stale); err == nil {
t.Fatal("stale release succeeded")
}
if err := frontier.finish(
context.Background(), raw, stale, nil,
); err == nil {
t.Fatal("stale completion succeeded")
}
if err := frontier.finish(
context.Background(), raw, current, nil,
); err != nil {
t.Fatalf("current owner could not complete: %v", err)
}
}
func TestHeartbeatStopsCleanly(t *testing.T) {
frontier, err := openFrontier(
filepath.Join(t.TempDir(), "frontier.db"))
if err != nil {
t.Fatal(err)
}
defer frontier.db.Close()
const job = "https://example.test/heartbeat"
if err := frontier.enqueue(context.Background(), job); err != nil {
t.Fatal(err)
}
raw, token, ok, err := frontier.claim(context.Background())
if err != nil || !ok {
t.Fatalf("claim: ok=%v err=%v", ok, err)
}
crawler := &Crawler{
frontier: frontier,
leaseEvery: time.Millisecond,
}
_, cancelJob := context.WithCancel(context.Background())
stopLease := crawler.maintainLease(cancelJob, raw, token)
time.Sleep(5 * time.Millisecond)
if err := stopLease(); err != nil {
t.Fatalf("deliberate heartbeat stop reported loss: %v", err)
}
cancelJob()
if err := frontier.finish(
context.Background(), raw, token, nil,
); err != nil {
t.Fatalf("completion after heartbeat stop failed: %v", err)
}
}
func TestWorkerReportsDatabaseFailure(t *testing.T) {
frontier, err := openFrontier(
filepath.Join(t.TempDir(), "frontier.db"))
if err != nil {
t.Fatal(err)
}
if err := frontier.db.Close(); err != nil {
t.Fatal(err)
}
crawler := &Crawler{
frontier: frontier,
logger: slog.New(slog.NewTextHandler(io.Discard, nil)),
}
ctx, cancel := context.WithTimeout(
context.Background(), time.Second)
defer cancel()
if err := runWorkers(ctx, crawler, 1); err == nil {
t.Fatal("database failure was mistaken for successful completion")
}
}
Извршете го обичното однесување, инструментацијата за истовременост, статичката анализа и финалната изградба:
go test ./...
go test -race ./...
go vet ./...
go build -trimpath -o durable-crawler .
Стартувајте и набљудувајте го crawler-от
Секое почетно семе одобрува еден точен канонички авторитет, вклучувајќи нестандардна порта. Откриените врски надвор од тие авторитети се отфрлаат. Пренасочувањата се построги: секој скок мора да го задржи авторитетот на претходното барање, дури и ако друго семе ја одобрило дестинацијата.
./durable-crawler \
-seed https://docs.example.org/ \
-workers 24 \
-per-host 2 \
-retries 3 \
-db ./state/frontier.db \
-metrics 127.0.0.1:9090
curl --fail --silent http://127.0.0.1:9090/debug/vars
Структурираните логови изложуваат повторни обиди, терминални неуспеси при преземање, загубени закупи, неуспеси на базата на податоци и неуспеси на контролни точки. Expvar објавува бројачи за завршувања, повторни обиди, неуспеси и обновени закупи. Неговата крајна точка исто така изложува информации за извршувањето, затоа чувајте ја приватна и алармирајте при траен раст на повторните обиди, повторена загуба на закупи, рестартирања на услугата или фронта што престанува да завршува.
Перформанси, безбедност и учтивост
Дополнителните работници помагаат само кога работата опфаќа доволно авторитети. Голем број работници насочени кон една локација главно создаваат чекачи на семафор. Држете го -per-host конзервативен, објавете значаен user agent и имплементирајте ја robots-политиката на целта пред да индексирате системи што не ги контролирате.
SQLite е погоден за ограничени документациски локации, ревизии и оперативни инвентари. За распределен crawler е потребна серверска база на податоци, но состојбената машина останува иста: атомски преземете со единствен токен, потврдете пред HTTP, обновувајте преку ограничени барања, оградете ги застарените сопственици и задржете го чесниот договор „најмалку еднаш“.
Дозволената листа на авторитети не е целосна SSRF-заштита. DNS-одговорите може да се променат, одобрените имиња може да се разрешат кон приватна инфраструктура, а проксите ја менуваат ефективната мрежна патека. Наметнете политика за излезен сообраќај со контролиран прокси или заштитен ѕид. Задржете ја стандардната TLS-верификација на Go; попустливите TLS-поставки ги претвораат неуспесите на сертификати во безбедносни неуспеси.
Телото на одговорот е ограничено, но HTML-парсирањето сè уште конструира дрво на документот. Страниците со многу врски може да алоцираат значителна меморија. За непријателски или исклучително големи влезови, намалете го ограничувањето на телото и наметнете максимален број прифатени врски по страница.
Распоредување со systemd
Оваа единица за хост користи динамичен идентитет и доделува пристап за запишување само до нејзиниот директориум за состојба управуван од systemd. Прегледајте ги семето и патеката до извршната датотека пред да ја зачувате како durable-crawler.service.
[Unit]
Description=Durable bounded Go crawler
After=network-online.target
Wants=network-online.target
[Service]
Type=simple
DynamicUser=yes
StateDirectory=durable-crawler
ExecStart=/usr/local/bin/durable-crawler -seed https://docs.example.org/ -workers 24 -per-host 2 -retries 3 -db /var/lib/durable-crawler/frontier.db -metrics 127.0.0.1:9090
Restart=on-failure
RestartSec=5s
TimeoutStopSec=25s
NoNewPrivileges=yes
PrivateTmp=yes
ProtectSystem=strict
ProtectHome=yes
ProtectKernelTunables=yes
ProtectKernelModules=yes
ProtectControlGroups=yes
RestrictAddressFamilies=AF_INET AF_INET6
LockPersonality=yes
MemoryDenyWriteExecute=yes
[Install]
WantedBy=multi-user.target
Ова се привилегирани команди за хост, а не команди за контејнер. Проверете ги сите постојни одредишни датотеки пред инсталацијата; install ги заменува датотеките на тие точни патеки.
sudo stat /usr/local/bin/durable-crawler 2>/dev/null || true
sudo stat /etc/systemd/system/durable-crawler.service 2>/dev/null || true
sudo install -o root -g root -m 0755 \
./durable-crawler /usr/local/bin/durable-crawler
sudo install -o root -g root -m 0644 \
./durable-crawler.service \
/etc/systemd/system/durable-crawler.service
sudo systemctl daemon-reload
sudo systemctl enable --now durable-crawler.service
sudo systemctl status durable-crawler.service
sudo journalctl -u durable-crawler.service --since today
Применете политика за излезен заштитен ѕид одделно, без да ја изложите портата за loopback метрики. За време на пуштањето, тестирајте DNS, TLS-валидација, пренасочувања и однесување на проксито бидејќи блокирана зависност често прво се појавува како обичен истек на време.
Вообичаени продукциски неуспеси
- Дупликатни барања по рестартирање: очекувано кога процесот умира пред да го потврди завршувањето. Клучирајте ги контролираните надолни ефекти според каноничка URL-адреса или друг стабилен идемпотентен клуч.
- Застарен работник завршува повторно преземена работа: токените за преземање биле повторно употребени или изоставени од проверка на сопственост. Третирајте го токенот како идентитет за оградување на резервацијата, не како идентификатор на процесот.
- Успешна задача останува под закуп: запирањето на heartbeat со откажување на контекстот што се користи за обновување создава трка при откажување. Користете одделен сигнал за запирање на heartbeat, приклучете го, па дури потоа откажете го контекстот на задачата.
- Работниците исчезнуваат додека редовите остануваат: грешка во базата на податоци била третирана како празна редица. Пренесете ги инфраструктурните неуспеси и излезете со вредност различна од нула, или повторете ги според експлицитна ограничена политика.
- Исклучувањето започнува нова работа: повторно проверете го откажувањето веднаш по секое преземање. Ако исклучувањето пристигнало за време на повикот до базата на податоци, ослободете ја токму таа резервација без преземање.
- Пренасочувањата ги заобиколуваат ограничувањата на авторитет: бројот на конекции не ги ограничува HTTP/2 тековите. Задржете семафор на ниво на апликација и отфрлете или одделно регулирајте ги пренасочувањата меѓу авторитети.
- Повторните обиди го засилуваат прекинот: ограничете ги обидите, држете го семафорот низ повторните обиди, додајте случајно отстапување и ограничете го
Retry-After. - SQLite развива притисок од заклучувања: чувајте ги трансакциите без мрежна работа и избегнувајте непотребни конекции. Континуирано зафатен работен товар со повеќе процеси ја надраснал оваа опција за складирање.
Конечна листа за проверка
- Успешно извршете
go test -race ./.... - Потврдете дека глобалната истовременост и истовременоста по авторитет остануваат во рамките на конфигурацијата.
- Проверете дека пренасочувањата меѓу авторитети никогаш не ја контактираат својата дестинација.
- Запрете heartbeat за време на обновување и потврдете дека задачата во сопственост сè уште може да заврши.
- Затворете ја базата на податоци или направете ја недостапна и потврдете дека процесот излегува со вредност различна од нула.
- Испратете
SIGTERMза време на бавно барање и потврдете дека резервацијата во сопственост е ослободена. - Нагло прекинете го процесот, оставете закупот да истече и потврдете дека URL-адресата е повторно преземена.
- Проверете дека застарените токени не можат да обноват, ослободат или завршат повторно преземена работа.
- Тестирајте дупликатни врски, преголеми тела, привремени статуси и терминални статуси.
- Потврдете дека директориумот за состојба и слушачот за метрики се достапни само за наменетите идентитети.
- Прегледајте ги контролите за излезен сообраќај и овластувањето за секоја цел за индексирање.
Издржливоста не е само запишување URL-адреси на диск. Таа е зачувување на сопственоста низ прекин, разликување празна редица од расипана инфраструктура и оградување на вчерашните работници од денешните преземања. Штом тие граници се експлицитни, истовременоста престанува да биде коцкање и станува оперативна алатка.