batchx

package module
v0.1.1 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 22, 2026 License: MIT Imports: 20 Imported by: 0

README

batchx

ci Go Reference

Chunk-oriented batch processing for Go, in the spirit of Spring Batch. Generic reader, processor and writer, jobs made of steps, restart from the last committed chunk, skip and retry policies, listeners, and a pluggable job repository. Standard library only, no CGO. Requires Go 1.24+.

go get github.com/JiaBao-do/batchx

Quick start

// Quickstart: read CSV, process each row, write JSON lines, in chunks of 2.
// Run: go run ./examples/quickstart
package main

import (
	"context"
	"log"
	"os"
	"strings"

	"github.com/JiaBao-do/batchx"
)

type User struct {
	ID   string `json:"id"`
	Name string `json:"name"`
}

func main() {
	in := strings.NewReader("id,name\n1,ann\n2,bob\n3,cy\n")
	toUser := batchx.ProcessorFunc[[]string, User](func(_ context.Context, r []string) (User, error) {
		return User{ID: r[0], Name: strings.ToUpper(r[1])}, nil
	})
	step, err := batchx.NewStep("import", batchx.NewCSVReader(in, true), toUser,
		batchx.NewJSONLWriter[User](os.Stdout), batchx.WithChunkSize(2))
	if err != nil {
		log.Fatal(err)
	}
	if _, err := batchx.NewJob("quickstart", nil).Then(step).Run(context.Background(), nil); err != nil {
		log.Fatal(err)
	}
}

This is examples/quickstart (a test keeps them identical). If the process dies or ctx is cancelled, run the same job with the same params again: completed steps are skipped and the failed step resumes from its last committed checkpoint. A completed instance returns ErrAlreadyCompleted.

Examples

Runnable programs, each with an expected_output.txt that CI compares against:

Program Shows
go run ./examples/quickstart CSV to processor to JSON lines, chunked
go run ./examples/restart crash mid-run, rerun, resume from the checkpoint
go run ./examples/skipretry skipping bad rows, retry with backoff, listeners
go run ./examples/filerepo the JSON job state file before and after a restart
go run ./examples/demo multi-step job: parallel loads, join, retry, skip, report

Things to care about

Full list with wrong/right snippets, each backed by a test: docs/PITFALLS.md. The important ones:

  • At-least-once. A crash after Write but before the checkpoint re-writes that chunk. Make writers idempotent (upsert by key).
  • Readers must replay the same order on restart, and non-seekable readers (lines, CSV, JSONL, channels) must be created fresh for every run, or the restart silently skips data.
  • Retry retries skippable errors by default. Set Retry.If so bad rows are not retried.
  • Skip limit is per job instance across restarts; default 0 means any error fails the step.
  • Chunk size is memory plus one repository save per chunk; the file repository fsyncs on each save.
  • Parallel steps run concurrently; share nothing without your own locking.
  • No cross-process locking in FileRepository; step names are checkpoint keys (renaming restarts the step).

Concepts

Spring Batch batchx
ItemReader / ItemProcessor / ItemWriter Reader[I], Processor[I,O], Writer[O]
Step (chunk / tasklet) NewStep, NewCopyStep, NewTaskletStep
Job, JobParameters Job (Then, ThenParallel), Params
JobRepository Repository interface (two methods), MemoryRepository, FileRepository
Skip / retry policy WithSkipLimit, WithSkipIf, Skippable(err), Retry
Listeners Listener (struct of optional funcs)

Guarantees and limits:

  • Delivery is at-least-once per chunk. A crash after Write returns but before the checkpoint is saved re-writes that chunk on restart; make writers idempotent. A chunk Write should be all-or-nothing.
  • When a skippable error fails a chunk write, batchx re-writes the chunk item by item to isolate the bad item.
  • Restart checkpoint is a position count. Readers implementing Seeker restart in O(1); others are replayed (read and discarded), so the input must be deterministic between runs.
  • FileRepository is atomic per write but does not lock across processes; run one instance per job key.
  • Not in v0.1: partitioning of one step across workers, multi-process coordination, SQL repository (implement the two-method Repository).

