| 525 | } |
| 526 | |
| 527 | static inline size_t onUpdateData(const Row & rows_data, Block & buffer, const std::vector<size_t> & unique_columns_index) |
| 528 | { |
| 529 | if (rows_data.size() % 2 != 0) |
| 530 | throw Exception("LOGICAL ERROR: It is a bug.", ErrorCodes::LOGICAL_ERROR); |
| 531 | |
| 532 | size_t prev_bytes = buffer.bytes(); |
| 533 | std::vector<bool> writeable_rows_mask(rows_data.size()); |
| 534 | |
| 535 | for (size_t index = 0; index < rows_data.size(); index += 2) |
| 536 | { |
| 537 | writeable_rows_mask[index + 1] = true; |
| 538 | writeable_rows_mask[index] = differenceUniqueKeys( |
| 539 | DB::get<const Tuple &>(rows_data[index]), DB::get<const Tuple &>(rows_data[index + 1]), unique_columns_index); |
| 540 | } |
| 541 | |
| 542 | for (size_t column = 0; column < buffer.columns() - 1; ++column) |
| 543 | { |
| 544 | MutableColumnPtr col_to = IColumn::mutate(std::move(buffer.getByPosition(column).column)); |
| 545 | |
| 546 | writeFieldsToColumn(*col_to, rows_data, column, writeable_rows_mask); |
| 547 | buffer.getByPosition(column).column = std::move(col_to); |
| 548 | } |
| 549 | |
| 550 | MutableColumnPtr sign_mutable_column = IColumn::mutate(std::move(buffer.getByPosition(buffer.columns() - 1).column)); |
| 551 | |
| 552 | ColumnInt8::Container & sign_column_data = assert_cast<ColumnInt8 &>(*sign_mutable_column).getData(); |
| 553 | |
| 554 | for (size_t index = 0; index < rows_data.size(); index += 2) |
| 555 | { |
| 556 | if (likely(!writeable_rows_mask[index])) |
| 557 | { |
| 558 | sign_column_data.emplace_back(0); |
| 559 | } |
| 560 | else |
| 561 | { |
| 562 | /// If the sorting keys is modified, we should cancel the old data, but this should not happen frequently |
| 563 | sign_column_data.emplace_back(1); |
| 564 | sign_column_data.emplace_back(0); |
| 565 | } |
| 566 | } |
| 567 | |
| 568 | buffer.getByPosition(buffer.columns() - 1).column = std::move(sign_mutable_column); |
| 569 | return buffer.bytes() - prev_bytes; |
| 570 | } |
| 571 | |
| 572 | void MaterializeMySQLSyncThread::onEvent(Buffers & buffers, const BinlogEventPtr & receive_event, MaterializeMetadata & metadata) |
| 573 | { |
no test coverage detected