(ev *replication.BinlogEvent, rowsEvent *replication.RowsEvent, entriesChannel chan<- *common.EntryContext)
| 1833 | } |
| 1834 | |
| 1835 | func (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) |
no test coverage detected