Vodiči

Go Stream Processing: Resilient Systems with Bounded Memory and Checkpoints

Obrada tokova u Go-u: otporni sustavi s ograničenom memorijom i kontrolnim točkama

Procesor toka postaje zanimljiv onoga trenutka kada „pročitaj zapis, zapiši zapis” više nije dovoljno. Ulazi mogu biti veći od RAM-a, procesi mogu umrijeti između zapisivanja, neispravni zapisi mogu zaustaviti napredak, a uredno gašenje može stići dok je rad u tijeku.

Ovaj vodič izrađuje produkcijski orijentiran Go procesor za nepromjenjive datoteke JSON Lines. Održava međuspremnike aplikacije ograničenima, trajno zapisuje pakete, bilježi kontrolnu točku samo za dovršeni ulaz i sigurno se oporavlja kada dođe do kvara između potvrđivanja izlaza i pomicanja kontrolne točke.

Procesor namjerno ostaje jednonitan. To čuva redoslijed izvora i čini kontrolnu točku najvišćim dovršenim neprekinutim pomakom bajtova. Konkurentnost može poboljšati transformacije koje intenzivno koriste CPU, ali također zahtijeva ograničeni međuspremnik za preslagivanje i može ostaviti kasnije zapise da čekaju iza jednog sporog zapisa.

Preduvjeti i jamstva

Potrebni su vam Go 1.22 ili noviji i Linux datotečni sustav čije ste ponašanje atomskog povezivanja, preimenovanja i sinkronizacije direktorija provjerili. Lokalni ext4 i XFS tipični su izbori za implementaciju; mrežni datotečni sustavi zahtijevaju izričito testiranje.

Ulazni ugovor je strog:

  • Izvor je obična, nepromjenjiva JSONL datoteka.
  • Svaki zapis završava novim retkom.
  • Jedan zapis ne može premašiti 1 MiB.
  • Izlazni direktorij pripada jednom zadatku obrade.

Procesor prije pokretanja izračunava sažetak izvora te u kontrolnu točku pohranjuje njegovu veličinu i SHA-256 sažetak. To pri pokretanju zahtijeva jedno potpuno čitanje, ali sprječava slučajan nastavak kontrolne točke nad drugim ulazom. Za vrlo velike izvore, nepromjenjivi identifikatori objekata ili manifesti mogu zamijeniti izračunavanje sažetka pri pokretanju.

Obrada iskreno ostaje najmanje jednom. Stvaranje trajnih segmenata je idempotentno, ali svaka nizvodna baza podataka, API ili objava poruke mora koristiti vlastiti ključ idempotentnosti, kao što je sažetak izvora uz source_offset.

Arhitektura i granica kvara

Međuspremnički čitač prihvaća jedan ograničeni zapis odjednom. Transformirani zapisi prikupljaju se dok se ne dosegne 500 zapisa ili 4 MiB kodiranog izlaza. Budući da čitanje prestaje dok se paket trajno pohranjuje, povratni pritisak je automatski, a memorija ne raste s veličinom izvora.

  1. Pročitajte i provjerite zapis.
  2. Transformirajte ga deterministički.
  3. Dodajte ga u ograničeni paket.
  4. Zapišite paket u privremenu datoteku i pozovite fsync.
  5. Atomski ga povežite s nazivom datoteke izvedenim iz pomaka izvora.
  6. Sinkronizirajte izlazni direktorij.
  7. Atomski zamijenite i sinkronizirajte kontrolnu točku.

Ako proces umre nakon petog, ali prije sedmog koraka, oporavak ponovno stvara isti naziv datoteke i sadržaj. Postojeći istovjetni segment prihvaća se. Drugačiji sadržaj pri istim pomacima tretira se kao oštećenje ili nedeterministička transformacija.

Struktura projekta

jsonstream/
├── cmd/
│   └── jsonstream/
│       └── main.go
├── deploy/
│   └── jsonstream.service
└── go.mod