Comparison with gobatch

Checked 2026-09-21 from github.com/chararch/gobatch source and metadata (I did not run it): gobatch (MIT, 56 stars, last push 2025-02-27) uses interface{} item types with reflect, go 1.16, and depends on go-sql-driver/mysql, ants, pkg/errors and an FTP client. Its job repository is package-level functions over a SQL database and the repository ships only a MySQL schema. batchx differs by being generic, dependency-free and storage-agnostic (memory and JSON file built in). gobatch has features batchx lacks: SQL persistence, file/FTP components, TSV, task pools. See STATUS.md for the evidence log.

Development

git config core.hooksPath .githooks   # pre-push runs fmt, tidy, vet, race tests, build, lint

License

MIT

Documentation

Overview

Package batchx is a chunk-oriented batch processing framework in the spirit of Spring Batch, written with generics and only the standard library.

A Job is a list of Steps. A chunk step reads items from a Reader, transforms them with a Processor and hands them to a Writer in chunks, then commits a checkpoint to a Repository. If the process stops, running the same Job with the same Params resumes from the last committed chunk.

Skip limits with typed skippable errors, retry with backoff, listeners, context cancellation and parallel steps are supported. Repositories are pluggable: MemoryRepository and a JSON FileRepository ship in the box.

Delivery guarantee: at-least-once per chunk. A crash between a successful Writer.Write and the checkpoint save re-writes that chunk on restart, so writers should be idempotent.

Example

A chunk step copies items with a processor, and the job records progress.

package main

import (
	"context"
	"fmt"
	"strings"

	"github.com/JiaBao-do/batchx"
)

func main() {
	var out batchx.SliceWriter[string]
	upper := batchx.ProcessorFunc[string, string](func(_ context.Context, s string) (string, error) {
		return strings.ToUpper(s), nil
	})
	step, _ := batchx.NewStep("shout", batchx.NewSliceReader([]string{"a", "b", "c"}), upper, &out, batchx.WithChunkSize(2))
	rec, err := batchx.NewJob("demo", nil).Then(step).Run(context.Background(), nil)
	s, _ := rec.Step("shout")
	fmt.Println(err, rec.Status, out.Items(), s.Commits)
}
Output:
<nil> completed [A B C] 2

Index

Examples

Constants

This section is empty.

Variables

View Source
var (
	// ErrFilter is returned by a Processor to drop an item without failing.
	// Filtered items are counted in StepRecord.Filtered and never written.
	ErrFilter = errors.New("batchx: item filtered")

	// ErrSkipLimitExceeded wraps the error that pushed a step past its skip limit.
	ErrSkipLimitExceeded = errors.New("batchx: skip limit exceeded")

	// ErrAlreadyCompleted is returned by Job.Run when a job instance with the
	// same name and parameters already completed.
	ErrAlreadyCompleted = errors.New("batchx: job instance already completed")

	// ErrAlreadyRunning is returned by Job.Run when the same job instance is
	// already running in this process.
	ErrAlreadyRunning = errors.New("batchx: job instance already running")

	// ErrNotFound is returned by Repository.Load when no record exists.
	ErrNotFound = errors.New("batchx: job record not found")

	// ErrCheckpointBeyondEnd is returned when a restart checkpoint lies past
	// the end of the input, meaning the input shrank between runs.
	ErrCheckpointBeyondEnd = errors.New("batchx: checkpoint beyond end of input")

	// ErrInvalidConfig is returned for invalid step or job configuration.
	ErrInvalidConfig = errors.New("batchx: invalid configuration")
)

Functions

func ConstantBackoff

func ConstantBackoff(d time.Duration) func(int) time.Duration

ConstantBackoff waits d between attempts.

func ExponentialBackoff

func ExponentialBackoff(base, max time.Duration) func(int) time.Duration

ExponentialBackoff doubles the wait from base each attempt, capped at max (max <= 0 means no cap). There is no jitter, keeping behaviour deterministic.

func IsSkippable

