MergeSlicesParallel merge sorted slices in parallel using the MergeSortedSlices function
(ctx context.Context, parallelism int, a ...[]string)
| 37 | // MergeSlicesParallel merge sorted slices in parallel |
| 38 | // using the MergeSortedSlices function |
| 39 | func MergeSlicesParallel(ctx context.Context, parallelism int, a ...[]string) ([]string, error) { |
| 40 | if parallelism <= 1 { |
| 41 | return MergeSortedSlices(ctx, a...) |
| 42 | } |
| 43 | if len(a) == 0 { |
| 44 | return nil, nil |
| 45 | } |
| 46 | if len(a) == 1 { |
| 47 | return a[0], nil |
| 48 | } |
| 49 | c := make(chan []string, len(a)) |
| 50 | errCh := make(chan error, 1) |
| 51 | wg := sync.WaitGroup{} |
| 52 | var r [][]string |
| 53 | p := min(parallelism, len(a)/2) |
| 54 | batchSize := len(a) / p |
| 55 | |
| 56 | for i := 0; i < len(a); i += batchSize { |
| 57 | wg.Add(1) |
| 58 | go func(i int) { |
| 59 | m := min(len(a), i+batchSize) |
| 60 | r, e := MergeSortedSlices(ctx, a[i:m]...) |
| 61 | if e != nil { |
| 62 | errCh <- e |
| 63 | wg.Done() |
| 64 | return |
| 65 | } |
| 66 | c <- r |
| 67 | wg.Done() |
| 68 | }(i) |
| 69 | } |
| 70 | |
| 71 | go func() { |
| 72 | wg.Wait() |
| 73 | close(c) |
| 74 | close(errCh) |
| 75 | }() |
| 76 | |
| 77 | if err := <-errCh; err != nil { |
| 78 | return nil, err |
| 79 | } |
| 80 | for s := range c { |
| 81 | r = append(r, s) |
| 82 | } |
| 83 | |
| 84 | return MergeSortedSlices(ctx, r...) |
| 85 | } |
| 86 | |
| 87 | func NewStringListIter(s []string) *StringListIter { |
| 88 | return &StringListIter{l: s} |