Izradite go.mod bez ovisnosti trećih strana:

module example.com/jsonstream

go 1.22

Implementirajte procesor

Primjer ulaza predstavlja događaje računa. Transformacija provjerava identifikatore i klasificira svaki iznos, uz zadržavanje izvornog pomaka bajtova za nizvodnu deduplikaciju.

package main

import (
	"bufio"
	"bytes"
	"context"
	"crypto/sha256"
	"encoding/hex"
	"encoding/json"
	"errors"
	"flag"
	"fmt"
	"io"
	"log/slog"
	"os"
	"os/signal"
	"path/filepath"
	"syscall"
)

const (
	maxRecordBytes = 1 << 20
	maxBatchBytes  = 4 << 20
	maxBatchEvents = 500
)

type config struct {
	input, out, state string
	failAfterSegment bool
}

type event struct {
	ID          string `json:"id"`
	Account     string `json:"account"`
	AmountCents int64  `json:"amount_cents"`
}

type output struct {
	ID           string `json:"id"`
	Account      string `json:"account"`
	AmountCents  int64  `json:"amount_cents"`
	Class        string `json:"class"`
	SourceOffset int64  `json:"source_offset"`
}

type checkpoint struct {
	SourceSize   int64  `json:"source_size"`
	SourceSHA256 string `json:"source_sha256"`
	Offset       int64  `json:"offset"`
}

type batch struct {
	start, end int64
	count      int
	data       bytes.Buffer
}

func main() {
	var cfg config
	flag.StringVar(&cfg.input, "input", "events.jsonl", "immutable JSONL source")
	flag.StringVar(&cfg.out, "out", "run/out", "segment directory")
	flag.StringVar(&cfg.state, "state", "run/checkpoint.json", "checkpoint path")
	flag.BoolVar(&cfg.failAfterSegment, "fail-after-segment", false, "test recovery boundary")
	flag.Parse()

	logger := slog.New(slog.NewJSONHandler(os.Stderr, nil))
	ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
	defer stop()

	if err := run(ctx, cfg, logger); err != nil {
		logger.Error("stream_failed", "error", err)
		os.Exit(1)
	}
}

func run(ctx context.Context, cfg config, logger *slog.Logger) error {
	if err := os.MkdirAll(cfg.out, 0750); err != nil {
		return err
	}
	if err := os.MkdirAll(filepath.Dir(cfg.state), 0750); err != nil {
		return err
	}

	f, err := os.Open(cfg.input)
	if err != nil {
		return err
	}
	defer f.Close()

	before, err := f.Stat()
	if err != nil || !before.Mode().IsRegular() {
		return fmt.Errorf("input must be a regular file")
	}

	h := sha256.New()
	if _, err := io.Copy(h, f); err != nil {
		return fmt.Errorf("hash input: %w", err)
	}
	after, err := f.Stat()
	if err != nil {
		return err
	}
	if before.Size() != after.Size() || !before.ModTime().Equal(after.ModTime()) {
		return fmt.Errorf("input changed while hashing")
	}
	digest := hex.EncodeToString(h.Sum(nil))

	cp, found, err := loadCheckpoint(cfg.state)
	if err != nil {
		return err
	}
	if !found {
		cp = checkpoint{SourceSize: after.Size(), SourceSHA256: digest}
		if err := saveCheckpoint(cfg.state, cp); err != nil {
			return err
		}
	} else if cp.SourceSize != after.Size() || cp.SourceSHA256 != digest {
		return fmt.Errorf("checkpoint belongs to different input")
	}
	if cp.Offset < 0 || cp.Offset > cp.SourceSize {
		return fmt.Errorf("invalid checkpoint offset %d", cp.Offset)
	}
	if _, err := f.Seek(cp.Offset, io.SeekStart); err != nil {
		return err
	}

	reader := bufio.NewReaderSize(f, maxRecordBytes+1)
	offset := cp.Offset
	var current *batch

	flush := func() error {
		if current == nil {
			return nil
		}
		if err := commitBatch(current, cfg, &cp, logger); err != nil {
			return err
		}
		current = nil
		return nil
	}

	for {
		if ctx.Err() != nil {
			return flush()
		}

		line, readErr := reader.ReadSlice('\n')

		// A stop arriving during the blocking read must not start new work.
		if ctx.Err() != nil {
			return flush()
		}
		if errors.Is(readErr, bufio.ErrBufferFull) {
			return fmt.Errorf("record at offset %d exceeds %d bytes", offset, maxRecordBytes)
		}
		if errors.Is(readErr, io.EOF) {
			if len(line) != 0 {
				return fmt.Errorf("final record at offset %d lacks newline", offset)
			}
			if err := flush(); err != nil {
				return err
			}
			logger.Info("stream_complete", "offset", cp.Offset, "source_size", cp.SourceSize)
			return nil
		}
		if readErr != nil {
			return fmt.Errorf("read at offset %d: %w", offset, readErr)
		}
		if len(line) > maxRecordBytes {
			return fmt.Errorf("record at offset %d exceeds %d bytes", offset, maxRecordBytes)
		}

		encoded, err := transform(line, offset)
		if err != nil {
			return fmt.Errorf("record at offset %d: %w", offset, err)
		}
		if len(encoded) > maxBatchBytes {
			return fmt.Errorf("encoded record at offset %d exceeds batch limit", offset)
		}

		if current != nil &&
			(current.count == maxBatchEvents ||
				current.data.Len()+len(encoded) > maxBatchBytes) {
			if err := flush(); err != nil {
				return err
			}
		}
		if current == nil {
			current = &batch{start: offset}
		}

		current.data.Write(encoded)
		current.count++
		offset += int64(len(line))
		current.end = offset
	}
}

