Туториали

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

Стабилни Go работнички пулови: Објаснети контекст, повторни обиди и идемпотентност

Пулот на работници лесно се демонстрира, а изненадувачки тешко се управува. Започнете неколку goroutine-и, доставувајте им задачи и среќниот пат изгледа завршен. Продукцијата ги воведува непријатните прашања: Што се случува кога процес ќе згасне на половина задача? Може ли исклучувањето случајно да започне нова работа? Дали повторниот обид ќе повтори неповратен ефект? Кој е сопственик на задача откако ќе ѝ истече закупот?

Овој туторијал гради Go пул на работници поддржан од PostgreSQL кој експлицитно одговара на тие прашања. Користи кратки резервации, обновливи закупи, ограничени контексти, експоненцијални повторни обиди, проверки на сопственост и граница на идемпотентност наметната од базата на податоци. Неговиот договор за испорака е најмалку еднаш. Задачата може повторно да се изврши поради неизвесност, но нејзиниот потврден ефект се применува еднаш по клуч за идемпотентност.

Предуслови и распоред на проектот

Ви треба Go 1.22 или понов, PostgreSQL 14 или понов и улога во базата на податоци на која ѝ е дозволено да ги чита и менува табелите на апликацијата. Единствената Go зависност е github.com/jackc/pgx/v5; овој пример ја фиксира верзијата 5.7.2 и е наменет за компатибилната линија 5.7.x.

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

Создадете go.mod:

module example.com/reliable-worker

go 1.22

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

Архитектура: одделете ја сопственоста од извршувањето

Работниците преземаат задачи со една атомска PostgreSQL наредба користејќи FOR UPDATE SKIP LOCKED. Наредбата се потврдува пред да започне деловната обработка, па ниедна трансакција на базата на податоци не останува отворена додека се извршува задача.

Преземениот ред бележи единствен сопственик-работник и рок на закуп. Отчукувањето на срцето го продолжува тој рок. Ако процесот исчезне, друг работник може повторно да ја преземе задачата по истекувањето. Секое завршување, повторен обид, ослободување и обновување го вклучува сопственикот во својот предикат. Оваа проверка за оградување спречува работник со истечен закуп да измени задача што сега е во сопственост на друг.

Важните премини се:

  • ready во running при преземање, со зголемување на бројот на обиди.
  • running во succeeded кога ефектот и завршувањето се потврдуваат заедно.
  • running во ready по неуспех, со одложен available_at.
  • Задача во running со истечен рок повторно во running под нов сопственик.
  • Исцрпена задача во dead.

Создадете трајна редица

Клучот за идемпотентност е единствен и во редицата и во табелите со ефекти. Тука деловниот ефект е намерно едноставен: бележење обработена порака. Во реален систем, табелата со ефекти би можела да претставува фактура, барање за известување, промена на сметка или запис во outbox.

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()
);

Применете ја шемата преку постојна, експлицитно избрана врска со базата на податоци:

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

Користете TLS потврда соодветна за вашата околина. sslmode=require го шифрира транспортот, но не ја обезбедува истата потврда на идентитетот на серверот како verify-full со доверлив CA.

Имплементирајте го пулот на работници

Имплементацијата подолу користи закуп од 30 секунди, го обновува на секои 10 секунди, ги ограничува преземањата на две секунди и ги ограничува наредбите кон базата на податоци на три секунди. Овие граници се удобно под закупот. Временското ограничување на врската е конфигурирано одделно бидејќи не го ограничува извршувањето на прашањето.

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]
}

Зошто е важна трансакцијата на ефектот

processed_messages и ажурирањето за завршување на задачата се потврдуваат во една трансакција. Ако работникот повеќе не е сопственик на закупот, ажурирањето засега нула редови и трансакцијата се поништува, вклучувајќи го и секој нов ефект. Ако претходен обид го потврдил ефектот, но редот во редицата некако останал повторно достапен, ON CONFLICT DO NOTHING го прави повторното извршување безопасно.

Овој образец штити само ефекти во истата база на податоци. За надворешно плаќање или HTTP API, испратете го стабилниот клуч за идемпотентност до надолна услуга што го поддржува, или потврдете ред во трансакциски outbox и дозволете друг идемпотентен диспечер да го изврши мрежниот повик. Локален знак „обработено“ не може атомски да докаже дека се случил неповрзан оддалечен спореден ефект.

Тестирајте повторни обиди, закупи и исклучување

Започнете со еден работник за премините лесно да се прегледаат:

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

Од друга школка, ставете во редица една успешна задача и еден детерминистички неуспех:

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;"

За време на осумсекундната задача, испратете SIGTERM со Ctrl-C. Работникот ја откажува обработката и условно ја ослободува својата резервација. Рестартирајте го и потврдете дека важечката задача на крај успева со еден ред на ефект. Нејзиниот број на обиди може да надмине еден; тоа се очекува при испорака најмалку еднаш. Неважечката задача прави повторни обиди со ограничени експоненцијални одложувања и станува dead по третото преземање.

