(entry *common.DumpEntry)
| 1434 | return nil |
| 1435 | } |
| 1436 | func (e *Extractor) encodeAndSendDumpEntry(entry *common.DumpEntry) error { |
| 1437 | bs, err := entry.Marshal(nil) |
| 1438 | if err != nil { |
| 1439 | return err |
| 1440 | } |
| 1441 | e.logger.Debug("encodeAndSendDumpEntry. after Marshal", "size", len(bs)) |
| 1442 | txMsg, err := common.Compress(bs) |
| 1443 | if err != nil { |
| 1444 | return errors.Wrap(err, "common.Compress") |
| 1445 | } |
| 1446 | e.logger.Debug("encodeAndSendDumpEntry. after Compress", "size", len(txMsg)) |
| 1447 | if err := e.publish(fmt.Sprintf("%s_full", e.subject), txMsg, 0); err != nil { |
| 1448 | return err |
| 1449 | } |
| 1450 | e.mysqlContext.Stage = common.StageSendingData |
| 1451 | return nil |
| 1452 | } |
| 1453 | |
| 1454 | func (e *Extractor) Stats() (*common.TaskStatistics, error) { |
| 1455 | totalRowsCopied := atomic.LoadInt64(&e.TotalRowsCopied) |
no test coverage detected