reset for a new epoch of samples
(self)
| 247 | raise StopIteration() |
| 248 | |
| 249 | def reset(self): |
| 250 | """ reset for a new epoch of samples |
| 251 | """ |
| 252 | assert not self._exit, "cannot reset for already stopped dataset" |
| 253 | |
| 254 | if self._epoch < 0: |
| 255 | self._epoch = 0 |
| 256 | for w in self._consumers: |
| 257 | w.start() |
| 258 | self._producer.start() |
| 259 | else: |
| 260 | assert self._consumer_healthy(), "cannot start another pass of data" \ |
| 261 | " for some consumers exited abnormally before!!!" |
| 262 | |
| 263 | if not self.drained(): |
| 264 | logger.warn("reset before epoch[{}] finishes".format( |
| 265 | self._epoch)) |
| 266 | self._produced = self._produced - self._consumed |
| 267 | else: |
| 268 | self._produced = 0 |
| 269 | |
| 270 | self._epoch += 1 |
| 271 | |
| 272 | assert len(self._consumer_endsig.keys()) == 0, "some consumers already exited," \ |
| 273 | + " cannot start another epoch" |
| 274 | |
| 275 | self._source.reset() |
| 276 | self._souce_drained = False |
| 277 | self._consumed = 0 |
| 278 | self._feeding_ev.set() |
| 279 | |
| 280 | |
| 281 | # FIXME(dengkaipeng): fix me if you have better impliment |
no test coverage detected