Skip to content

Commit

Permalink
Refactor CDCJob and MigrationJob (#29361)
Browse files Browse the repository at this point in the history
* Refactor CDCJob

* Refactor MigrationJob
  • Loading branch information
terrymanu authored Dec 11, 2023
1 parent 4cb0cd0 commit bc68578
Show file tree
Hide file tree
Showing 2 changed files with 9 additions and 19 deletions.
Original file line number Diff line number Diff line change
Expand Up @@ -71,24 +71,20 @@
@Slf4j
public final class CDCJob extends AbstractInseparablePipelineJob<CDCJobItemContext> {

@Getter
private final PipelineSink sink;
private final CDCJobAPI jobAPI = (CDCJobAPI) TypedSPILoader.getService(TransmissionJobAPI.class, "STREAMING");

private final CDCJobAPI jobAPI;
private final PipelineJobItemManager<TransmissionJobItemProgress> jobItemManager = new PipelineJobItemManager<>(new CDCJobType().getYamlJobItemProgressSwapper());

private final PipelineJobItemManager<TransmissionJobItemProgress> jobItemManager;
private final PipelineProcessConfigurationPersistService processConfigPersistService = new PipelineProcessConfigurationPersistService();

private final PipelineProcessConfigurationPersistService processConfigPersistService;
private final CDCJobPreparer jobPreparer = new CDCJobPreparer();

private final CDCJobPreparer jobPreparer;
@Getter
private final PipelineSink sink;

public CDCJob(final PipelineSink sink) {
super(new PipelineJobRunnerManager(new CDCJobRunnerCleaner(sink)));
this.sink = sink;
jobAPI = (CDCJobAPI) TypedSPILoader.getService(TransmissionJobAPI.class, "STREAMING");
jobItemManager = new PipelineJobItemManager<>(new CDCJobType().getYamlJobItemProgressSwapper());
processConfigPersistService = new PipelineProcessConfigurationPersistService();
jobPreparer = new CDCJobPreparer();
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -63,18 +63,12 @@
@Slf4j
public final class MigrationJob extends AbstractSeparablePipelineJob<MigrationJobItemContext> {

private final PipelineJobItemManager<TransmissionJobItemProgress> jobItemManager;
private final PipelineJobItemManager<TransmissionJobItemProgress> jobItemManager = new PipelineJobItemManager<>(new MigrationJobType().getYamlJobItemProgressSwapper());

private final PipelineProcessConfigurationPersistService processConfigPersistService;
private final PipelineProcessConfigurationPersistService processConfigPersistService = new PipelineProcessConfigurationPersistService();

// Shared by all sharding items
private final MigrationJobPreparer jobPreparer;

public MigrationJob() {
jobItemManager = new PipelineJobItemManager<>(new MigrationJobType().getYamlJobItemProgressSwapper());
processConfigPersistService = new PipelineProcessConfigurationPersistService();
jobPreparer = new MigrationJobPreparer();
}
private final MigrationJobPreparer jobPreparer = new MigrationJobPreparer();

@Override
protected MigrationJobItemContext buildJobItemContext(final ShardingContext shardingContext) {
Expand Down

0 comments on commit bc68578

Please sign in to comment.