Vodiči

Robust Go Worker Pools: Context, Retries, and Idempotency Unpacked

Robusni skupovi radnika u Gou: objašnjeni kontekst, ponovni pokušaji i idempotentnost

Skup radnika lako je demonstrirati, ali iznenađujuće ga je teško održavati u radu. Pokrenite nekoliko gorutina, dodijelite im poslove i sretni put izgleda dovršen. Produkcija donosi neugodna pitanja: Što se događa kada proces umre usred posla? Može li gašenje slučajno pokrenuti novi rad? Hoće li ponovni pokušaj ponoviti nepovratan učinak? Tko posjeduje posao nakon isteka njegovog zakupa?

Ovaj vodič izrađuje Go skup radnika s PostgreSQL pozadinom koji izričito odgovara na ta pitanja. Koristi kratke rezervacije, obnovljive zakupe, ograničene kontekste, eksponencijalne ponovne pokušaje, provjere vlasništva i granicu idempotentnosti koju provodi baza podataka. Njegov ugovor isporuke je najmanje jednom. Posao se može ponovno izvršiti nakon neizvjesnosti, ali njegov se potvrđeni učinak primjenjuje jednom po ključu idempotentnosti.

Preduvjeti i raspored projekta

Potrebni su vam Go 1.22 ili noviji, PostgreSQL 14 ili noviji i uloga baze podataka kojoj je dopušteno čitati i mijenjati tablice aplikacije. Jedina Go ovisnost je github.com/jackc/pgx/v5; ovaj primjer prikvačuje verziju 5.7.2 i namijenjen je kompatibilnoj liniji 5.7.x.

reliable-worker/
├── go.mod
├── main.go
├── schema.sql
└── deploy/
    └── reliable-worker.service

Izradite go.mod:

module example.com/reliable-worker

go 1.22

require github.com/jackc/pgx/v5 v5.7.2

Arhitektura: odvojite vlasništvo od izvršavanja

Radnici preuzimaju poslove jednom atomskom PostgreSQL naredbom koristeći FOR UPDATE SKIP LOCKED. Naredba se potvrđuje prije početka poslovne obrade, tako da nijedna transakcija baze podataka ne ostaje otvorena dok se posao izvršava.

Preuzeti redak bilježi jedinstvenog vlasnika radnika i rok zakupa. Otkuucaj srca produljuje taj rok. Ako proces nestane, drugi radnik može ponovno preuzeti posao nakon isteka. Svako dovršavanje, ponovni pokušaj, otpuštanje i obnova u svom predikatu uključuje vlasnika. Ova provjera ograđivanja sprječava radnika s isteklim zakupom da promijeni posao koji je sada u vlasništvu drugog.

Važni prijelazi su:

  • ready u running pri preuzimanju, uz povećanje broja pokušaja.
  • running u succeeded kada se učinak i dovršavanje potvrde zajedno.
  • running u ready nakon neuspjeha, s odgođenim available_at.
  • Istekli posao running ponovno u running pod novim vlasnikom.
  • Iscrpljeni posao u dead.

Izradite trajni red

Ključ idempotentnosti jedinstven je i u tablicama reda i učinaka. Ovdje je poslovni učinak namjerno jednostavan: bilježenje obrađene poruke. U stvarnom sustavu tablica učinaka može predstavljati račun, zahtjev za obavijest, promjenu računa ili unos izlaznog spremnika.

CREATE TABLE jobs (
    id              bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
    idempotency_key text NOT NULL UNIQUE,
    payload         jsonb NOT NULL,
    state           text NOT NULL DEFAULT 'ready'
                    CHECK (state IN ('ready', 'running', 'succeeded', 'dead')),
    attempts        integer NOT NULL DEFAULT 0 CHECK (attempts >= 0),
    max_attempts     integer NOT NULL DEFAULT 5 CHECK (max_attempts > 0),
    available_at    timestamptz NOT NULL DEFAULT now(),
    lease_owner     text,
    lease_until     timestamptz,
    last_error      text,
    created_at      timestamptz NOT NULL DEFAULT now(),
    completed_at    timestamptz
);

CREATE INDEX jobs_claimable_idx
    ON jobs (available_at, id)
    WHERE state = 'ready';

CREATE INDEX jobs_expired_idx
    ON jobs (lease_until)
    WHERE state = 'running';