func transform(line []byte, offset int64) ([]byte, error) {
	var in event
	decoder := json.NewDecoder(bytes.NewReader(line))
	decoder.DisallowUnknownFields()
	if err := decoder.Decode(&in); err != nil {
		return nil, err
	}
	if err := decoder.Decode(&struct{}{}); !errors.Is(err, io.EOF) {
		return nil, fmt.Errorf("record contains trailing JSON")
	}
	if in.ID == "" || in.Account == "" {
		return nil, fmt.Errorf("id and account are required")
	}

	class := "zero"
	if in.AmountCents > 0 {
		class = "credit"
	} else if in.AmountCents < 0 {
		class = "debit"
	}

	encoded, err := json.Marshal(output{
		ID: in.ID, Account: in.Account, AmountCents: in.AmountCents,
		Class: class, SourceOffset: offset,
	})
	if err != nil {
		return nil, err
	}
	return append(encoded, '\n'), nil
}

func commitBatch(b *batch, cfg config, cp *checkpoint, logger *slog.Logger) error {
	if err := writeSegment(cfg.out, b); err != nil {
		return err
	}
	if cfg.failAfterSegment {
		return fmt.Errorf("injected failure after durable segment")
	}

	next := *cp
	next.Offset = b.end
	if err := saveCheckpoint(cfg.state, next); err != nil {
		return err
	}
	*cp = next
	logger.Info("batch_committed",
		"start_offset", b.start, "end_offset", b.end,
		"records", b.count, "bytes", b.data.Len())
	return nil
}