За да тестирате опоравување по пад наместо уредно ослободување, насилно прекинете го процесот во изолирана тест-околина. Редот останува running сè додека не истече неговиот закуп од 30 секунди, по што може да се преземе. Никогаш не користете насилен сигнал како вообичаен механизам за исклучување.

Набљудливост и перформанси

Работникот емитува структурирани JSON дневници што содржат ID на задачи, обиди, сопственици и грешки. Избегнувајте да ги евидентирате целосните payload-и: тие може да содржат акредитиви или лични податоци. Корисни мерења во базата на податоци вклучуваат број на подготвени задачи, старост на најстарата подготвена задача, активни закупи, истечени закупи, мртви задачи, стапка на повторни обиди, траење на обработката и број на загубени закупи.

Скалирајте ги работниците според измерениот капацитет на базата на податоци и надолните системи, а не само според бројот на процесори. Секој работник може да држи една врска за време на прашање, додека на отчукувањата на срцето и чистачот им треба резервен капацитет; затоа пулот дозволува workers + 2 врски. Одржувајте ги payload-ите на задачите мали, архивирајте ги старите успешни редови и проверувајте ги плановите на прашањата како што расте табелата. Делумните индекси ги држат скенирањата за преземање фокусирани, но не го заменуваат редовното одржување на PostgreSQL.

Безбедност и распоредување

Користете посветена сметка на оперативниот систем и улога во базата на податоци со најмали привилегии. На работникот не му е потребна влезна мрежна порта, затоа не отворајте таква во заштитниот ѕид на хостот. Ограничете го излезниот пристап на PostgreSQL и на вистинските надолни зависности таму каде што архитектурата на вашиот заштитен ѕид поддржува политика за излезен сообраќај. Чувајте ја URL-адресата на базата на податоци надвор од unit-датотеката, заштитете ја со root сопственост и режим 0600 и ротирајте ги акредитивите без да ги вградувате во бинарната датотека.

[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

На хостот, администратор може да ја изгради бинарната датотека, да ја постави под /opt/reliable-worker, да ја создаде сервисната сметка и заштитената датотека со околински променливи, да ја инсталира единицата, а потоа да изврши systemctl daemon-reload и systemctl enable --now reliable-worker. Осигурете се дека TimeoutStopSec останува поголем од ограниченото прашање за ослободување. Контејнерите треба да ги користат истите принципи на непривилегиран корисник и датотечен систем само за читање и мора да добијат доволно време за уредно прекинување пред присилно отстранување.

Вообичаени неуспеси во продукција

  • Закуп пократок од реалното време за обработка: обновувајте многу пред истекувањето и ограничете ја секоја операција. Долгите паузи сè уште може да доведат до губење сопственост, па завршувањето мора да остане оградено.
  • Повторување обиди додека се држи трансакција: прво преземете и потврдете. Извршувањето деловен код во трансакција за резервација создава заклучувања, изгладнување на врски и тешко опоравување.
  • Започнување работа при исклучување: секогаш проверувајте откажување веднаш по блокирачкото преземање. Ако преземањето се потврдило за време на исклучувањето, ослободете го со ажурирање што ја проверува сопственоста.
  • Претпоставка дека временското ограничување на врската ги ограничува прашањата: конфигурирајте ги роковите за врска, наредба, заклучување и контекст независно.
  • Нарекување на повторните обиди точно еднаш: закупите решаваат напуштање, не неизвесност. Идемпотентноста на границата на ефектот е тоа што ја прави повторената испорака безбедна.

Конечна контролна листа за потврда

  • Шемата и делумните индекси се применуваат без грешки.
  • Истовремените работници никогаш не го преземаат истиот активен закуп.
  • Уредното исклучување не започнува нова деловна работа по завршено преземање.
  • Задачата на паднат работник станува достапна за преземање по истекување на закупот.
  • Отчукувањата на срцето ја запираат обработката кога се губи сопственоста.
  • Неуспешните задачи се обидуваат повторно со одложување и на крај влегуваат во dead.
  • Повторената испорака создава еден потврден ефект по клуч за идемпотентност.
  • Временските ограничувања за базата на податоци и исклучувањето остануваат под буџетот на закупот.
  • Дневниците и прашањата кон редицата откриваат заостаток, неуспеси и здравје на закупите.

Робустен пул на работници не се дефинира според тоа колку брзо троши канал. Се дефинира според она што останува вистинито кога времето станува непријателско: резервациите се кратки, сопственоста истекува, исклучувањето е намерно, повторните обиди се ограничени, а ефектите толерираат повторување. Штом тие инваријанти се видливи во шемата и наметнати во секое ажурирање, неуспехот престанува да биде исклучителна патека. Станува обична промена на состојба што системот веќе знае како да ја преживее.

Портрет на автор на блогот

Mihajlo

Јас сум Михајло - развивач поттикнат од љубопитност, дисциплина и постојаната желба да создадам нешто значајно. Споделувам увиди, упатства и бесплатни услуги за да им помогнам на другите да ја поедностават својата работа и да растат во постојано развивачкиот свет на софтверот и вештачката интелигенција.