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 ¶
- Variables
- func ConstantBackoff(d time.Duration) func(int) time.Duration
- func ExponentialBackoff(base, max time.Duration) func(int) time.Duration
- func IsSkippable(err error) bool
- func Skippable(err error) error
- type CSVReader
- type CSVWriter
- type ChanReader
- type FileRepository
- type JSONLReader
- type JSONLWriter
- type Job
- type JobOption
- type JobRecord
- type LinesReader
- type Listener
- type MemoryRepository
- type Params
- type Phase
- type Processor
- type ProcessorFunc
- type Reader
- type ReaderFunc
- type Repository
- type Retry
- type Seeker
- type SeqReader
- type SkippableError
- type SliceReader
- type SliceWriter
- type Status
- type Step
- type StepContext
- type StepError
- type StepOption
- type StepRecord
- type Writer
- type WriterFunc
Examples ¶
Constants ¶
This section is empty.
Variables ¶
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 ¶
ConstantBackoff waits d between attempts.
func ExponentialBackoff ¶
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 ¶
IsSkippable is the default skip predicate: it reports whether err wraps a *SkippableError.
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 ¶
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]
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 ¶
NewCSVWriter returns a writer to w. configure may adjust the *csv.Writer (Comma, UseCRLF) used to encode each chunk.
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.
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.
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.
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.
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) Run ¶
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) ThenParallel ¶
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 WithListener ¶
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.
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.
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.
type Params ¶
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 Processor ¶
Processor transforms one item. Return ErrFilter to drop the item. Implementations must be goroutine-safe only if shared between steps.
type ProcessorFunc ¶
ProcessorFunc adapts a function to Processor.
type Reader ¶
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 ¶
ReaderFunc adapts a function to 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 ¶
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 ¶
NewSeqReader returns a reader over seq.
type SkippableError ¶
type SkippableError struct{ Err error }
SkippableError marks an error as eligible for skipping (see WithSkipLimit).
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).
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.
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 ¶
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.
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.
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 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 ¶
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.
Source Files
¶
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. |