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

Method NewRunner

driver/handle.go:180–246  ·  view source on GitHub ↗
(d *Driver)

Source from the content-addressed store, hash-verified

178}
179
180func (h *taskHandle) NewRunner(d *Driver) (runner DriverHandle, err error) {
181 ctx := &common.ExecContext{
182 Subject: h.taskConfig.JobName,
183 StateDir: d.config.DataDir,
184 }
185
186 logger := h.logger
187 // create a new log file for debug job
188 if d.config.DebugJob == h.taskConfig.JobName && hclog.LevelFromString(d.config.LogLevel) != hclog.Debug {
189 logger = setupLogger(d.config.LogFile, "dtle_debug_job.log")
190 logger.SetLevel(hclog.Debug)
191 }
192
193 switch common.TaskTypeFromString(h.taskConfig.Name) {
194 case common.TaskTypeSrc:
195 if h.driverConfig.SrcOracleConfig != nil {
196 h.logger.Debug("found oracle src", "SrcOracleConfig", h.driverConfig.SrcOracleConfig)
197 runner, err = extractor.NewExtractorOracle(ctx, h.driverConfig, logger, d.storeManager, h.waitCh, h.ctx)
198 if err != nil {
199 return nil, errors.Wrap(err, "NewExtractor")
200 }
201 } else {
202 e, err := mysql.NewExtractor(ctx, h.driverConfig, logger, d.storeManager, h.waitCh, h.ctx)
203 if err != nil {
204 return nil, errors.Wrap(err, "NewOracleExtractor")
205 }
206 runner = e
207 if h.driverConfig.TwoWaySync {
208 ctx2 := &common.ExecContext{
209 Subject: ctx.Subject + "_dtrev",
210 StateDir: d.config.DataDir,
211 }
212 cfg2 := &common.MySQLDriverConfig{
213 DtleTaskConfig: common.DtleTaskConfig{
214 DestType: "mysql",
215 },
216 }
217 e.RevApplier, err = mysql.NewApplier(ctx2, cfg2, logger, d.storeManager, d.config.NatsAdvertise,
218 h.waitCh, d.eventer, h.taskConfig, h.ctx)
219 }
220
221 }
222 case common.TaskTypeDest:
223 h.logger.Debug("found dest", "allConfig", h.driverConfig)
224 switch strings.ToLower(h.driverConfig.DestType) {
225 case "kafka":
226 runner, err = kafka.NewKafkaRunner(ctx, logger,
227 d.storeManager, d.config.NatsAdvertise, h.waitCh, h.ctx)
228 if err != nil {
229 return nil, errors.Wrap(err, "NewKafkaRunner")
230 }
231 case "mysql":
232 runner, err = mysql.NewApplier(ctx, h.driverConfig, logger, d.storeManager,
233 d.config.NatsAdvertise, h.waitCh, d.eventer, h.taskConfig, h.ctx)
234 if err != nil {
235 return nil, errors.Wrap(err, "NewApplier")
236 }
237 case "":

Callers 1

resumeTaskMethod · 0.95

Calls 6

TaskTypeFromStringFunction · 0.92
NewExtractorOracleFunction · 0.92
NewExtractorFunction · 0.92
NewApplierFunction · 0.92
NewKafkaRunnerFunction · 0.92
setupLoggerFunction · 0.85

Tested by

no test coverage detected