func IsSkippable(err error) bool

IsSkippable is the default skip predicate: it reports whether err wraps a *SkippableError.

func Skippable

func Skippable(err error) error

Skippable wraps err so the default skip predicate accepts it. It returns nil for a nil err.

Types

type CSVReader

type CSVReader struct {
	// contains filtered or unexported fields
}

CSVReader reads records with encoding/csv. Malformed records (csv.ParseError, including wrong field counts) are returned as Skippable errors and the reader continues with the next record. It is not safe for concurrent use.

func NewCSVReader

func NewCSVReader(r io.Reader, hasHeader bool, configure ...func(*csv.Reader)) *CSVReader

NewCSVReader returns a reader over r. If hasHeader is true the first record is consumed on the first Read and available from Header. configure may adjust the underlying *csv.Reader (Comma, Comment, LazyQuotes, ...) and is called once.

Example

CSV in, JSON lines out.

package main

import (
	"context"
	"fmt"
	"strings"

	"github.com/JiaBao-do/batchx"
)

func main() {
	r := batchx.NewCSVReader(strings.NewReader("id,name\n1,ann\n2,bob\n"), true)
	rec, _ := r.Read(context.Background())
	fmt.Println(r.Header(), rec)
}
Output:
[id name] [1 ann]

func (*CSVReader) Header

func (c *CSVReader) Header() []string

Header returns the header record, valid after the first Read.

func (*CSVReader) Read

func (c *CSVReader) Read(ctx context.Context) ([]string, error)

Read implements Reader.

type CSVWriter

type CSVWriter struct {
	// contains filtered or unexported fields
}

CSVWriter writes records with encoding/csv. Each chunk is encoded into an in-memory buffer first; only once the whole chunk has encoded successfully does CSVWriter issue a single Write call to the underlying io.Writer with the complete, already-formed chunk bytes.

This matters because encoding/csv.Writer wraps a bufio.Writer (4KB by default) that auto-flushes to its destination as soon as the buffer fills, independent of record boundaries. If CSVWriter encoded directly into the caller's io.Writer, a chunk that failed partway through (an unencodable value, or the underlying writer rejecting a later write) could already have leaked a row split mid-record into the sink, even though the chunk as a whole is reported as failed and gets retried -- so a retry would append the resumed data after a truncated, malformed row instead of the clean "the chunk's rows appear a second time" duplicate that at-least-once retry semantics promise. Buffering the whole chunk first means the real sink sees exactly one Write call per chunk, carrying either all of the chunk's rows or (if the encode itself fails) none of them.

This guarantee assumes the underlying io.Writer's Write call is itself atomic (it either accepts the full slice or accepts none of it before returning an error), which holds for the common cases (files, in-memory buffers, most network writers for chunk-sized payloads). A sink that performs a genuine short write -- returning n < len(p) together with an error partway through consuming that one slice -- can still end up with a truncated tail; that is a property of the sink itself and is outside what any single io.Writer.Write call can guarantee.

It is not safe for concurrent use.

func NewCSVWriter

func NewCSVWriter(w io.Writer, configure ...func(*csv.Writer)) *CSVWriter

NewCSVWriter returns a writer to w. configure may adjust the *csv.Writer (Comma, UseCRLF) used to encode each chunk.

func (*CSVWriter) Write

func (c *CSVWriter) Write(_ context.Context, recs [][]string) error

Write implements Writer. It encodes the whole chunk into an in-memory buffer and, only on success, writes that buffer to the underlying io.Writer in a single call. See the CSVWriter doc comment for why.

type ChanReader

type ChanReader[T any] struct {
	// contains filtered or unexported fields
}

ChanReader reads from a channel until it is closed. Restart replays by discarding, so the channel producer must be deterministic to restart.

func NewChanReader

func NewChanReader[T any](ch <-chan T) *ChanReader[T]

NewChanReader returns a reader over ch.

func (*ChanReader[T]) Read

func (c *ChanReader[T]) Read(ctx context.Context) (T, error)

Read implements Reader; it returns ctx.Err() if ctx ends first.

