(
id: number,
inputA: DifferenceStreamReader<[K, V]>,
output: DifferenceStreamWriter<[K, number]>,
)
| 9 | */ |
| 10 | export class CountOperator<K, V> extends ReduceOperator<K, V, number> { |
| 11 | constructor( |
| 12 | id: number, |
| 13 | inputA: DifferenceStreamReader<[K, V]>, |
| 14 | output: DifferenceStreamWriter<[K, number]>, |
| 15 | ) { |
| 16 | const countInner = (vals: Array<[V, number]>): Array<[number, number]> => { |
| 17 | let totalCount = 0 |
| 18 | for (const [_, diff] of vals) { |
| 19 | totalCount += diff |
| 20 | } |
| 21 | return [[totalCount, 1]] |
| 22 | } |
| 23 | |
| 24 | super(id, inputA, output, countInner) |
| 25 | } |
| 26 | } |
| 27 | |
| 28 | /** |
nothing calls this directly
no outgoing calls
no test coverage detected