| 570 | } |
| 571 | |
| 572 | void MaterializeMySQLSyncThread::onEvent(Buffers & buffers, const BinlogEventPtr & receive_event, MaterializeMetadata & metadata) |
| 573 | { |
| 574 | if (receive_event->type() == MYSQL_WRITE_ROWS_EVENT) |
| 575 | { |
| 576 | WriteRowsEvent & write_rows_event = static_cast<WriteRowsEvent &>(*receive_event); |
| 577 | Buffers::BufferAndUniqueColumnsPtr buffer = buffers.getTableDataBuffer(write_rows_event.table, getContext()); |
| 578 | size_t bytes = onWriteOrDeleteData<0>(write_rows_event.rows, buffer->first); |
| 579 | buffers.add(buffer->first.rows(), buffer->first.bytes(), write_rows_event.rows.size(), bytes); |
| 580 | } |
| 581 | else if (receive_event->type() == MYSQL_UPDATE_ROWS_EVENT) |
| 582 | { |
| 583 | UpdateRowsEvent & update_rows_event = static_cast<UpdateRowsEvent &>(*receive_event); |
| 584 | Buffers::BufferAndUniqueColumnsPtr buffer = buffers.getTableDataBuffer(update_rows_event.table, getContext()); |
| 585 | size_t bytes = onUpdateData(update_rows_event.rows, buffer->first, buffer->second); |
| 586 | buffers.add(buffer->first.rows(), buffer->first.bytes(), update_rows_event.rows.size(), bytes); |
| 587 | } |
| 588 | else if (receive_event->type() == MYSQL_DELETE_ROWS_EVENT) |
| 589 | { |
| 590 | DeleteRowsEvent & delete_rows_event = static_cast<DeleteRowsEvent &>(*receive_event); |
| 591 | Buffers::BufferAndUniqueColumnsPtr buffer = buffers.getTableDataBuffer(delete_rows_event.table, getContext()); |
| 592 | size_t bytes = onWriteOrDeleteData<1>(delete_rows_event.rows, buffer->first); |
| 593 | buffers.add(buffer->first.rows(), buffer->first.bytes(), delete_rows_event.rows.size(), bytes); |
| 594 | } |
| 595 | else if (receive_event->type() == MYSQL_QUERY_EVENT) |
| 596 | { |
| 597 | QueryEvent & query_event = static_cast<QueryEvent &>(*receive_event); |
| 598 | Position position_before_ddl; |
| 599 | position_before_ddl.update(metadata.binlog_position, metadata.binlog_file, metadata.executed_gtid_set); |
| 600 | metadata.transaction(position_before_ddl, [&]() { buffers.commit(getContext(), position_before_ddl); }); |
| 601 | metadata.transaction(client.getPosition(),[&](){ executeDDLAtomic(query_event); }); |
| 602 | } |
| 603 | else |
| 604 | { |
| 605 | /// MYSQL_UNHANDLED_EVENT |
| 606 | if (receive_event->header.type == ROTATE_EVENT) |
| 607 | { |
| 608 | /// Some behaviors(such as changing the value of "binlog_checksum") rotate the binlog file. |
| 609 | /// To ensure that the synchronization continues, we need to handle these events |
| 610 | metadata.fetchMasterVariablesValue(pool.get()); |
| 611 | client.setBinlogChecksum(metadata.binlog_checksum); |
| 612 | } |
| 613 | else if (receive_event->header.type != HEARTBEAT_EVENT) |
| 614 | { |
| 615 | const auto & dump_event_message = [&]() |
| 616 | { |
| 617 | WriteBufferFromOwnString buf; |
| 618 | receive_event->dump(buf); |
| 619 | return buf.str(); |
| 620 | }; |
| 621 | |
| 622 | LOG_DEBUG(log, "Skip MySQL event: \n {}", dump_event_message()); |
| 623 | } |
| 624 | } |
| 625 | } |
| 626 | |
| 627 | void MaterializeMySQLSyncThread::commitPosition(const MySQLReplication::Position & position) |
| 628 | { |
nothing calls this directly
no test coverage detected