func writeSegment(outDir string, b *batch) error {
	final := filepath.Join(outDir,
		fmt.Sprintf("%020d-%020d.jsonl", b.start, b.end))

	tmp, err := os.CreateTemp(outDir, ".batch-*")
	if err != nil {
		return err
	}
	tmpName := tmp.Name()
	defer os.Remove(tmpName)

	if err := tmp.Chmod(0640); err != nil {
		tmp.Close()
		return err
	}
	n, err := tmp.Write(b.data.Bytes())
	if err != nil || n != b.data.Len() {
		tmp.Close()
		return fmt.Errorf("write temporary segment: %w", err)
	}
	if err := tmp.Sync(); err != nil {
		tmp.Close()
		return err
	}
	if err := tmp.Close(); err != nil {
		return err
	}

	if err := os.Link(tmpName, final); err != nil {
		if !errors.Is(err, os.ErrExist) {
			return fmt.Errorf("publish segment: %w", err)
		}
		same, compareErr := sameContent(final, b.data.Bytes())
		if compareErr != nil {
			return compareErr
		}
		if !same {
			return fmt.Errorf("existing segment %s has different content", final)
		}
	}
	if err := os.Remove(tmpName); err != nil && !errors.Is(err, os.ErrNotExist) {
		return err
	}
	return syncDir(outDir)
}

func sameContent(path string, expected []byte) (bool, error) {
	f, err := os.Open(path)
	if err != nil {
		return false, err
	}
	defer f.Close()

	info, err := f.Stat()
	if err != nil || !info.Mode().IsRegular() {
		return false, fmt.Errorf("existing segment is not a regular file")
	}
	if info.Size() != int64(len(expected)) {
		return false, nil
	}
	actualHash := sha256.New()
	if _, err := io.Copy(actualHash, f); err != nil {
		return false, err
	}
	expectedHash := sha256.Sum256(expected)
	return bytes.Equal(actualHash.Sum(nil), expectedHash[:]), nil
}

func loadCheckpoint(path string) (checkpoint, bool, error) {
	f, err := os.Open(path)
	if errors.Is(err, os.ErrNotExist) {
		return checkpoint{}, false, nil
	}
	if err != nil {
		return checkpoint{}, false, err
	}
	defer f.Close()

	var cp checkpoint
	decoder := json.NewDecoder(io.LimitReader(f, 64<<10))
	decoder.DisallowUnknownFields()
	if err := decoder.Decode(&cp); err != nil {
		return cp, false, fmt.Errorf("decode checkpoint: %w", err)
	}
	if err := decoder.Decode(&struct{}{}); !errors.Is(err, io.EOF) {
		return cp, false, fmt.Errorf("checkpoint contains trailing data")
	}
	return cp, true, nil
}

func saveCheckpoint(path string, cp checkpoint) error {
	data, err := json.Marshal(cp)
	if err != nil {
		return err
	}
	data = append(data, '\n')
	dir := filepath.Dir(path)

	tmp, err := os.CreateTemp(dir, ".checkpoint-*")
	if err != nil {
		return err
	}
	tmpName := tmp.Name()
	defer os.Remove(tmpName)

	if err := tmp.Chmod(0600); err != nil {
		tmp.Close()
		return err
	}
	if _, err := tmp.Write(data); err != nil {
		tmp.Close()
		return err
	}
	if err := tmp.Sync(); err != nil {
		tmp.Close()
		return err
	}
	if err := tmp.Close(); err != nil {
		return err
	}
	if err := os.Rename(tmpName, path); err != nil {
		return err
	}
	return syncDir(dir)
}

func syncDir(path string) error {
	dir, err := os.Open(path)
	if err != nil {
		return err
	}
	defer dir.Close()
	return dir.Sync()
}

Kontrolna točka napreduje tek nakon što je odgovarajući segment trajan. Pri SIGTERM-u petlja ponovno provjerava otkazivanje neposredno nakon čitanja, odbija transformirati novopročitani zapis i potvrđuje samo paket koji je već bio u tijeku.

Izradite i testirajte uspješnu putanju

mkdir -p sample/run/out
cat > sample/events.jsonl <<'EOF'
{"id":"evt-001","account":"alpha","amount_cents":1250}
{"id":"evt-002","account":"alpha","amount_cents":-300}
{"id":"evt-003","account":"beta","amount_cents":0}
EOF

go test ./...
go build -trimpath -o jsonstream ./cmd/jsonstream
./jsonstream \
  -input sample/events.jsonl \
  -out sample/run/out \
  -state sample/run/checkpoint.json

