Documentation
¶
Overview ¶
Package mr provides a generic MapReduce with a bounded worker pool, guarded channel writes and first-error cancellation. It is modeled on go-zero's core/mr so that fan-out/fan-in data processing can be expressed without manual goroutine and channel plumbing.
Index ¶
- Variables
- func Finish(fns ...func() error) error
- func FinishVoid(fns ...func())
- func ForEach[T any](generate GenerateFunc[T], mapper ForEachFunc[T], opts ...Option)
- func MapReduce[T, U, V any](generate GenerateFunc[T], mapper MapperFunc[T, U], reducer ReducerFunc[U, V], ...) (V, error)
- func MapReduceChan[T, U, V any](source <-chan T, mapper MapperFunc[T, U], reducer ReducerFunc[U, V], ...) (V, error)
- type ForEachFunc
- type GenerateFunc
- type MapFunc
- type MapperFunc
- type Option
- type ReducerFunc
- type VoidReducerFunc
- type Writer
Constants ¶
This section is empty.
Variables ¶
var ( // ErrCancelWithNil is returned when a mapper/reducer cancels with a nil error. ErrCancelWithNil = errors.New("mr: cancelled with nil") // ErrReduceNoOutput is returned when the reducer writes no value. ErrReduceNoOutput = errors.New("mr: reducer wrote no value") )
Functions ¶
func FinishVoid ¶
func FinishVoid(fns ...func())
FinishVoid runs fns in parallel and ignores errors.
func ForEach ¶
func ForEach[T any](generate GenerateFunc[T], mapper ForEachFunc[T], opts ...Option)
ForEach maps all elements generated by generate with no output.
func MapReduce ¶
func MapReduce[T, U, V any](generate GenerateFunc[T], mapper MapperFunc[T, U], reducer ReducerFunc[U, V], opts ...Option) (V, error)
MapReduce maps all elements generated by generate, then reduces the outputs with reducer.
func MapReduceChan ¶
func MapReduceChan[T, U, V any](source <-chan T, mapper MapperFunc[T, U], reducer ReducerFunc[U, V], opts ...Option) (V, error)
MapReduceChan is MapReduce over an existing source channel.
Types ¶
type ForEachFunc ¶
type ForEachFunc[T any] func(item T)
ForEachFunc processes an element without producing output.
type GenerateFunc ¶
type GenerateFunc[T any] func(source chan<- T)
GenerateFunc sends input elements into source.
type MapFunc ¶
type MapFunc[T, U any] func(item T, writer Writer[U])
MapFunc processes an element and writes output to the writer.
type MapperFunc ¶
MapperFunc processes an element and writes output, or cancels processing.
type Option ¶
type Option func(opts *mapReduceOptions)
Option customizes a MapReduce.
func WithContext ¶
WithContext makes the processing honor the given context.
func WithWorkers ¶
func WithWorkers(workers int) Option
WithWorkers sets the number of worker goroutines (minimum 1).
type ReducerFunc ¶
ReducerFunc reduces the mapper outputs into a single value, or cancels.
type VoidReducerFunc ¶
VoidReducerFunc reduces the mapper outputs without producing a value.