| 16 | } |
| 17 | |
| 18 | bool Skip::getNextTuplesInternal(ExecutionContext* context) { |
| 19 | auto numTupleSkippedBefore = 0u; |
| 20 | auto numTuplesAvailable = 1u; |
| 21 | do { |
| 22 | restoreSelVector(*dataChunkToSelect->state); |
| 23 | // end of execution due to no more input |
| 24 | if (!children[0]->getNextTuple(context)) { |
| 25 | return false; |
| 26 | } |
| 27 | saveSelVector(*dataChunkToSelect->state); |
| 28 | numTuplesAvailable = resultSet->getNumTuples(dataChunksPosInScope); |
| 29 | numTupleSkippedBefore = counter->fetch_add(numTuplesAvailable); |
| 30 | } while (numTupleSkippedBefore + numTuplesAvailable <= skipNumber); |
| 31 | auto numTupleToSkipInCurrentResultSet = (int64_t)(skipNumber - numTupleSkippedBefore); |
| 32 | if (numTupleToSkipInCurrentResultSet <= 0) { |
| 33 | // Other thread has finished skipping. Process everything in current result set. |
| 34 | metrics->numOutputTuple.increase(numTuplesAvailable); |
| 35 | } else { |
| 36 | // If all dataChunks are flat, numTupleAvailable = 1 which means numTupleSkippedBefore = |
| 37 | // skipNumber. So execution is handled in above if statement. |
| 38 | DASSERT(!dataChunkToSelect->state->isFlat()); |
| 39 | auto buffer = dataChunkToSelect->state->getSelVectorUnsafe().getMutableBuffer(); |
| 40 | if (dataChunkToSelect->state->getSelVector().isUnfiltered()) { |
| 41 | for (uint64_t i = numTupleToSkipInCurrentResultSet; |
| 42 | i < dataChunkToSelect->state->getSelVector().getSelSize(); ++i) { |
| 43 | buffer[i - numTupleToSkipInCurrentResultSet] = i; |
| 44 | } |
| 45 | dataChunkToSelect->state->getSelVectorUnsafe().setToFiltered(); |
| 46 | } else { |
| 47 | for (uint64_t i = numTupleToSkipInCurrentResultSet; |
| 48 | i < dataChunkToSelect->state->getSelVector().getSelSize(); ++i) { |
| 49 | buffer[i - numTupleToSkipInCurrentResultSet] = buffer[i]; |
| 50 | } |
| 51 | } |
| 52 | dataChunkToSelect->state->getSelVectorUnsafe().setSelSize( |
| 53 | dataChunkToSelect->state->getSelVector().getSelSize() - |
| 54 | numTupleToSkipInCurrentResultSet); |
| 55 | metrics->numOutputTuple.increase(dataChunkToSelect->state->getSelVector().getSelSize()); |
| 56 | } |
| 57 | return true; |
| 58 | } |
| 59 | |
| 60 | } // namespace processor |
| 61 | } // namespace lbug |
nothing calls this directly
no test coverage detected