CREATE TABLE processed_messages (
    idempotency_key text PRIMARY KEY,
    message         text NOT NULL,
    processed_at    timestamptz NOT NULL DEFAULT now()
);

Primijenite shemu putem postojeće, izričito odabrane veze s bazom podataka:

cd reliable-worker
export DATABASE_URL='postgres://worker_user:[email protected]:5432/workerdb?sslmode=require'
psql "$DATABASE_URL" -v ON_ERROR_STOP=1 -f schema.sql
go mod tidy

Upotrijebite TLS provjeru primjerenu svojem okruženju. sslmode=require šifrira prijenos, ali ne pruža istu provjeru identiteta poslužitelja kao verify-full s pouzdanim CA-om.

Implementirajte skup radnika

Donja implementacija koristi zakup od 30 sekundi, obnavlja ga svakih 10 sekundi, ograničava preuzimanja na dvije sekunde i naredbe baze podataka na tri sekunde. Ta su ograničenja znatno ispod trajanja zakupa. Vremensko ograničenje veze konfigurirano je zasebno jer ne ograničava izvršavanje upita.

package main

import (
	"context"
	"crypto/rand"
	"encoding/hex"
	"encoding/json"
	"errors"
	"fmt"
	"log/slog"
	"os"
	"os/signal"
	"strconv"
	"sync"
	"syscall"
	"time"

	"github.com/jackc/pgx/v5"
	"github.com/jackc/pgx/v5/pgxpool"
)

var errLeaseLost = errors.New("lease lost")

type config struct {
	Workers, MaxConns int
	Lease, ClaimTimeout, QueryTimeout, JobTimeout time.Duration
}

type job struct {
	ID int64
	Key string
	Payload []byte
	Attempts, MaxAttempts int
}

type payload struct {
	Message string `json:"message"`
	SleepMS int    `json:"sleep_ms"`
}

func main() {
	ctx, stop := signal.NotifyContext(context.Background(),
		syscall.SIGINT, syscall.SIGTERM)
	defer stop()

	cfg := config{
		Workers: envInt("WORKERS", 4),
		Lease: 30 * time.Second,
		ClaimTimeout: 2 * time.Second,
		QueryTimeout: 3 * time.Second,
		JobTimeout: 20 * time.Second,
	}
	cfg.MaxConns = cfg.Workers + 2

	poolCfg, err := pgxpool.ParseConfig(mustEnv("DATABASE_URL"))
	if err != nil {
		panic(err)
	}
	poolCfg.MaxConns = int32(cfg.MaxConns)
	poolCfg.MinConns = 1
	poolCfg.MaxConnLifetime = 30 * time.Minute
	poolCfg.MaxConnIdleTime = 5 * time.Minute
	poolCfg.HealthCheckPeriod = 30 * time.Second
	poolCfg.ConnConfig.ConnectTimeout = 3 * time.Second
	poolCfg.AfterConnect = func(ctx context.Context, c *pgx.Conn) error {
		for _, q := range []string{
			"SET statement_timeout = '3s'",
			"SET lock_timeout = '1s'",
			"SET idle_in_transaction_session_timeout = '5s'",
		} {
			if _, err := c.Exec(ctx, q); err != nil {
				return err
			}
		}
		return nil
	}

	pool, err := pgxpool.NewWithConfig(ctx, poolCfg)
	if err != nil {
		panic(err)
	}
	defer pool.Close()

	logger := slog.New(slog.NewJSONHandler(os.Stdout, nil))
	processID := randomID()

	var wg sync.WaitGroup
	for i := 0; i < cfg.Workers; i++ {
		wg.Add(1)
		owner := fmt.Sprintf("%s-%d", processID, i)
		go func() {
			defer wg.Done()
			worker(ctx, pool, logger, cfg, owner)
		}()
	}

	wg.Add(1)
	go func() {
		defer wg.Done()
		reaper(ctx, pool, logger, cfg.QueryTimeout)
	}()

	<-ctx.Done()
	logger.Info("shutdown_started")
	wg.Wait()
	logger.Info("shutdown_complete")
}