type FileRepository

type FileRepository struct {
	// contains filtered or unexported fields
}

FileRepository stores one JSON file per job instance in a directory. Writes are atomic (temp file, fsync, rename). It is safe for concurrent use within one process; it does not lock across processes.

func NewFileRepository

func NewFileRepository(dir string) (*FileRepository, error)

NewFileRepository returns a FileRepository rooted at dir, creating it if needed.

func (*FileRepository) Load

func (f *FileRepository) Load(_ context.Context, key string) (JobRecord, error)

Load implements Repository.

func (*FileRepository) Save

func (f *FileRepository) Save(_ context.Context, rec JobRecord) error

Save implements Repository.

type JSONLReader

type JSONLReader[T any] struct {
	// contains filtered or unexported fields
}

JSONLReader reads one JSON value per line into T. Blank lines are ignored; invalid lines are Skippable errors. Lines have no length limit. It is not safe for concurrent use.

func NewJSONLReader

func NewJSONLReader[T any](r io.Reader) *JSONLReader[T]

NewJSONLReader returns a reader over the JSON lines of r.

func (*JSONLReader[T]) Read

func (j *JSONLReader[T]) Read(ctx context.Context) (T, error)

Read implements Reader.

type JSONLWriter

type JSONLWriter[T any] struct {
	// contains filtered or unexported fields
}

JSONLWriter writes each item as one JSON line. It is not safe for concurrent use.

func NewJSONLWriter

func NewJSONLWriter[T any](w io.Writer) *JSONLWriter[T]

NewJSONLWriter returns a writer to w.

func (*JSONLWriter[T]) Write

func (j *JSONLWriter[T]) Write(_ context.Context, items []T) error

Write implements Writer. Items are encoded before any byte is written, so a value that cannot be marshalled fails the chunk without a partial write.

type Job

type Job struct {
	// contains filtered or unexported fields
}

Job is an ordered list of stages; each stage is one step or several steps run in parallel. Build it with NewJob, Then and ThenParallel, then call Run. A Job is safe for concurrent Run calls on different parameters.

func NewJob

func NewJob(name string, repo Repository, opts ...JobOption) *Job

NewJob creates a job named name persisting to repo. A nil repo uses a new MemoryRepository, which makes restarts work only within the process.

func (*Job) Name

func (j *Job) Name() string

Name returns the job name.

func (*Job) Run

func (j *Job) Run(ctx context.Context, params Params) (JobRecord, error)

Run executes the job instance identified by the job name and params.

If the repository holds a completed record for the instance, Run returns it with ErrAlreadyCompleted. If it holds a failed, stopped or running record (a previous crash), Run restarts: completed steps are skipped and the failed step resumes from its last committed checkpoint. The returned record is the final persisted state, valid even when err is non-nil.

Example (Restart)

A failed job resumes from the last committed chunk when run again with the same parameters.

package main

import (
	"context"
	"errors"
	"fmt"

	"github.com/JiaBao-do/batchx"
)

func main() {
	repo := &batchx.MemoryRepository{}
	var out batchx.SliceWriter[int]
	calls := 0
	w := batchx.WriterFunc[int](func(ctx context.Context, items []int) error {
		if calls++; calls == 2 {
			return errors.New("connection lost")
		}
		return out.Write(ctx, items)
	})
	build := func() batchx.Step {
		s, _ := batchx.NewCopyStep("copy", batchx.NewSliceReader([]int{1, 2, 3, 4, 5, 6}), w, batchx.WithChunkSize(2))
		return s
	}
	params := batchx.Params{"day": "2026-09-21"}

	rec, err := batchx.NewJob("nightly", repo).Then(build()).Run(context.Background(), params)
	s, _ := rec.Step("copy")
	fmt.Println(rec.Status, "checkpoint:", s.Position, err != nil)

	rec, err = batchx.NewJob("nightly", repo).Then(build()).Run(context.Background(), params)
	fmt.Println(rec.Status, rec.Attempts, out.Items(), err)
}
Output:
failed checkpoint: 2 true
completed 2 [1 2 3 4 5 6] <nil>

