How to groupBy and then debounce each group correctly in RxGo?

Viewed 81

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!

Live demo

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)
    }
}
0 Answers
Related