mr

package
v0.0.5 Latest Latest
Warning

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

Go to latest
Published: Oct 2, 2026 License: MIT Imports: 7 Imported by: 0

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

Constants

This section is empty.

Variables

View Source
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 Finish

func Finish(fns ...func() error) error

Finish runs fns in parallel, cancelling on the first error.

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

type MapperFunc[T, U any] func(item T, writer Writer[U], cancel func(error))

MapperFunc processes an element and writes output, or cancels processing.

type Option

type Option func(opts *mapReduceOptions)

Option customizes a MapReduce.

func WithContext

func WithContext(ctx context.Context) Option

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

type ReducerFunc[U, V any] func(pipe <-chan U, writer Writer[V], cancel func(error))

ReducerFunc reduces the mapper outputs into a single value, or cancels.

type VoidReducerFunc

type VoidReducerFunc[U any] func(pipe <-chan U, cancel func(error))

VoidReducerFunc reduces the mapper outputs without producing a value.

type Writer

type Writer[T any] interface {
	Write(v T)
}

Writer wraps Write, the sink used by mappers and reducers.

Jump to

Keyboard shortcuts

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