func (*Job) Then

func (j *Job) Then(s Step) *Job

Then appends a sequential stage with one step and returns j.

func (*Job) ThenParallel

func (j *Job) ThenParallel(steps ...Step) *Job

ThenParallel appends a stage whose steps run concurrently. The stage ends when all steps end; if one fails the others are cancelled.

type JobOption

type JobOption func(*Job)

JobOption configures a Job.

func WithClock

func WithClock(now func() time.Time) JobOption

WithClock replaces time.Now (useful for deterministic tests).

func WithListener

func WithListener(l Listener) JobOption

WithListener adds a job-level listener; it also receives chunk, skip and retry events from every step.

type JobRecord

type JobRecord struct {
	Version   int          `json:"version"`
	Name      string       `json:"name"`
	Key       string       `json:"key"`
	Params    Params       `json:"params,omitempty"`
	Status    Status       `json:"status"`
	Attempts  int          `json:"attempts"`
	Steps     []StepRecord `json:"steps"`
	Err       string       `json:"err,omitempty"`
	StartedAt time.Time    `json:"startedAt,omitzero"`
	UpdatedAt time.Time    `json:"updatedAt,omitzero"`
	EndedAt   time.Time    `json:"endedAt,omitzero"`
}

JobRecord is the persisted state of one job instance. It is the only thing a Repository has to store, as JSON if it wants.

func (JobRecord) Clone

func (r JobRecord) Clone() JobRecord

Clone returns a deep copy.

func (JobRecord) Step

func (r JobRecord) Step(name string) (StepRecord, bool)

Step returns the record of the named step.

type LinesReader

type LinesReader struct {
	// contains filtered or unexported fields
}

LinesReader reads lines from an io.Reader, without a line length limit. Line terminators ("\n" or "\r\n") are removed. It is not safe for concurrent use.

func NewLinesReader

func NewLinesReader(r io.Reader) *LinesReader

NewLinesReader returns a reader over the lines of r.

func (*LinesReader) Read

func (l *LinesReader) Read(ctx context.Context) (string, error)

Read implements Reader. An empty final segment after the last newline is not a line.

type Listener

type Listener struct {
	// BeforeJob runs after the record is loaded, before any step.
	BeforeJob func(ctx context.Context, rec JobRecord)
	// AfterJob runs after the final record was saved.
	AfterJob func(ctx context.Context, rec JobRecord)
	// BeforeStep runs before a step executes (not for steps skipped as completed).
	BeforeStep func(ctx context.Context, job, step string)
	// AfterStep runs after a step ends, whatever the outcome.
	AfterStep func(ctx context.Context, job string, rec StepRecord)
	// AfterChunk runs after a chunk was committed.
	AfterChunk func(ctx context.Context, job string, rec StepRecord)
	// OnSkip runs when an item is skipped. item is nil for read errors.
	OnSkip func(ctx context.Context, step string, phase Phase, item any, err error)
	// OnRetry runs before each retry sleep; attempt is the failed attempt number (1-based).
	OnRetry func(ctx context.Context, step string, phase Phase, attempt int, err error)
}

Listener is a set of optional hooks; nil hooks are ignored, so the zero value is a no-op. Hooks run synchronously on the step's goroutine (parallel steps call them concurrently, so hooks must be goroutine-safe) and must not panic.

type MemoryRepository

type MemoryRepository struct {
	// contains filtered or unexported fields
}

MemoryRepository is a Repository held in memory. The zero value is ready to use and safe for concurrent use.

func (*MemoryRepository) Load

Load implements Repository.

func (*MemoryRepository) Save

Save implements Repository.

type Params

type Params map[string]string

Params identify a job instance together with the job name: running the same job with the same Params is a restart of one logical instance.

type Phase

type Phase string

Phase identifies where in a chunk an event happened.

const (
	PhaseRead    Phase = "read"
	PhaseProcess Phase = "process"
	PhaseWrite   Phase = "write"
)

Phase values.

type Processor

