Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 6 additions & 1 deletion link-move/src/main/java/com/nhl/link/move/Execution.java
Original file line number Diff line number Diff line change
Expand Up @@ -18,15 +18,17 @@ public class Execution implements AutoCloseable {
protected Map<String, ?> parameters;
protected Map<String, Object> attributes;
protected ExecutionStats stats;
private boolean dryRun;

public Execution(String name, Map<String, ?> params) {
public Execution(String name, Map<String, ?> params, boolean dryRun) {
this.name = name;
this.parameters = params;

// a "parallel" execution should have this turned into a ConcurrentMap
this.attributes = new HashMap<>();

this.stats = new ExecutionStats();
this.dryRun = dryRun;

stats.executionStarted();
}
Expand Down Expand Up @@ -101,4 +103,7 @@ public ExecutionStats getStats() {
return stats;
}

public boolean isDryRun() {
return dryRun;
}
}
47 changes: 47 additions & 0 deletions link-move/src/main/java/com/nhl/link/move/LmTask.java
Original file line number Diff line number Diff line change
Expand Up @@ -44,4 +44,51 @@ public interface LmTask {
*/
Execution run(SyncToken token, Map<String, ?> params);

/**
* Executes the task returning {@link Execution} object that can be used by
* the caller to analyze the results. Currently all task implementations are
* synchronous, so this method returns only on task completion.
*
* Any changes will not be committed to the target datasource.
*
* @since 2.3
*/
Execution dryRun();

/**
* Executes the task with a map of parameters returning {@link Execution}
* object that can be used by the caller to analyze the results. Currently
* all task implementations are synchronous, so this method returns only on
* task completion.
*
* Any changes will not be committed to the target datasource.
*
* @since 2.3
*/
Execution dryRun(Map<String, ?> params);

/**
* Executes the task with a map of parameters returning {@link Execution}
* object that can be used by the caller to analyze the results. Currently
* all task implementations are synchronous, so this method returns only on
* task completion.
*
* Any changes will not be committed to the target datasource.
*
* @since 2.3
*/
Execution dryRun(SyncToken token);

/**
* Executes the task with a map of parameters returning {@link Execution}
* object that can be used by the caller to analyze the results. Currently
* all task implementations are synchronous, so this method returns only on
* task completion.
*
* Any changes will not be committed to the target datasource.
*
* @since 2.3
*/
Execution dryRun(SyncToken token, Map<String, ?> params);

}
Original file line number Diff line number Diff line change
Expand Up @@ -21,22 +21,62 @@ public BaseTask(ITokenManager tokenManager) {
this.tokenManager = tokenManager;
}

@Override
public abstract Execution run(Map<String, ?> params);
protected abstract void run(Execution execution, Map<String, ?> params);

protected abstract String name();

@Override
public Execution run() {
return run(Collections.<String, Object> emptyMap());
return run(Collections.emptyMap(), false);
}

@Override
public Execution run(Map<String, ?> params) {
return run(params, false);
}

@Override
public Execution run(SyncToken token) {
return run(token, Collections.<String, Object> emptyMap());
return run(token, Collections.emptyMap(), false);
}

@Override
public Execution run(SyncToken token, Map<String, ?> params) {
return run(token, params, false);
}

@Override
public Execution dryRun() {
return run(Collections.emptyMap(), true);
}

@Override
public Execution dryRun(Map<String, ?> params) {
return run(params, true);
}

@Override
public Execution dryRun(SyncToken token) {
return run(token, Collections.emptyMap(), true);
}

@Override
public Execution dryRun(SyncToken token, Map<String, ?> params) {
return run(token, params, true);
}

private Execution run(Map<String, ?> params, boolean dryRun) {
if (params == null) {
throw new NullPointerException("Null params");
}

try (Execution execution = new Execution(name(), params, dryRun)) {
run(execution, params);
return execution;
}
}

private Execution run(SyncToken token, Map<String, ?> params, boolean dryRun) {
Map<String, Object> combinedParams = new HashMap<>();

SyncToken startToken = tokenManager.previousToken(token);
Expand All @@ -45,13 +85,12 @@ public Execution run(SyncToken token, Map<String, ?> params) {

combinedParams.putAll(params);

Execution exec = run(combinedParams);
Execution exec = run(combinedParams, dryRun);

// if we ever start using delayed executions, token should be
// saved inside the execution...
tokenManager.saveToken(token);

return exec;
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -14,4 +14,6 @@ public abstract class BaseTaskBuilder {
public BaseTaskBuilder() {
this.batchSize = DEFAULT_BATCH_SIZE;
}


}
Original file line number Diff line number Diff line change
Expand Up @@ -41,14 +41,16 @@ public CreateOrUpdateSegmentProcessor(RowConverter rowConverter, SourceMapper ma
}

public void process(Execution exec, CreateOrUpdateSegment<T> segment) {

// execute create-or-update pipeline stages
convertSrc(exec, segment);
mapSrc(exec, segment);
matchTarget(exec, segment);
mapToTarget(exec, segment);
mergeToTarget(exec, segment);
commitTarget(exec, segment);

if (!exec.isDryRun()) {
commitTarget(exec, segment);
}
}

private void convertSrc(Execution exec, CreateOrUpdateSegment<T> segment) {
Expand Down
Original file line number Diff line number Diff line change
@@ -1,11 +1,5 @@
package com.nhl.link.move.runtime.task.createorupdate;

import java.util.List;
import java.util.Map;

import org.apache.cayenne.DataObject;
import org.apache.cayenne.ObjectContext;

import com.nhl.link.move.CountingRowReader;
import com.nhl.link.move.Execution;
import com.nhl.link.move.Row;
Expand All @@ -18,6 +12,11 @@
import com.nhl.link.move.runtime.extractor.IExtractorService;
import com.nhl.link.move.runtime.task.BaseTask;
import com.nhl.link.move.runtime.token.ITokenManager;
import org.apache.cayenne.DataObject;
import org.apache.cayenne.ObjectContext;

import java.util.List;
import java.util.Map;

/**
* A task that reads streamed source data and creates/updates records in a
Expand Down Expand Up @@ -46,22 +45,17 @@ public CreateOrUpdateTask(ExtractorName extractorName, int batchSize, ITargetCay
}

@Override
public Execution run(Map<String, ?> params) {
protected void run(Execution execution, Map<String, ?> params) {
BatchProcessor<Row> batchProcessor = createBatchProcessor(execution);

if (params == null) {
throw new NullPointerException("Null params");
try (RowReader data = getRowReader(execution, params)) {
BatchRunner.create(batchProcessor).withBatchSize(batchSize).run(data);
}
}

try (Execution execution = new Execution("CreateOrUpdateTask:" + extractorName, params);) {

BatchProcessor<Row> batchProcessor = createBatchProcessor(execution);

try (RowReader data = getRowReader(execution, params)) {
BatchRunner.create(batchProcessor).withBatchSize(batchSize).run(data);
}

return execution;
}
@Override
protected String name() {
return "CreateOrUpdateTask:" + extractorName;
}

protected BatchProcessor<Row> createBatchProcessor(final Execution execution) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,10 @@ public void process(Execution exec, DeleteSegment<T> segment) {
extractSourceKeys(exec, segment);
filterMissingTargets(exec, segment);
deleteTarget(segment);
commitTarget(segment);

if (!exec.isDryRun()) {
commitTarget(segment);
}
}

private void mapTarget(Execution exec, DeleteSegment<T> segment) {
Expand Down
Original file line number Diff line number Diff line change
@@ -1,20 +1,18 @@
package com.nhl.link.move.runtime.task.delete;

import java.util.List;
import java.util.Map;

import org.apache.cayenne.DataObject;
import org.apache.cayenne.ObjectContext;
import org.apache.cayenne.ResultIterator;
import org.apache.cayenne.exp.Expression;
import org.apache.cayenne.query.ObjectSelect;

import com.nhl.link.move.Execution;
import com.nhl.link.move.batch.BatchProcessor;
import com.nhl.link.move.batch.BatchRunner;
import com.nhl.link.move.runtime.cayenne.ITargetCayenneService;
import com.nhl.link.move.runtime.task.BaseTask;
import com.nhl.link.move.runtime.token.ITokenManager;
import org.apache.cayenne.DataObject;
import org.apache.cayenne.ObjectContext;
import org.apache.cayenne.ResultIterator;
import org.apache.cayenne.exp.Expression;
import org.apache.cayenne.query.ObjectSelect;

import java.util.Map;

/**
* A task that allows to delete target objects not present in the source.
Expand Down Expand Up @@ -45,41 +43,35 @@ public DeleteTask(String extractorName, int batchSize, Class<T> type, Expression
}

@Override
public Execution run(Map<String, ?> params) {

try (Execution execution = new Execution("DeleteTask:" + extractorName, params);) {

BatchProcessor<T> batchProcessor = createBatchProcessor(execution);
ResultIterator<T> data = createTargetSelect();

try {
BatchRunner.create(batchProcessor).withBatchSize(batchSize).run(data);
} finally {
// TODO: ResultIterator should be made Closeable in Cayenne,
// then we can use try-with-resources.
data.close();
}

return execution;
protected void run(Execution execution, Map<String, ?> params) {
BatchProcessor<T> batchProcessor = createBatchProcessor(execution);
ResultIterator<T> data = createTargetSelect();

try {
BatchRunner.create(batchProcessor).withBatchSize(batchSize).run(data);
} finally {
// TODO: ResultIterator should be made Closeable in Cayenne,
// then we can use try-with-resources.
data.close();
}
}

@Override
protected String name() {
return "DeleteTask:" + extractorName;
}

protected ResultIterator<T> createTargetSelect() {
ObjectSelect<T> query = ObjectSelect.query(type).where(targetFilter);
return targetCayenneService.newContext().iterator(query);
}

protected BatchProcessor<T> createBatchProcessor(final Execution execution) {
return new BatchProcessor<T>() {

@Override
public void process(List<T> segment) {

// executing in the select context..
ObjectContext context = segment.get(0).getObjectContext();
processor.process(execution, new DeleteSegment<>(context, segment));
}
};
return segment -> {
// executing in the select context..
ObjectContext context = segment.get(0).getObjectContext();
processor.process(execution, new DeleteSegment<>(context, segment));
};
}

}
Loading