cat sample/run/out/*.jsonl
cat sample/run/checkpoint.json

Izlaz bi trebao sadržavati tri klasificirana zapisa. Ponovno pokretanje iste naredbe čita EOF na kontrolnoj točki i ne stvara dodatni segment.

Isprobajte oporavak nakon kvara

Zastavica za test umeće kvar nakon objavljivanja segmenta, ali prije ažuriranja kontrolne točke. Modelira trajno stanje koje ostaje nakon rušenja na najvažnijoj granici.

mkdir -p sample/recovery/out

if ./jsonstream \
  -input sample/events.jsonl \
  -out sample/recovery/out \
  -state sample/recovery/checkpoint.json \
  -fail-after-segment
then
  echo "expected injected failure" >&2
  exit 1
fi

segment=$(find sample/recovery/out -maxdepth 1 -type f -name '*.jsonl')
before=$(sha256sum "$segment" | cut -d' ' -f1)

./jsonstream \
  -input sample/events.jsonl \
  -out sample/recovery/out \
  -state sample/recovery/checkpoint.json

after=$(sha256sum "$segment" | cut -d' ' -f1)
test "$before" = "$after"
test "$(find sample/recovery/out -maxdepth 1 -type f -name '*.jsonl' | wc -l)" -eq 1
cat sample/recovery/checkpoint.json

Oporavak nailazi na postojeći segment, provjerava njegovu veličinu i sažetak te zatim pomiče kontrolnu točku. Neusklađeni segment zaustavlja obradu umjesto da tiho prepiše dokaze.

Implementirajte kao ojačanu Linux uslugu

Izradite deploy/jsonstream.service sa sljedećom konfiguracijom:

[Unit]
Description=Bounded JSON stream processor
After=local-fs.target
StartLimitIntervalSec=300
StartLimitBurst=3

[Service]
Type=exec
DynamicUser=yes
StateDirectory=jsonstream
StateDirectoryMode=0750
ExecStart=/usr/local/libexec/jsonstream \
  -input /srv/jsonstream/events.jsonl \
  -out /var/lib/jsonstream/out \
  -state /var/lib/jsonstream/checkpoint.json
Restart=on-failure
RestartSec=5s
TimeoutStopSec=30s
UMask=0027
Environment=GOMEMLIMIT=96MiB
MemoryMax=128M

NoNewPrivileges=yes
PrivateTmp=yes
PrivateDevices=yes
ProtectSystem=strict
ProtectHome=yes
ProtectKernelTunables=yes
ProtectKernelModules=yes
ProtectControlGroups=yes
RestrictSUIDSGID=yes
LockPersonality=yes
CapabilityBoundingSet=
RestrictAddressFamilies=AF_UNIX
ReadOnlyPaths=/srv/jsonstream/events.jsonl

StandardOutput=journal
StandardError=journal

[Install]
WantedBy=multi-user.target

Ovo su naredbe za poslužitelj i zahtijevaju administrativne ovlasti. Prije pokretanja usluge pripremite stvarni izvor; nemojte ga poslije mijenjati na mjestu.

sudo install -d -o root -g root -m 0755 /usr/local/libexec
sudo install -d -o root -g root -m 0755 /srv/jsonstream
sudo install -o root -g root -m 0755 jsonstream /usr/local/libexec/jsonstream
sudo install -o root -g root -m 0644 sample/events.jsonl /srv/jsonstream/events.jsonl
sudo install -o root -g root -m 0644 \
  deploy/jsonstream.service /etc/systemd/system/jsonstream.service

sudo systemctl daemon-reload
sudo systemctl enable --now jsonstream.service
sudo systemctl status jsonstream.service
sudo journalctl -u jsonstream.service --no-pager

Usluga ne otvara mrežne utičnice, stoga nije potrebno pravilo vatrozida. Izostavljanje mrežnog pristupa iz ovog radnika također sužava njegovu površinu kvara i sigurnosnu površinu. Dinamički korisnik može pisati samo unutar direktorija stanja kojim upravlja systemd.

Promatranje i performanse

Svaki potvrđeni paket u strukturiranom JSON-u bilježi pomake, broj zapisa i izlazne bajtove. Događaj dovršetka prijavljuje završni pomak i veličinu izvora. Operativni nadzor trebao bi upozoravati na ponavljane događaje stream_failed i mjeriti zaostatak kontrolne točke kao source_size - offset.

Veličina paketa središnji je kompromis između trajnosti i propusnosti. Veći paketi amortiziraju stvaranje datoteka, izračunavanje sažetka, zamjenu kontrolne točke i fsync, ali povećavaju posao ponavljanja i latenciju vidljivosti. Manji paketi čine obrnuto. Mjerite s reprezentativnim veličinama zapisa i pohranom, uključujući prisilna ponovna pokretanja.

Ograničenja paketa od 4 MiB i zapisa od 1 MiB ograničavaju međuspremnik korisnog tereta aplikacije, a ne cijeli Go proces. Nizovi dekodera, metapodaci izvršnog okruženja, međuspremnici datotečnog sustava i skupljanje smeća zahtijevaju dodatnu memoriju. GOMEMLIMIT usmjerava izvršno okruženje, dok systemd-ov MemoryMax pruža čvrstu granicu ograničenja.

Uobičajeni produkcijski kvarovi

  • Nepodudaranje sažetka izvora: ulaz se promijenio ili je odabrana pogrešna kontrolna točka. Pokrenite novi zadatak s novim direktorijem stanja i izlaza.
  • Otrovni zapis: neispravan JSON, nepoznata polja ili nedostajući identifikatori zaustavljaju se na stabilnom pomaku bajtova. Provjerite podatke prije pripreme; nemojte nevidljivo preskakati zapise.
  • Nepodudaranje segmenta: kôd se promijenio, izlaz je uređen ili su direktoriji ponovno korišteni u više zadataka. Sačuvajte datoteke radi istrage.
  • Kvar dozvole ili sinkronizacije: potvrdite vlasništvo, opcije montiranja, slobodan prostor i dostupnost inodeova. Neuspjele pozive trajnosti tretirajte kao neuspjele pakete.
  • Petlja ponovnog pokretanja: determinističke pogreške ulaza ponavljat će se. Ograničenje pokretanja systemd-a sprječava neograničenu usku petlju, ali operateri i dalje trebaju upozorenje.
  • Prisilno gašenje: SIGKILL može prekinuti trajno spremanje, ali usklađivanje pri ponovnom pokretanju ostaje sigurno. Neka proračun za gašenje bude dulji od uočenog najgoreg vremena sinkronizacije.

Završni kontrolni popis za provjeru

  • Izvor je nepromjenjiv, završava novim retkom i pohranjen je odvojeno od izlaza.
  • Ograničenja memorije uzimaju u obzir dodatno opterećenje izvršnog okruženja, kao i konstante paketa.
  • Uobičajeno ponovno pokretanje ne stvara duplicirani segment.
  • Umetnuti kvar nakon segmenta oporavlja se bez promjene sadržaja segmenta.
  • SIGTERM zaustavlja nove transformacije i potvrđuje samo prethodni rad u tijeku.
  • Dnevnici izlažu potvrđene pomake, kvarove i dovršetak.
  • Nizvodni učinci koriste ključ idempotentnosti i pretpostavljaju isporuku najmanje jednom.
  • Testirano je ponašanje poveznice, preimenovanja i sinkronizacije direktorija implementiranog datotečnog sustava.

Pouzdana obrada toka manje se odnosi na impresivnu petlju, a više na odabir jedne precizne granice: najprije izlaz, zatim kontrolna točka, uz deterministički oporavak između njih. Kada je ta granica trajna, ograničena, vidljiva i testabilna, rušenja prestaju biti zagonetni događaji. Postaju još jedan ulaz koji sustav već zna obraditi.

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.