I have a stream "1", "1", "3", "1", "3" which has 1s delay between each two. I hope to groupBy the value first. Then for each group debounce by 2.4s. At the end, I hope to get the result "1", "3", "1", "3".
Here is a RxJS demo I found doing similar thing.
And below is my code using RxGo:
However, the current code still output "1", "1", "3", "1", "3".
Any idea? Thanks!
package main
import (
"fmt"
"github.com/reactivex/rxgo/v2"
"time"
)
func main() {
// Simulate a stream
ch := make(chan rxgo.Item)
go func() {
items := []string{"1", "1", "3", "1", "3"}
for i := 0; i < len(items); i++ {
ch <- rxgo.Item{V: items[i]}
time.Sleep(time.Second * 1) // 1s delay between each two
}
close(ch)
}()
// Group by the value
observable := rxgo.
FromChannel(ch).
GroupByDynamic(func(i rxgo.Item) string {
return i.V.(string)
}, rxgo.WithBufferedChannel(10))
// Debounce for each group
list := []rxgo.Observable{}
for i := range observable.Observe() {
obs := i.V.(rxgo.Observable)
obs.Debounce(rxgo.WithDuration(2400 * time.Millisecond)) // debounce every 2.4s
list = append(list, obs)
}
// Merge the group of observables to one observable
observable = rxgo.Merge(list)
// Print the result
for i := range observable.Observe() {
// Currently it prints "1", "1", "3", "1", "3"
// I hope it prints "1", "3", "1", "3"
fmt.Printf("item: %v\n", i.V)
}
}