Обработка на текови во Go: Отпорни системи со ограничена меморија и контролни точки
Процесорот на текови станува интересен во моментот кога „прочитај запис, запиши запис“ веќе не е доволно. Влезовите може да бидат поголеми од RAM, процесите може да прекинат меѓу запишувања, невалидните записи може да го запрат напредокот, а уредно гасење може да пристигне додека работата е во тек.
Овој туторијал создава производно ориентиран Go процесор за непроменливи JSON Lines датотеки. Тој ги одржува ограничени апликациските бафери, трајно запишува серии, прави контролни точки само за завршениот влез и безбедно се опоравува кога ќе дојде до неуспех меѓу потврдувањето на излезот и унапредувањето на контролната точка.
Процесорот намерно останува еднонишковен. Тоа го зачувува редоследот на изворот и ја прави контролната точка највисоката завршена непрекината бајтова поместеност. Паралелноста може да ги подобри трансформациите со интензивна употреба на CPU, но бара и ограничен бафер за преуредување и може да остави подоцнежни записи да чекаат зад еден бавен запис.
Предуслови и гаранции
Потребни ви се Go 1.22 или понов и Linux датотечен систем чиe атомско поврзување, преименување и однесување при синхронизација на директориуми сте ги потврдиле. Локалните ext4 и XFS се вообичаени избори за распоредување; мрежните датотечни системи бараат изречно тестирање.
Договорот за влезот е строг:
- Изворот е обична, непроменлива JSONL датотека.
- Секој запис завршува со нов ред.
- Еден запис не може да надмине 1 MiB.
- Излезниот директориум припаѓа на една задача за обработка.
Процесорот го хешира изворот пред почетокот и ги зачувува неговата големина и SHA-256 отпечаток во контролната точка. Ова чини едно целосно читање при стартување, но спречува случајно продолжување на контролна точка со различен влез. За многу големи извори, непроменливите идентификатори на објекти или манифести може да го заменат хеширањето при стартување.
Обработката искрено останува барем еднаш. Создавањето трајни сегменти е идемпотентно, но секое објавување во надворешна база на податоци, API или порака мора да користи сопствен клуч за идемпотентност, како што се отпечатокот на изворот плус source_offset.
Архитектура и граница на неуспех
Баферизиран читач прифаќа по еден ограничен запис. Трансформираните записи се собираат додека не се достигнат или 500 записи или 4 MiB кодиран излез. Бидејќи читањето запира додека се зачувува серија, повратниот притисок е автоматски и меморијата не расте со големината на изворот.
- Прочитајте и проверете запис.
- Трансформирајте го детерминистички.
- Додајте го во ограничената серија.
- Запишете ја серијата во привремена датотека и повикајте
fsync. - Атомски поврзете ја со име на датотека изведено од изворните поместувања.
- Синхронизирајте го излезниот директориум.
- Атомски заменете ја и синхронизирајте ја контролната точка.
Ако процесот прекине по чекор пет, но пред чекор седум, опоравувањето го создава истото име на датотека и содржина. Постоечки идентичен сегмент се прифаќа. Различна содржина на истите поместувања се смета за оштетување или недетерминистичка трансформација.
Структура на проектот
jsonstream/
├── cmd/
│ └── jsonstream/
│ └── main.go
├── deploy/
│ └── jsonstream.service
└── go.mod
Создадете go.mod без зависности од трети страни:
module example.com/jsonstream
go 1.22
Имплементирајте го процесорот
Примерот за влез претставува настани на сметки. Трансформацијата ги проверува идентификаторите и го класифицира секој износ, притоа задржувајќи ја оригиналната бајтова поместеност за дедупликација надолу по текот.
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()
}
Контролната точка напредува само откако соодветниот сегмент е траен. При SIGTERM, јамката повторно ја проверува откажаноста веднаш по читањето, одбива да го трансформира новопрочитаниот запис и ја потврдува само серијата што веќе била во тек.
Изградете и тестирајте ја успешната патека
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
Излезот треба да содржи три класифицирани записи. Повторното извршување на истата команда чита EOF на контролната точка и не создава дополнителен сегмент.
Испробајте опоравување по неуспех
Тест-знаменцето вметнува неуспех по објавувањето на сегментот, но пред ажурирањето на контролната точка. Тоа ја моделира трајната состојба што ја остава пад на најважната граница.
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
Опоравувањето се среќава со постоечкиот сегмент, ги проверува неговата големина и хеш, а потоа ја унапредува контролната точка. Несовпаѓачки сегмент ја запира обработката наместо тивко да ги презапише доказите.
Распоредете како зајакната Linux услуга
Создадете deploy/jsonstream.service со следнава конфигурација:
[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
Ова се команди за домаќинот и бараат административни привилегии. Поставете го вистинскиот извор пред да ја стартувате услугата; не го менувајте на место потоа.
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
Услугата не отвора мрежни приклучници, па не е потребно правило за заштитен ѕид. Исклучувањето на мрежниот пристап од овој работник исто така ја стеснува неговата површина за неуспех и безбедност. Динамичкиот корисник може да запишува само под директориумот за состојба управуван од systemd.
Набљудливост и перформанси
Секоја потврдена серија ги запишува поместувањата, бројот на записи и излезните бајти како структуриран JSON. Настанот за завршување ја пријавува конечната поместеност и големината на изворот. Оперативното следење треба да алармира за повторени настани stream_failed и да го мери заостанувањето на контролната точка како source_size - offset.
Големината на серијата е централниот компромис меѓу трајност и пропусност. Поголемите серии го амортизираат создавањето датотеки, хеширањето, замената на контролната точка и fsync, но ја зголемуваат работата при повторување и латентноста на видливоста. Помалите серии го прават спротивното. Мерете со репрезентативни големини на записи и складиште, вклучително и со присилни рестартирања.
Ограничувањата од 4 MiB за серија и 1 MiB за запис го ограничуваат баферирањето на товарот на апликацијата, а не целиот Go процес. Низите на декодерот, метаподатоците на извршувачката околина, баферите на датотечниот систем и собирањето отпад бараат дополнителна меморија. GOMEMLIMIT ја насочува извршувачката околина, додека MemoryMax на systemd ја обезбедува тврдата граница за ограничување.
Вообичаени производни неуспеси
- Несовпаѓање на хешот на изворот: влезот се променил или е избрана погрешна контролна точка. Започнете нова задача со нова состојба и излезен директориум.
- Токсичен запис: невалиден JSON, непознати полиња или идентификатори што недостигаат запираат на стабилна бајтова поместеност. Проверете ги податоците пред поставување; не прескокнувајте записи невидливо.
- Несовпаѓање на сегмент: кодот се променил, излезот бил уреден или директориумите биле повторно употребени меѓу задачи. Зачувајте ги датотеките за истрага.
- Неуспех со дозволи или синхронизација: потврдете ги сопственоста, опциите за монтирање, слободниот простор и достапноста на inode-и. Третирајте ги неуспешните повици за трајност како неуспешни серии.
- Јамка на рестартирање: детерминистичките грешки во влезот ќе се повторат. Ограничувањето на стартување на systemd спречува неограничена тесна јамка, но на операторите сè уште им е потребен аларм.
- Присилно гасење: SIGKILL може да го прекине зачувувањето, но усогласувањето по рестарт останува безбедно. Одржувајте го буџетот за гасење подолг од забележаното најлошо време за синхронизација.
Контролна листа за конечна проверка
- Изворот е непроменлив, завршува со нов ред и се чува одделно од излезот.
- Ограничувањата на меморијата го земаат предвид трошокот на извршувачката околина, како и константите за серијата.
- Нормалното повторно извршување не создава дупликат сегмент.
- Вметнатиот неуспех по сегментот се опоравува без да ја промени содржината на сегментот.
- SIGTERM запира нови трансформации и потврдува само претходната работа во тек.
- Дневниците ги изложуваат потврдените поместувања, неуспесите и завршувањето.
- Ефектите надолу по текот користат клуч за идемпотентност и претпоставуваат испорака барем еднаш.
- Тестирано е однесувањето на датотечниот систем при распоредување за поврзување, преименување и синхронизација на директориуми.
Сигурната обработка на текови е помалку прашање на импресивна јамка, а повеќе на избор на една прецизна граница: прво излезот, потоа контролната точка, со детерминистичко опоравување меѓу нив. Штом таа граница е трајна, ограничена, набљудлива и тестирувачка, падовите престануваат да бидат мистериозни настани. Тие стануваат уште еден влез што системот веќе знае како да го обработи.