(seq dag.Seq, replicas int)
| 73 | } |
| 74 | |
| 75 | func (o *Optimizer) parallelizeFileScan(seq dag.Seq, replicas int) (dag.Seq, error) { |
| 76 | // Prepend a pass so we can parallelize seq[0]. |
| 77 | seq = append(dag.Seq{dag.Pass}, seq...) |
| 78 | n, sortExprs, _, err := o.concurrentPath(seq, nil) |
| 79 | if err != nil { |
| 80 | return nil, err |
| 81 | } |
| 82 | if n < len(seq) { |
| 83 | switch seq[n].(type) { |
| 84 | case *dag.AggregateOp, *dag.SortOp, *dag.TopOp: |
| 85 | return parallelizeHead(seq, n, sortExprs, replicas), nil |
| 86 | } |
| 87 | } |
| 88 | return nil, nil |
| 89 | } |
| 90 | |
| 91 | func (o *Optimizer) parallelizeSeqScan(seq dag.Seq, replicas int) (dag.Seq, error) { |
| 92 | scan := seq[0].(*dag.SeqScan) |
no test coverage detected