func worker(ctx context.Context, db *pgxpool.Pool, log *slog.Logger,
	cfg config, owner string) {
	for {
		claimCtx, cancel := context.WithTimeout(ctx, cfg.ClaimTimeout)
		j, err := claim(claimCtx, db, owner, cfg.Lease)
		cancel()

		// Shutdown may have arrived while the blocking claim committed.
		if ctx.Err() != nil {
			if err == nil {
				release(db, j.ID, owner, cfg.QueryTimeout)
			}
			return
		}
		if errors.Is(err, pgx.ErrNoRows) {
			if !wait(ctx, 500*time.Millisecond) {
				return
			}
			continue
		}
		if err != nil {
			log.Error("claim_failed", "owner", owner, "error", err)
			if !wait(ctx, time.Second) {
				return
			}
			continue
		}

		jobCtx, timeoutCancel := context.WithTimeout(ctx, cfg.JobTimeout)
		runCtx, causeCancel := context.WithCancelCause(jobCtx)
		renewDone := make(chan struct{})
		stopRenew := make(chan struct{})
		go renew(runCtx, db, j.ID, owner, cfg, causeCancel,
			stopRenew, renewDone)

		err = handle(runCtx, db, j, owner, cfg.QueryTimeout)
		close(stopRenew)
		<-renewDone
		cause := context.Cause(runCtx)
		causeCancel(nil)
		timeoutCancel()

		switch {
		case err == nil:
			log.Info("job_succeeded", "job_id", j.ID,
				"attempt", j.Attempts)
		case errors.Is(cause, errLeaseLost):
			log.Warn("job_abandoned", "job_id", j.ID,
				"reason", "lease_lost")
		case ctx.Err() != nil:
			release(db, j.ID, owner, cfg.QueryTimeout)
			return
		default:
			if failErr := fail(db, j, owner, err, cfg.QueryTimeout); failErr != nil {
				log.Error("failure_update_failed", "job_id", j.ID,
					"error", failErr)
			}
			log.Warn("job_failed", "job_id", j.ID,
				"attempt", j.Attempts, "error", err)
		}
	}
}

func claim(ctx context.Context, db *pgxpool.Pool, owner string,
	lease time.Duration) (job, error) {
	const q = `
WITH candidate AS (
	SELECT id
	FROM jobs
	WHERE attempts < max_attempts
	  AND (
		(state = 'ready' AND available_at <= now())
		OR (state = 'running' AND lease_until < now())
	  )
	ORDER BY available_at, id
	FOR UPDATE SKIP LOCKED
	LIMIT 1
)
UPDATE jobs AS j
SET state = 'running',
    lease_owner = $1,
    lease_until = now() + ($2 * interval '1 millisecond'),
    attempts = attempts + 1
FROM candidate
WHERE j.id = candidate.id
RETURNING j.id, j.idempotency_key, j.payload,
          j.attempts, j.max_attempts`
	var j job
	err := db.QueryRow(ctx, q, owner, lease.Milliseconds()).Scan(
		&j.ID, &j.Key, &j.Payload, &j.Attempts, &j.MaxAttempts)
	return j, err
}

func renew(ctx context.Context, db *pgxpool.Pool, id int64, owner string,
	cfg config, cancel context.CancelCauseFunc, stop <-chan struct{},
	done chan<- struct{}) {
	defer close(done)
	ticker := time.NewTicker(cfg.Lease / 3)
	defer ticker.Stop()

	for {
		select {
		case <-ctx.Done():
			return
		case <-stop:
			return
		case <-ticker.C:
			qctx, qcancel := context.WithTimeout(ctx, cfg.QueryTimeout)
			tag, err := db.Exec(qctx, `
UPDATE jobs
SET lease_until = now() + ($3 * interval '1 millisecond')
WHERE id = $1 AND lease_owner = $2 AND state = 'running'`,
				id, owner, cfg.Lease.Milliseconds())
			qcancel()
			if err != nil || tag.RowsAffected() != 1 {
				cancel(errLeaseLost)
				return
			}
		}
	}
}

