| 35 | { |
| 36 | |
| 37 | void CountingBlockOutputStream::write(const Block & block) |
| 38 | { |
| 39 | stopwatch.start(); |
| 40 | |
| 41 | stream->write(block); |
| 42 | |
| 43 | Progress local_progress(block.rows(), block.bytes(), 0); |
| 44 | |
| 45 | if (constraints_filter_stream) |
| 46 | { |
| 47 | if (const auto *counting_stream = dynamic_cast<const CheckConstraintsFilterBlockOutputStream *>(constraints_filter_stream.get())) |
| 48 | { |
| 49 | const auto & written_block = counting_stream->getWrittenBlock(); |
| 50 | local_progress.read_rows = written_block.rows(); |
| 51 | local_progress.read_bytes = written_block.bytes(); |
| 52 | } |
| 53 | } |
| 54 | |
| 55 | progress.incrementPiecewiseAtomically(local_progress); |
| 56 | |
| 57 | ProfileEvents::increment(ProfileEvents::InsertedRows, local_progress.read_rows); |
| 58 | ProfileEvents::increment(ProfileEvents::InsertedBytes, local_progress.read_bytes); |
| 59 | |
| 60 | if (process_elem) |
| 61 | process_elem->updateProgressOut(local_progress); |
| 62 | |
| 63 | if (progress_callback) |
| 64 | progress_callback(local_progress); |
| 65 | } |
| 66 | |
| 67 | void CountingBlockOutputStream::writePrefix() |
| 68 | { |
nothing calls this directly
no test coverage detected