This is package which make use of the Generics feature introduced in Go 1.18 for providing a functional oriented rich Map/Reduce API in Go. As this is still an experimental package, some of its apis may change slightly in future releases.
This api leverages Go Generics for providing a set of higher-order Map/Reduce functions. These functions, when chained together, allow for functional programming techniques that ultimately reduce code duplication and make it easier to transform and iterate over collections of elements. The current Generics' implementation in Go 1.18 doesn't allow methods themselves to have additional type parameters. This limitation forces mapping functions, whose input and return type differs, to be defined as top-level, limiting their chaining capacity and compromising readability. Despite this limitation, some of the most common functional programming techniques are still possible to be implemented in Go, as samples bellow illustrate.
All available Api operations are enumerated here.
Standard go get
$ go get
This api is a wrapper of its inputs, either a slice or a set of elements of a generic type. This wrapper is internally backed from a slice - either passed explicitly by a constructor or created on the fly if built by a set of elements - providing a set of operations applied over the contents of its backing slice. All operations have been implemented aiming minimal memory allocation, hence never allocating intermediate slices.
stringStrm := strm.Of("Hey!", "Hello!", "Hi!")
intStrm := strm.Of(1, 2, 3, 4)
sliceStrm := strm.Of([]int{1}, []int{1, 2}, []int{1, 2, 3})
backingSlice := []int {1, 2, 3}
intStrm := strm.From(backingSlice)
The strm will create its own backing slice
which will be a copy of the initSlice
The initSlice
state remain unmodified, independent of the operation applied to the strm.
initSlice := []int {1, 2, 3}
intStrm := strm.CopyFrom(initSlice)
slice := strm.Of(1, 2, 3, 4).ToSlice()
isEven := func(n int) bool { return n%2 == 0 }
// Filtering even elements
// evenSlice -> [2 4]
evenSlice := strm.Of(1, 2, 3, 4, 5).
// filter chaining
// evenSlice -> [4]
evenSlice := strm.Of(1, 2, 3, 4, 5).
Filter(func(n int) bool { return n > 2 }).
// iterating over all elements
// prints -> n: 2; n: 4;
strm.Of(1, 2, 3, 4, 5).
ForEach(func(n int) { fmt.Printf("n: %v;\t", n) })
// printing each element and filtering
// prints -> n: 1
// n: 2
// even: 2
// slice -> [2]
slice := strm.Of(1, 2).
OnEach(func(n int) { fmt.Printf("n: %v\n", n) }).
OnEach(func(n int) { fmt.Printf("even: %v\n", n) }).
// applies a transformation (n->n) to each element
// sqrSlice -> [1 4 9 16 25 36]
sqrSlice := strm.Of(1, 2, 3, 4, 5, 6).
ApplyOnEach(func(n int) int { return n * n }).
type Person struct { name string; age int }
people := []Person{{"Peter", 30}, {"John", 18}, {"Sarah", 16}, {"Kate", 16}}
// maps a go-strm of (Person) to a slice of (string)
// names -> [Peter John Sarah Kate]
names := Map(From(people), func(p Person) string { return }).
slices := [][]int{{1}, {1, 2}, {1, 2, 3}}
// sums -> [1 3 6]
sums := Map(
func(it []int) int { return Reduce(From(it), func(a int, b int) int { return a + b }) },
// sums -> [1 3 6]
sums := Map(
func(it []int) int { return Sum(From(it)) },
// mins -> [1 1 1]
mins := Map(
func(it []int) int { return Min(From(it)) },
// maxs -> [1 2 3]
maxs := Map(
func(it []int) int { return Max(From(it)) },
// flatSlice -> [1 1 2 1 2 3]
flatSlice := FlatMap(
func(it []int) *Stream[int] { return From(it) },
A PMap
function is available for applying the given mapping function over all stream elements in parallel, leveraging
goroutines. The PMap
usage is similar to Map
. By default, PMap
will launch a new goroutine per each element
present in the given Stream. If the batching
flag is provided, the parallel work is batched by number of
available logical CPUs.
people := []Person{{"Peter", 30}, {"John", 18}, {"Sarah", 16}, {"Kate", 16}}
// maps a go-strm of (Person) to a slice of (string) in parallel
// names -> [Peter John Sarah Kate]
names := strm.
PMap(strm.From(people), func(p Person) string { return }).
// maps a go-strm of (Person) to a slice of (string) in parallel without batching
// names -> [Peter John Sarah Kate]
names := strm.
PMap(strm.From(people), func(p Person) string { return }, true).
people := strm.Of(Person{"Tim", 30}, Person{"Bil", 40}, Person{"John", 30}, Person{"Tim", 35})
// byAge -> map[30:[{Tim 30} {John 30}] 35:[{Tim 35}] 40:[{Bil 40}]]
byAge := strm.GroupBy(people, func(it Person) int { return it.age })
// byName -> map[Bil:[{Bil 40}] John:[{John 30}] Tim:[{Tim 30} {Tim 35}]]
byName := strm.GroupBy(people, func(it Person) string { return })
de-dupes strms of both Comparable
and Non-Comparable
// deduped -> [2 3 4 5 6]
deduped := strm.Of(2, 2, 3, 4, 4, 5, 6, 6).
// dedupedStruct -> [{Peter 18} {Bruce 48}]
dedupedStruct := strm.Of(Person{"Peter", 18}, Person{"Peter", 18}, Person{"Bruce", 48}).
// dedupedSlices -> [[1 2] [3 4]]
dedupedSlices := strm.Of([]int{1, 2}, []int{1, 2}, []int{3, 4}).
// reversed -> [6 5 4 3 2 1]
reversed := strm.Of(1, 2, 3, 4, 5, 6).
// all -> false
all := strm.Of(1, 2, 3, 4, 5).
All(func(n int) bool { return n < 5 })
// none -> false
none := strm.Of(1, 2, 3, 4, 5).
None(func(n int) bool { return n < 5 })
// any -> true
any := strm.Of("Hey!", "Hello!", "Hi!").
Any(func(n string) bool { return n == "Hi!" })
// count -> 3
count := strm.Of(2, 2, 3, 4, 4, 5, 6, 6).
// Sum only even numbers
// sumBy -> 6
sumBy := strm.Of(1, 2, 3, 4, 5).
SumBy(func(n int) int { if n%2 == 0 { return n } else { return 0 } })
// Count only even numbers
// countBy -> 2
countBy := strm.Of(1, 2, 3, 4, 5).
// contains -> true
contains := strm.Of(Person{"Peter", 18}, Person{"John", 30}, Person{"Bruce", 48}).
Contains(Person{"Bruce", 48})
// contains -> true
contains := strm.Of([]int{1, 2}, []int{3, 4}).
Contains([]int{1, 2})
// names -> Peter,John,Sarah,Kate
people := From([]Person{{"Peter", 30}, {"John", 18}, {"Sarah", 16}, {"Kate", 16}})
names := Map(people, func(p Person) string { return }).
// converting to batches of 2 elements each
// batches -> [[1 2] [4 6] [14 1] [2]]
batches := strm.Of(1, 2, 4, 6, 14, 1, 2).
// converting to windows of 5 elements with a step of 3, without partial windows at the end
// windows -> [[1 2 3 4 5] [4 5 6 7 8] [7 8 9 10 11] [10 11 12 13 14]]
windows := strm.Of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15).
Windowed(5, 3)
// converting to windows of 5 elements with a step of 3, preserving all partial windows at the end
// partWindows -> [[1 2 3 4 5] [4 5 6 7 8] [7 8 9 10 11] [10 11 12 13 14] [13 14 15]]
partWindows := strm.Of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15).
Windowed(5, 3, true)
// first -> {Peter 18}
first := strm.Of(Person{"Peter", 18}, Person{"John", 30}, Person{"Bruce", 18}).
// firstBy -> {John 30}
firstBy := strm.Of(Person{"Peter", 18}, Person{"John", 30}, Person{"Bruce", 18}).
FirstBy(func(p Person) bool { return p.age > 18 })
// last -> {"Bruce", 18}
last := strm.Of(Person{"Peter", 18}, Person{"John", 30}, Person{"Bruce", 18}).
// take -> [0 1]
take := strm.Of(0, 1, 2, 3).
// drop -> [1 2 3]
drop := strm.Of(0, 1, 2, 3).
An int range, represented by the IntStream
type, is also available for convenience of use Sum
, Min
, Max
, Avg
and Sorted
This type encloses a *Stream[int]
exclusively, leveraging all methods already available in the Stream type.
// sum -> 15
sum := strm.Range(1,5).Sum()
// min -> 1
min := strm.RangeOf(1, 2, 3, 4, 5).Min()
// max -> 5
max := strm.RangeFrom([]int{1, 2, 3, 4, 5}).Max()
// avg -> 2
avg := strm.Range(1,4).Avg()
// sorted -> [1 2 3 4 5]
sorted := strm.RangeOf(5, 4, 2, 1, 3).
// mappedSlice -> [2 3 4 5 6]
mappedSlice := Map(
RangeOf(5, 4, 2, 1, 3).Sorted().ToStrm(),
func(it int) int { return it+1 },
// batches -> [[2 3] [4 5] [6]]
batches := strm.RangeOf(5, 4, 2, 1, 3, 6).
Filter(func(it int) bool { return it > 1 }).
Performance-wise, single mapping and filtering ops perform very well. Chained operations like applying several mappings and filters over the strm, can be slower than just performing native for loops. The following benchmarks can be found at strm_test.go.
goos: darwin
goarch: amd64
pkg: strm/strm
cpu: Intel(R) Core(TM) i9-8950HK CPU @ 2.90GHz
name old time/op new time/op delta
Filter-12 114ns ± 2% 98ns ± 5% -13.81% (p=0.008 n=5+5)
Distinct-12 526ns ± 3% 461ns ± 5% -12.48% (p=0.008 n=5+5)
Map-12 1.36µs ± 3% 1.55µs ± 3% +13.35% (p=0.008 n=5+5)
ChainedFilters-12 117ns ± 4% 138ns ± 5% +17.97% (p=0.008 n=5+5)
MapFilter-12 1.29µs ± 5% 1.65µs ± 3% +27.77% (p=0.008 n=5+5)
// Constructors
func Of[T any](elems ...T) *Stream[T]
func From[T any](backingSlice []T) *Stream[T]
func CopyFrom[T any](slice []T) *Stream[T]
// Top-Level functions
func Map[IN any, OUT any](s *Stream[IN], f func(IN) OUT) *Stream[OUT]
func PMap[IN any, OUT any](s *Stream[IN], f func(IN) OUT) *Stream[OUT]
func FlatMap[IN any, OUT any](s *Stream[IN], f func(v IN) *Stream[OUT]) *Stream[OUT]
func Reduce[IN any, OUT any](s *Stream[IN], f reducer[OUT, IN], start ...OUT) OUT
func GroupBy[K comparable, V any](s *Stream[V], keySelector func(V) K) map[K][]V
func Max[O Ordered](s *Stream[O]) O
func Min[O Ordered](s *Stream[O]) O
func Sum[O Ordered](s *Stream[O]) O
func Merge[T any](streams ...*Stream[T]) *Stream[T]
// go-strm operations
func Filter(predicate func(T) bool) *Stream[T]
func ApplyOnEach(action func(T) T) *Stream[T]
func OnEach(f func(T)) *Stream[T]
func Plus(other *Stream[T]) *Stream[T]
func Append(elems []T) *Stream[T]
func Take(n int) *Stream[T]
func Drop(n int) *Stream[T]
func Reversed() *Stream[T]
func Distinct() *Stream[T]
// Terminal go-strm operations
func ToSlice() []T
func ForEach(action func(T))
func Any(predicate func(T) bool) bool
func All(predicate func(T) bool) bool
func None(predicate func(T) bool) bool
func Count() int
func CountBy(predicate func(T) bool) int
func SumBy(selector func(T) int) int
func FirstBy(predicate func(T) bool) T
func First() T
func Last() T
func Contains(element T) bool
func JoinToString(delimiter string) string
func Chunked(batchSize int) [][]T
func Windowed(size int, step int, partialWindows ...bool) [][]T
// Int Ranges operations
func Range(from int, to int) *IntStream
func RangeOf(elems *IntStream
func RangeFrom(backingIntSlice []int) *IntStream
func RangeCopyFrom(backingIntSlice []int) *IntStream
func Sorted() *IntStream
func Sum() int
func Min() int
func Max() int
func Avg() int
func ToStrm() *Stream[int]