func handle(ctx context.Context, db *pgxpool.Pool, j job,
	owner string, queryTimeout time.Duration) error {
	var p payload
	if err := json.Unmarshal(j.Payload, &p); err != nil {
		return fmt.Errorf("decode payload: %w", err)
	}
	if p.Message == "" {
		return errors.New("message is required")
	}
	if p.SleepMS < 0 || p.SleepMS > 15000 {
		return errors.New("sleep_ms must be between 0 and 15000")
	}
	if !wait(ctx, time.Duration(p.SleepMS)*time.Millisecond) {
		return context.Cause(ctx)
	}

	qctx, cancel := context.WithTimeout(ctx, queryTimeout)
	defer cancel()
	tx, err := db.Begin(qctx)
	if err != nil {
		return err
	}
	defer tx.Rollback(context.Background())

	_, err = tx.Exec(qctx, `
INSERT INTO processed_messages (idempotency_key, message)
VALUES ($1, $2)
ON CONFLICT (idempotency_key) DO NOTHING`, j.Key, p.Message)
	if err != nil {
		return err
	}

	tag, err := tx.Exec(qctx, `
UPDATE jobs
SET state = 'succeeded', completed_at = now(),
    lease_owner = NULL, lease_until = NULL, last_error = NULL
WHERE id = $1 AND lease_owner = $2 AND state = 'running'`,
		j.ID, owner)
	if err != nil {
		return err
	}
	if tag.RowsAffected() != 1 {
		return errLeaseLost
	}
	return tx.Commit(qctx)
}

func fail(db *pgxpool.Pool, j job, owner string, jobErr error,
	timeout time.Duration) error {
	delay := time.Second << min(j.Attempts-1, 6)
	ctx, cancel := context.WithTimeout(context.Background(), timeout)
	defer cancel()
	_, err := db.Exec(ctx, `
UPDATE jobs
SET state = CASE WHEN attempts >= max_attempts THEN 'dead'
                 ELSE 'ready' END,
    available_at = now() + ($4 * interval '1 millisecond'),
    lease_owner = NULL, lease_until = NULL, last_error = $3
WHERE id = $1 AND lease_owner = $2 AND state = 'running'`,
		j.ID, owner, truncate(jobErr.Error(), 1000), delay.Milliseconds())
	return err
}

func release(db *pgxpool.Pool, id int64, owner string, timeout time.Duration) {
	ctx, cancel := context.WithTimeout(context.Background(), timeout)
	defer cancel()
	_, _ = db.Exec(ctx, `
UPDATE jobs
SET state = 'ready', available_at = now(),
    lease_owner = NULL, lease_until = NULL
WHERE id = $1 AND lease_owner = $2 AND state = 'running'`,
		id, owner)
}

func reaper(ctx context.Context, db *pgxpool.Pool, log *slog.Logger,
	timeout time.Duration) {
	ticker := time.NewTicker(10 * time.Second)
	defer ticker.Stop()
	for {
		qctx, cancel := context.WithTimeout(ctx, timeout)
		_, err := db.Exec(qctx, `
UPDATE jobs
SET state = 'dead', lease_owner = NULL, lease_until = NULL,
    last_error = COALESCE(last_error, 'lease expired after final attempt')
WHERE state = 'running' AND lease_until < now()
  AND attempts >= max_attempts`)
		cancel()
		if err != nil && ctx.Err() == nil {
			log.Error("reaper_failed", "error", err)
		}
		select {
		case <-ctx.Done():
			return
		case <-ticker.C:
		}
	}
}

func wait(ctx context.Context, d time.Duration) bool {
	timer := time.NewTimer(d)
	defer timer.Stop()
	select {
	case <-ctx.Done():
		return false
	case <-timer.C:
		return true
	}
}

func envInt(name string, fallback int) int {
	if value := os.Getenv(name); value != "" {
		n, err := strconv.Atoi(value)
		if err == nil && n > 0 {
			return n
		}
	}
	return fallback
}

func mustEnv(name string) string {
	value := os.Getenv(name)
	if value == "" {
		panic(name + " is required")
	}
	return value
}

func randomID() string {
	b := make([]byte, 16)
	if _, err := rand.Read(b); err != nil {
		panic(err)
	}
	return hex.EncodeToString(b)
}

func truncate(s string, n int) string {
	if len(s) <= n {
		return s
	}
	return s[:n]
}

Zašto je transakcija učinka važna

processed_messages i ažuriranje dovršetka posla potvrđuju se u jednoj transakciji. Ako radnik više ne posjeduje zakup, ažuriranje utječe na nula redaka i transakcija se poništava, uključujući svaki novi učinak. Ako je prethodni pokušaj potvrdio učinak, ali je red u redu nekako ostao spreman za ponovni pokušaj, ON CONFLICT DO NOTHING čini ponovnu obradu bezopasnom.

