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

Method DataStreamEvents

driver/oracle/extractor/log_miner.go:521–579  ·  view source on GitHub ↗
(entriesChannel chan<- *common.EntryContext)

Source from the content-addressed store, hash-verified

519}
520
521func (e *ExtractorOracle) DataStreamEvents(entriesChannel chan<- *common.EntryContext) error {
522 e.logger.Debug("start oracle. DataStreamEvents")
523
524 if e.LogMinerStream.startScn == 0 {
525 scn, err := e.LogMinerStream.oracleDB.GetCurrentSnapshotSCN()
526 if err != nil {
527 e.logger.Error("GetCurrentSnapshotSCN", "err", err)
528 return err
529 }
530 e.logger.Debug("current scn", "scn", scn)
531 e.LogMinerStream.startScn = scn
532 }
533
534 e.LogMinerStream.txCache.Handler = func(tx *LogMinerTx) error {
535 if tx.endScn <= e.LogMinerStream.committedScn {
536 e.logger.Debug("skip SQLs", "transactionId", tx.transactionId,
537 "startSCN", tx.startScn, "endSCN", tx.endScn)
538 return nil
539 }
540 atomic.AddUint32(&e.LogMinerStream.OracleTxNum, 1)
541 Entry := e.handleSQLs(tx)
542 entriesChannel <- &common.EntryContext{
543 Entry: Entry,
544 TableItems: nil,
545 OriginalSize: 0,
546 }
547 return nil
548 }
549
550 err := e.LoopLogminerRecord()
551 if err != nil {
552 e.logger.Error("LogMinerStreamLoop", "err", err)
553 return err
554 }
555 return nil
556
557 //err := e.LogMinerStream.start()
558 //if err != nil {
559 // e.logger.Error("StartLogMiner", "err", err)
560 // return err
561 //}
562 //defer e.LogMinerStream.stop()
563 //
564 //for {
565 // time.Sleep(5 * time.Second)
566 // rs, err := e.LogMinerStream.queue()
567 // if err != nil {
568 // return err
569 // }
570 //
571 // Entry := e.handleSQLs(rs)
572 // e.logger.Debug("handle SQLs", "Entry", Entry)
573 // entriesChannel <- &common.BinlogEntryContext{
574 // Entry: Entry,
575 // TableItems: nil,
576 // OriginalSize: 0,
577 // }
578 //}

Callers 1

StreamEventsMethod · 0.95

Calls 3

handleSQLsMethod · 0.95
LoopLogminerRecordMethod · 0.95
GetCurrentSnapshotSCNMethod · 0.80

Tested by

no test coverage detected