(input <-chan baker.OutputRecord, _ chan<- string)
| 33 | } |
| 34 | |
| 35 | func (s *Shardable) Run(input <-chan baker.OutputRecord, _ chan<- string) error { |
| 36 | // Do something with `input` record. |
| 37 | // s.idx identifies the output process index and should |
| 38 | // be used to manage the sharding |
| 39 | for data := range input { |
| 40 | log.Printf(`Shard #%d: Getting "%s"`, s.idx, data.Record) |
| 41 | } |
| 42 | |
| 43 | return nil |
| 44 | } |
| 45 | |
| 46 | func (s *Shardable) Stats() baker.OutputStats { return baker.OutputStats{} } |
nothing calls this directly
no outgoing calls
no test coverage detected