Ovaj obrazac štiti samo učinke unutar iste baze podataka. Za vanjsko plaćanje ili HTTP API pošaljite stabilni ključ idempotentnosti nizvodnoj usluzi koja ga podržava ili potvrdite red transakcijskog izlaznog spremnika i prepustite drugom idempotentnom dispečeru mrežni poziv. Lokalna oznaka „obrađeno” ne može atomski dokazati da se dogodio nepovezan udaljeni sporedni učinak.

Testirajte ponovne pokušaje, zakupe i gašenje

Počnite s jednim radnikom kako biste prijelaze lako pregledali:

export DATABASE_URL='postgres://worker_user:[email protected]:5432/workerdb?sslmode=require'
export WORKERS=1
go run .

Iz druge ljuske stavite u red jedan uspješan posao i jedan deterministički neuspjeh:

psql "$DATABASE_URL" -v ON_ERROR_STOP=1 <<'SQL'
INSERT INTO jobs (idempotency_key, payload, max_attempts)
VALUES
  ('welcome-1001', '{"message":"welcome","sleep_ms":8000}', 5),
  ('invalid-1002', '{"sleep_ms":0}', 3);
SQL

psql "$DATABASE_URL" -c \
"SELECT id, state, attempts, lease_owner, lease_until, last_error FROM jobs ORDER BY id;"
psql "$DATABASE_URL" -c \
"SELECT idempotency_key, message FROM processed_messages ORDER BY idempotency_key;"

Tijekom osamsekundnog posla pošaljite SIGTERM pomoću Ctrl-C. Radnik otkazuje obradu i uvjetno oslobađa svoju rezervaciju. Ponovno ga pokrenite i potvrdite da valjani posao na kraju uspije s jednim retkom učinka. Njegov broj pokušaja može biti veći od jedan; to je očekivano pri isporuci najmanje jednom. Nevaljani posao ponavlja se uz ograničena eksponencijalna kašnjenja i postaje dead nakon trećeg preuzimanja.

Kako biste testirali oporavak nakon rušenja umjesto urednog otpuštanja, prisilno prekinite proces u izoliranom testnom okruženju. Redak ostaje running dok njegov zakup od 30 sekundi ne istekne, a zatim ga je moguće preuzeti. Nikada ne koristite prisilni signal kao uobičajeni mehanizam gašenja.

Vidljivost i performanse

Radnik emitira strukturirane JSON zapisnike koji sadržavaju ID-ove poslova, pokušaje, vlasnike i pogreške. Izbjegavajte bilježenje potpunih tereta: mogu sadržavati vjerodajnice ili osobne podatke. Korisna mjerenja baze podataka uključuju broj spremnih poslova, starost najstarijeg spremnog posla, aktivne zakupe, istekle zakupe, mrtve poslove, stopu ponovnih pokušaja, trajanje obrade i broj gubitaka zakupa.

Skalirajte radnike prema izmjerenom kapacitetu baze podataka i nizvodnih sustava, a ne samo prema broju procesora. Svaki radnik može držati jednu vezu tijekom upita, dok otkucaji srca i čistač trebaju slobodan kapacitet; zato spremište dopušta workers + 2 veza. Neka tereti poslova budu mali, arhivirajte stare uspješno dovršene retke i pregledavajte planove upita kako tablica raste. Djelomični indeksi održavaju skeniranja pri preuzimanju usredotočenima, ali ne zamjenjuju rutinsko održavanje PostgreSQL-a.

Sigurnost i implementacija

Koristite namjenski račun operacijskog sustava i ulogu baze podataka s najmanjim potrebnim ovlastima. Radniku nije potreban ulazni mrežni port, stoga ga nemojte otvoriti u vatrozidu glavnog računala. Ograničite izlazni pristup na PostgreSQL i stvarne nizvodne ovisnosti tamo gdje arhitektura vatrozida podržava pravila izlaznog prometa. URL baze podataka pohranite izvan datoteke jedinice, zaštitite ga vlasništvom roota i načinom 0600 te rotirajte vjerodajnice bez njihova ugrađivanja u binarnu datoteku.