type Processor[I, O any] interface {
	Process(ctx context.Context, in I) (O, error)
}

Processor transforms one item. Return ErrFilter to drop the item. Implementations must be goroutine-safe only if shared between steps.

type ProcessorFunc

type ProcessorFunc[I, O any] func(ctx context.Context, in I) (O, error)

ProcessorFunc adapts a function to Processor.

func (ProcessorFunc[I, O]) Process

func (f ProcessorFunc[I, O]) Process(ctx context.Context, in I) (O, error)

Process implements Processor.

type Reader

type Reader[I any] interface {
	Read(ctx context.Context) (I, error)
}

Reader produces items. Read returns io.EOF (or an error wrapping it) when the input is exhausted. A Reader is used by one goroutine at a time.

A read error that the step treats as skippable must leave the reader positioned at the next item, so the step can continue.

type ReaderFunc

type ReaderFunc[I any] func(ctx context.Context) (I, error)

ReaderFunc adapts a function to Reader.

func (ReaderFunc[I]) Read

func (f ReaderFunc[I]) Read(ctx context.Context) (I, error)

Read implements Reader.

type Repository

type Repository interface {
	// Load returns the record for key, or ErrNotFound.
	Load(ctx context.Context, key string) (JobRecord, error)
	// Save stores rec under rec.Key, replacing any previous record.
	Save(ctx context.Context, rec JobRecord) error
}

Repository persists job records. Implementations must copy records at the boundary (callers may mutate what they pass or receive) and be safe for concurrent use. SQL or KV implementations only need these two methods.

type Retry

type Retry struct {
	// MaxAttempts is the total number of attempts including the first.
	// Values below 2 disable retrying.
	MaxAttempts int
	// Backoff returns the wait before the next attempt; attempt is the number
	// of attempts made so far (1-based). Nil means no wait.
	Backoff func(attempt int) time.Duration
	// If decides whether an error is retryable. Nil retries every error
	// except ErrFilter and context cancellation.
	If func(error) bool
}

Retry configures retrying of a failing processor call (per item) or writer call (per chunk). The zero value never retries.

type Seeker

type Seeker interface {
	Seek(ctx context.Context, n int64) error
}

Seeker is an optional Reader capability: Seek positions the reader after its first n positions (items plus skipped read errors) so a restart does not have to re-read them. Without it, batchx reads and discards n positions.

type SeqReader

type SeqReader[T any] struct {
	// contains filtered or unexported fields
}

SeqReader reads from an iter.Seq. Close stops the underlying iterator; the step closes it automatically.

func NewSeqReader

func NewSeqReader[T any](seq iter.Seq[T]) *SeqReader[T]

NewSeqReader returns a reader over seq.

func (*SeqReader[T]) Close

func (s *SeqReader[T]) Close() error

Close implements io.Closer.

func (*SeqReader[T]) Read

func (s *SeqReader[T]) Read(context.Context) (T, error)

Read implements Reader.

type SkippableError

type SkippableError struct{ Err error }

SkippableError marks an error as eligible for skipping (see WithSkipLimit).

func (*SkippableError) Error

func (e *SkippableError) Error() string

Error implements error.

func (*SkippableError) Unwrap

func (e *SkippableError) Unwrap() error

Unwrap returns the wrapped error.

type SliceReader

type SliceReader[T any] struct {
	// contains filtered or unexported fields
}

SliceReader reads from a slice. It implements Seeker (O(1) restart). It is not safe for concurrent use.

func NewSliceReader

func NewSliceReader[T any](items []T) *SliceReader[T]

NewSliceReader returns a reader over items (not copied).

func (*SliceReader[T]) Read

func (s *SliceReader[T]) Read(context.Context) (T, error)

Read implements Reader.

func (*SliceReader[T]) Seek

func (s *SliceReader[T]) Seek(_ context.Context, n int64) error

Seek implements Seeker.

type SliceWriter

type SliceWriter[T any] struct {
	// contains filtered or unexported fields
}

SliceWriter collects written items in memory. It is safe for concurrent use.

