(ctx context.Context, sctx *super.Context, r io.Reader, p sbuf.Pushdown, concurrentReaders int)
| 30 | } |
| 31 | |
| 32 | func NewVectorReader(ctx context.Context, sctx *super.Context, r io.Reader, p sbuf.Pushdown, concurrentReaders int) (*VectorReader, error) { |
| 33 | if concurrentReaders < 1 { |
| 34 | panic(concurrentReaders) |
| 35 | } |
| 36 | ra, ok := r.(io.ReaderAt) |
| 37 | if !ok { |
| 38 | return nil, errors.New("Super Columnar requires a seekable input") |
| 39 | } |
| 40 | var metaFilters []*metafilter |
| 41 | if p != nil { |
| 42 | filter, _, err := p.MetaFilter() |
| 43 | if err != nil { |
| 44 | return nil, err |
| 45 | } |
| 46 | if filter != nil { |
| 47 | for range concurrentReaders { |
| 48 | filter, projection, err := p.MetaFilter() |
| 49 | if err != nil { |
| 50 | return nil, err |
| 51 | } |
| 52 | metaFilters = append(metaFilters, &metafilter{filter, projection}) |
| 53 | } |
| 54 | } |
| 55 | } |
| 56 | |
| 57 | return &VectorReader{ |
| 58 | ctx: ctx, |
| 59 | sctx: sctx, |
| 60 | activeReaders: &atomic.Int64{}, |
| 61 | stream: &stream{ctx: ctx, r: ra}, |
| 62 | pushdown: p, |
| 63 | metaFilters: metaFilters, |
| 64 | readerAt: ra, |
| 65 | }, nil |
| 66 | } |
| 67 | |
| 68 | type metafilter struct { |
| 69 | filter expr.Evaluator |
no test coverage detected