MCPcopy Create free account
hub / github.com/actiontech/dtle / handleRowsEvent

Method handleRowsEvent

driver/mysql/binlog/binlog_reader.go:1835–2035  ·  view source on GitHub ↗
(ev *replication.BinlogEvent, rowsEvent *replication.RowsEvent,
	entriesChannel chan<- *common.EntryContext)

Source from the content-addressed store, hash-verified

1833}
1834
1835func (b *BinlogReader) handleRowsEvent(ev *replication.BinlogEvent, rowsEvent *replication.RowsEvent,
1836 entriesChannel chan<- *common.EntryContext) (err error) {
1837
1838 schemaName := string(rowsEvent.Table.Schema)
1839 tableName := string(rowsEvent.Table.Table)
1840 coordinate := b.entryContext.Entry.Coordinates.(*common.MySQLCoordinateTx)
1841 b.logger.Trace("got rowsEvent", "schema", schemaName, "table", tableName,
1842 "gno", coordinate.GNO, "flags", rowsEvent.Flags, "tableFlags", rowsEvent.Table.Flags,
1843 "nRows", len(rowsEvent.Rows))
1844
1845 dml := common.ToEventDML(ev.Header.EventType)
1846 skip, table := b.skipRowEvent(rowsEvent, dml)
1847 if skip {
1848 b.logger.Debug("skip rowsEvent", "schema", schemaName, "table", tableName,
1849 "gno", coordinate.GNO)
1850 return nil
1851 }
1852
1853 if b.sqlFilter.NoDML ||
1854 (b.sqlFilter.NoDMLDelete && dml == common.DeleteDML) ||
1855 (b.sqlFilter.NoDMLInsert && dml == common.InsertDML) ||
1856 (b.sqlFilter.NoDMLUpdate && dml == common.UpdateDML) {
1857
1858 b.logger.Debug("skipped_a_dml_event.", "type", dml, "schema", schemaName, "table", tableName)
1859 return nil
1860 }
1861
1862 if dml == common.NotDML {
1863 return fmt.Errorf("unknown DML type: %s", ev.Header.EventType.String())
1864 }
1865 dmlEvent := common.NewDataEvent(
1866 schemaName,
1867 tableName,
1868 dml,
1869 rowsEvent.ColumnCount,
1870 ev.Header.Timestamp,
1871 )
1872 if table != nil {
1873 dmlEvent.FKParent = len(table.FKChildren) > 0
1874 }
1875 dmlEvent.Flags = make([]byte, 2)
1876 binary.LittleEndian.PutUint16(dmlEvent.Flags, rowsEvent.Flags)
1877 dmlEvent.LogPos = int64(ev.Header.LogPos - ev.Header.EventSize)
1878
1879 /*originalTableColumns, _, err := b.InspectTableColumnsAndUniqueKeys(string(rowsEvent.Table.Schema), string(rowsEvent.Table.Table))
1880 if err != nil {
1881 return err
1882 }
1883 dmlEvent.OriginalTableColumns = originalTableColumns*/
1884
1885 // It is hard to calculate exact row size. We use estimation.
1886 avgRowSize := len(ev.RawData) / len(rowsEvent.Rows)
1887
1888 if table != nil && table.Table.TableRename != "" {
1889 dmlEvent.TableName = table.Table.TableRename
1890 b.logger.Debug("dml table mapping", "from", dmlEvent.TableName, "to", table.Table.TableRename)
1891 }
1892 schemaContext := b.findCurrentSchema(schemaName)

Callers 1

handleEventMethod · 0.95

Calls 9

skipRowEventMethod · 0.95
findCurrentSchemaMethod · 0.95
sendEntryMethod · 0.95
ToEventDMLFunction · 0.92
NewDataEventFunction · 0.92
EncodeTableFunction · 0.92
NewBinlogEntryFunction · 0.92
WhereTrueMethod · 0.80
StringMethod · 0.45

Tested by

no test coverage detected