| 19 | } |
| 20 | |
| 21 | Block DistinctSortedBlockInputStream::readImpl() |
| 22 | { |
| 23 | /// Execute until end of stream or until |
| 24 | /// a block with some new records will be gotten. |
| 25 | for (;;) |
| 26 | { |
| 27 | /// Stop reading if we already reached the limit. |
| 28 | if (limit_hint && data.getTotalRowCount() >= limit_hint) |
| 29 | return Block(); |
| 30 | |
| 31 | Block block = children.back()->read(); |
| 32 | if (!block) |
| 33 | return Block(); |
| 34 | |
| 35 | const ColumnRawPtrs column_ptrs(getKeyColumns(block)); |
| 36 | if (column_ptrs.empty()) |
| 37 | return block; |
| 38 | |
| 39 | const ColumnRawPtrs clearing_hint_columns(getClearingColumns(block, column_ptrs)); |
| 40 | |
| 41 | if (data.type == ClearableSetVariants::Type::EMPTY) |
| 42 | data.init(ClearableSetVariants::chooseMethod(column_ptrs, key_sizes)); |
| 43 | |
| 44 | const size_t rows = block.rows(); |
| 45 | IColumn::Filter filter(rows); |
| 46 | |
| 47 | bool has_new_data = false; |
| 48 | switch (data.type) |
| 49 | { |
| 50 | case ClearableSetVariants::Type::EMPTY: |
| 51 | break; |
| 52 | case ClearableSetVariants::Type::bitmap64: |
| 53 | break; |
| 54 | #define M(NAME) \ |
| 55 | case ClearableSetVariants::Type::NAME: \ |
| 56 | has_new_data = buildFilter(*data.NAME, column_ptrs, clearing_hint_columns, filter, rows, data); \ |
| 57 | break; |
| 58 | APPLY_FOR_SET_VARIANTS(M) |
| 59 | #undef M |
| 60 | } |
| 61 | |
| 62 | /// Just go to the next block if there isn't any new record in the current one. |
| 63 | if (!has_new_data) |
| 64 | continue; |
| 65 | |
| 66 | if (!set_size_limits.check(data.getTotalRowCount(), data.getTotalByteCount(), "DISTINCT", ErrorCodes::SET_SIZE_LIMIT_EXCEEDED)) |
| 67 | return {}; |
| 68 | |
| 69 | prev_block.block = block; |
| 70 | prev_block.clearing_hint_columns = std::move(clearing_hint_columns); |
| 71 | |
| 72 | size_t all_columns = block.columns(); |
| 73 | for (size_t i = 0; i < all_columns; ++i) |
| 74 | block.safeGetByPosition(i).column = block.safeGetByPosition(i).column->filter(filter, -1); |
| 75 | |
| 76 | return block; |
| 77 | } |
| 78 | } |
nothing calls this directly
no test coverage detected