change cancel place (#82)
This commit is contained in:
parent
2037463c62
commit
245d1bd6b5
1 changed files with 2 additions and 2 deletions
|
|
@ -135,6 +135,8 @@ func (c *Consumer) Scan(ctx context.Context, fn func(*Record) ScanStatus) error
|
||||||
defer wg.Done()
|
defer wg.Done()
|
||||||
|
|
||||||
if err := c.ScanShard(ctx, shardID, fn); err != nil {
|
if err := c.ScanShard(ctx, shardID, fn); err != nil {
|
||||||
|
cancel()
|
||||||
|
|
||||||
select {
|
select {
|
||||||
case errc <- fmt.Errorf("shard %s error: %v", shardID, err):
|
case errc <- fmt.Errorf("shard %s error: %v", shardID, err):
|
||||||
// first error to occur
|
// first error to occur
|
||||||
|
|
@ -142,8 +144,6 @@ func (c *Consumer) Scan(ctx context.Context, fn func(*Record) ScanStatus) error
|
||||||
// error has already occured
|
// error has already occured
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
cancel()
|
|
||||||
}(shardID)
|
}(shardID)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue