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

Method Run

driver/oracle/extractor/extractor_oracle.go:156–255  ·  view source on GitHub ↗

Run executes the complete extract logic.

()

Source from the content-addressed store, hash-verified

154
155// Run executes the complete extract logic.
156func (e *ExtractorOracle) Run() {
157 var err error
158
159 // PutConfig before WatchNats
160 err = e.storeManager.PutConfig(e.subject, e.mysqlContext)
161 if err != nil {
162 e.onError(common.TaskStateDead, errors.Wrap(err, "PutConfig"))
163 return
164 }
165
166 e.logger.Info("src watch Nats")
167 e.natsAddr, err = e.storeManager.SrcWatchNats(e.subject, e.shutdownCh, func(err error) {
168 e.onError(common.TaskStateDead, err)
169 })
170 if err != nil {
171 e.onError(common.TaskStateDead, errors.Wrap(err, "SrcWatchNats"))
172 return
173 }
174 // init nats
175 e.logger.Info("initNatsPubClient")
176 e.logger.Debug("begin Connect nats server", "NatAddr", e.natsAddr)
177 sc, err := gonats.Connect(e.natsAddr)
178 if err != nil {
179 e.logger.Error("cannot connect nats server", "natsAddr", e.natsAddr, "err", err)
180 e.onError(common.TaskStateDead, errors.Wrap(err, "Connect Nats"))
181 return
182 }
183 e.logger.Info("Connect nats server", "natsAddr", e.natsAddr)
184 e.natsConn = sc
185
186 _, err = e.natsConn.Subscribe(fmt.Sprintf("%s_control2", e.subject), func(m *gonats.Msg) {
187 if m.Data == nil {
188 e.onError(common.TaskStateDead, fmt.Errorf("zero-byte control msg"))
189 return
190 }
191
192 ctrlMsg := &common.ControlMsg{}
193 _, err := ctrlMsg.Unmarshal(m.Data)
194 if err != nil {
195 e.onError(common.TaskStateDead, fmt.Errorf("failed to unmarshal a control msg"))
196 return
197 }
198
199 switch ctrlMsg.Type {
200 case common.ControlMsgError:
201 e.onError(common.TaskStateDead, fmt.Errorf("applier error: %v", ctrlMsg.Msg))
202 return
203 }
204 })
205 if err != nil {
206 e.onError(common.TaskStateDead, errors.Wrap(err, "Subscribe control2"))
207 return
208 }
209
210 e.initDBConnections()
211
212 startSCN, committedSCN, err := e.calculateSCNPos()
213 if err != nil {

Callers

nothing calls this directly

Calls 11

onErrorMethod · 0.95
UnmarshalMethod · 0.95
initDBConnectionsMethod · 0.95
calculateSCNPosMethod · 0.95
oracleDumpMethod · 0.95
sendFullCompleteMethod · 0.95
initiateStreamingMethod · 0.95
NewLogMinerStreamFunction · 0.85
PutConfigMethod · 0.80
SrcWatchNatsMethod · 0.80

Tested by

no test coverage detected