func (*SliceWriter[T]) Items

func (s *SliceWriter[T]) Items() []T

Items returns a copy of everything written so far.

func (*SliceWriter[T]) Write

func (s *SliceWriter[T]) Write(_ context.Context, items []T) error

Write implements Writer.

type Status

type Status string

Status is the lifecycle state of a job or step execution.

const (
	StatusRunning   Status = "running"
	StatusCompleted Status = "completed"
	StatusFailed    Status = "failed"
	// StatusStopped means the context was cancelled; the run is restartable.
	StatusStopped Status = "stopped"
)

Status values.

type Step

type Step interface {
	Name() string
	Execute(ctx context.Context, sc *StepContext) error
}

Step is one unit of a Job. Execute receives the step's persisted record in sc.Record (empty counters on a first run, the last committed state on a restart) and must return nil only when the step is fully done.

func NewCopyStep

func NewCopyStep[T any](name string, r Reader[T], w Writer[T], opts ...StepOption) (Step, error)

NewCopyStep is NewStep without a processor: items pass through unchanged.

func NewStep

func NewStep[I, O any](name string, r Reader[I], p Processor[I, O], w Writer[O], opts ...StepOption) (Step, error)

NewStep builds a chunk-oriented step: read up to a chunk of items, process each, write the chunk, then commit a checkpoint. A Step is single-use per run: it holds the reader and writer, so do not run one instance from two jobs at once. If r or w implement io.Closer they are closed when the step ends.

func NewTaskletStep

func NewTaskletStep(name string, fn func(context.Context) error) Step

NewTaskletStep builds a step that runs fn once. It is restarted from scratch if it did not complete.

type StepContext

type StepContext struct {
	// Job is the job name.
	Job string
	// Record is the step's state; steps update it and call Commit.
	Record StepRecord
	// contains filtered or unexported fields
}

StepContext is handed to Step.Execute.

func (*StepContext) Commit

func (sc *StepContext) Commit(ctx context.Context) error

Commit persists sc.Record. It uses a context detached from cancellation so a checkpoint for work already done is never lost to a shutdown.

type StepError

type StepError struct {
	Job  string
	Step string
	Err  error
}

StepError reports which step of which job failed.

func (*StepError) Error

func (e *StepError) Error() string

Error implements error.

func (*StepError) Unwrap

func (e *StepError) Unwrap() error

Unwrap returns the underlying error.

type StepOption

type StepOption func(*stepConfig)

StepOption configures a chunk step.

func WithChunkSize

func WithChunkSize(n int) StepOption

WithChunkSize sets the commit interval: items per chunk (default 10).

func WithRetry

func WithRetry(r Retry) StepOption

WithRetry sets the retry policy for processor and writer calls.

Example

Retry with exponential backoff for flaky writers.

package main

import (
	"context"
	"errors"
	"fmt"

	"github.com/JiaBao-do/batchx"
)

func main() {
	attempts := 0
	w := batchx.WriterFunc[int](func(context.Context, []int) error {
		if attempts++; attempts < 3 {
			return errors.New("busy")
		}
		return nil
	})
	step, _ := batchx.NewCopyStep("send", batchx.NewSliceReader([]int{1}), w,
		batchx.WithRetry(batchx.Retry{MaxAttempts: 5, Backoff: batchx.ConstantBackoff(0)}))
	_, err := batchx.NewJob("send", nil).Then(step).Run(context.Background(), nil)
	fmt.Println(attempts, err)
}
Output:
3 <nil>

func WithSkipIf

func WithSkipIf(f func(error) bool) StepOption

WithSkipIf sets which errors are skippable (default IsSkippable).

func WithSkipLimit

func WithSkipLimit(n int) StepOption

WithSkipLimit allows up to n skipped items per step across all runs of the instance (default 0: any error fails the step).

Example

Skippable errors are skipped up to the skip limit and reported to listeners.

package main

import (
	"context"
	"fmt"

	"github.com/JiaBao-do/batchx"
)