[Unit]
Description=Reliable Go worker pool
After=network-online.target
Wants=network-online.target

[Service]
Type=simple
User=reliable-worker
Group=reliable-worker
WorkingDirectory=/opt/reliable-worker
EnvironmentFile=/etc/reliable-worker/worker.env
ExecStart=/opt/reliable-worker/reliable-worker
Restart=on-failure
RestartSec=3s
TimeoutStopSec=15s
KillSignal=SIGTERM
NoNewPrivileges=true
PrivateTmp=true
ProtectSystem=strict
ProtectHome=true
RestrictAddressFamilies=AF_UNIX AF_INET AF_INET6

[Install]
WantedBy=multi-user.target

Na glavnom računalu administrator može izgraditi binarnu datoteku, smjestiti je pod /opt/reliable-worker, izraditi servisni račun i zaštićenu datoteku okruženja, instalirati jedinicu, a zatim pokrenuti systemctl daemon-reload i systemctl enable --now reliable-worker. Osigurajte da TimeoutStopSec ostane veći od ograničenog upita za otpuštanje. Kontejneri trebaju koristiti ista načela neprivilegiranog korisnika i datotečnog sustava samo za čitanje te moraju dobiti dovoljno vremena za uredno gašenje prije prisilnog uklanjanja.

Uobičajeni produkcijski kvarovi

  • Zakup kraći od realnog vremena obrade: obnavljajte ga mnogo prije isteka i ograničite svaku operaciju. Duga zastajkivanja i dalje mogu izgubiti vlasništvo, stoga dovršavanje mora ostati ograđeno.
  • Ponovni pokušaji uz držanje transakcije: prvo preuzmite i potvrdite. Pokretanje poslovnog koda unutar transakcije rezervacije stvara zaključavanja, izgladnjivanje veza i otežan oporavak.
  • Pokretanje rada tijekom gašenja: uvijek provjerite otkazivanje odmah nakon blokirajućeg preuzimanja. Ako je preuzimanje potvrđeno tijekom gašenja, otpustite ga ažuriranjem koje provjerava vlasništvo.
  • Pretpostavka da vremensko ograničenje veze ograničava upite: zasebno konfigurirajte rokove za vezu, naredbu, zaključavanje i kontekst.
  • Nazivanje ponovnih pokušaja točno jednom: zakupi rješavaju napuštanje, a ne neizvjesnost. Idempotentnost na granici učinka ono je što ponovljenu isporuku čini sigurnom.

Završni kontrolni popis za provjeru

  • Shema i djelomični indeksi primjenjuju se bez pogrešaka.
  • Istodobni radnici nikada ne preuzimaju isti aktivni zakup.
  • Uredno gašenje ne pokreće novi poslovni rad nakon dovršenog preuzimanja.
  • Posao srušenog radnika postaje moguće preuzeti nakon isteka zakupa.
  • Otkucaji srca zaustavljaju obradu kada se izgubi vlasništvo.
  • Neuspjeli poslovi ponavljaju se s odgodom i na kraju prelaze u dead.
  • Ponovljena isporuka proizvodi jedan potvrđeni učinak po ključu idempotentnosti.
  • Vremenska ograničenja baze podataka i gašenja ostaju ispod proračuna zakupa.
  • Zapisnici i upiti reda otkrivaju zaostatak, neuspjehe i zdravlje zakupa.

Robustan skup radnika nije definiran brzinom kojom troši kanal. Definiran je onime što ostaje istinito kada vremenski uvjeti postanu neprijateljski: rezervacije su kratke, vlasništvo istječe, gašenje je namjerno, ponovni pokušaji su ograničeni, a učinci podnose ponavljanje. Kada su te invarijante vidljive u shemi i provedene u svakom ažuriranju, neuspjeh prestaje biti iznimna putanja. Postaje uobičajeni prijelaz stanja koji sustav već zna preživjeti.

Portret autora bloga

Mihajlo

Ja sam Mihajlo — programer vođen znatiželjom, disciplinom i stalnom željom da stvorim nešto smisleno. Dijelim uvide, tutorijale i besplatne usluge kako bih pomogao drugima da pojednostave svoj rad i rastu u svijetu softvera i umjetne inteligencije koji se neprestano razvija.