func main() {
	parse := batchx.ProcessorFunc[string, int](func(_ context.Context, s string) (int, error) {
		var n int
		if _, err := fmt.Sscanf(s, "%d", &n); err != nil {
			return 0, batchx.Skippable(err)
		}
		return n, nil
	})
	var out batchx.SliceWriter[int]
	step, _ := batchx.NewStep("parse", batchx.NewSliceReader([]string{"1", "x", "3"}), parse, &out,
		batchx.WithSkipLimit(1),
		batchx.WithStepListener(batchx.Listener{OnSkip: func(_ context.Context, _ string, p batchx.Phase, item any, _ error) {
			fmt.Printf("skipped %v in %s\n", item, p)
		}}))
	_, err := batchx.NewJob("parse", nil).Then(step).Run(context.Background(), nil)
	fmt.Println(out.Items(), err)
}
Output:
skipped x in process
[1 3] <nil>

func WithSleep

func WithSleep(f func(context.Context, time.Duration) error) StepOption

WithSleep replaces the backoff sleeper (useful for deterministic tests).

func WithStepListener

func WithStepListener(l Listener) StepOption

WithStepListener adds a listener for this step only.

type StepRecord

type StepRecord struct {
	Name     string `json:"name"`
	Status   Status `json:"status"`
	Attempts int    `json:"attempts"`
	// Position is the restart checkpoint: the number of reader positions
	// (items plus skipped read errors) consumed by committed chunks.
	Position  int64     `json:"position"`
	Read      int64     `json:"read"`
	Written   int64     `json:"written"`
	Filtered  int64     `json:"filtered"`
	Skipped   int64     `json:"skipped"`
	Retries   int64     `json:"retries"`
	Commits   int64     `json:"commits"`
	Err       string    `json:"err,omitempty"`
	StartedAt time.Time `json:"startedAt,omitzero"`
	EndedAt   time.Time `json:"endedAt,omitzero"`
}

StepRecord is the persisted state of one step. Counters reflect committed chunks only.

type Writer

type Writer[O any] interface {
	Write(ctx context.Context, items []O) error
}

Writer writes one chunk. A chunk write should be all-or-nothing: if Write returns an error, none of the chunk may be considered written, because batchx may retry it, or (for skippable errors) re-write its items one at a time to isolate the bad one. Because a crash can happen after Write returns but before the checkpoint is saved, writers should be idempotent to get effectively-once results; batchx itself gives at-least-once per chunk.

type WriterFunc

type WriterFunc[O any] func(ctx context.Context, items []O) error

WriterFunc adapts a function to Writer.

func (WriterFunc[O]) Write

func (f WriterFunc[O]) Write(ctx context.Context, items []O) error

Write implements Writer.

Directories

Path Synopsis
examples
demo command
Realistic multi-step job: prepare, load two reference datasets in parallel, then turn orders into invoices (joining the reference data, retrying a flaky tax lookup, skipping bad rows), and finally report.
Realistic multi-step job: prepare, load two reference datasets in parallel, then turn orders into invoices (joining the reference data, retrying a flaky tax lookup, skipping bad rows), and finally report.
filerepo command
File repository demo: job state lives in one plain JSON file per job instance, so it survives process restarts and can be inspected, backed up or deleted with ordinary tools.
File repository demo: job state lives in one plain JSON file per job instance, so it survives process restarts and can be inspected, backed up or deleted with ordinary tools.
quickstart command
Quickstart: read CSV, process each row, write JSON lines, in chunks of 2.
Quickstart: read CSV, process each row, write JSON lines, in chunks of 2.
restart command
Restart demo: a job crashes in the middle of the input, then the same job is run again with the same parameters and resumes from the last committed chunk.
Restart demo: a job crashes in the middle of the input, then the same job is run again with the same parameters and resumes from the last committed chunk.
skipretry command
Skip and retry demo: bad CSV rows are skipped (up to a limit), a flaky processor call is retried with backoff, and listeners report both.
Skip and retry demo: bad CSV rows are skipped (up to a limit), a flaky processor call is retried with backoff, and listeners report both.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL