-
Notifications
You must be signed in to change notification settings - Fork 25k
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
[Transform] simplify class structure of indexer (#46306)
simplify transform task and indexer - remove redundant transform id - moving client data frame indexer (and builder) into a separate file
- Loading branch information
Hendrik Muhs
committed
Sep 6, 2019
1 parent
bb7bff5
commit c2194aa
Showing
6 changed files
with
739 additions
and
677 deletions.
There are no files selected for viewing
591 changes: 591 additions & 0 deletions
591
...me/src/main/java/org/elasticsearch/xpack/dataframe/transforms/ClientDataFrameIndexer.java
Large diffs are not rendered by default.
Oops, something went wrong.
123 changes: 123 additions & 0 deletions
123
...main/java/org/elasticsearch/xpack/dataframe/transforms/ClientDataFrameIndexerBuilder.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,123 @@ | ||
/* | ||
* Copyright Elasticsearch B.V. and/or licensed to Elasticsearch B.V. under one | ||
* or more contributor license agreements. Licensed under the Elastic License; | ||
* you may not use this file except in compliance with the Elastic License. | ||
*/ | ||
|
||
package org.elasticsearch.xpack.dataframe.transforms; | ||
|
||
import org.elasticsearch.client.Client; | ||
import org.elasticsearch.xpack.core.dataframe.transforms.DataFrameIndexerPosition; | ||
import org.elasticsearch.xpack.core.dataframe.transforms.DataFrameIndexerTransformStats; | ||
import org.elasticsearch.xpack.core.dataframe.transforms.DataFrameTransformCheckpoint; | ||
import org.elasticsearch.xpack.core.dataframe.transforms.DataFrameTransformConfig; | ||
import org.elasticsearch.xpack.core.dataframe.transforms.DataFrameTransformProgress; | ||
import org.elasticsearch.xpack.core.indexing.IndexerState; | ||
import org.elasticsearch.xpack.dataframe.checkpoint.CheckpointProvider; | ||
import org.elasticsearch.xpack.dataframe.checkpoint.DataFrameTransformsCheckpointService; | ||
import org.elasticsearch.xpack.dataframe.notifications.DataFrameAuditor; | ||
import org.elasticsearch.xpack.dataframe.persistence.DataFrameTransformsConfigManager; | ||
|
||
import java.util.Map; | ||
import java.util.concurrent.atomic.AtomicReference; | ||
|
||
class ClientDataFrameIndexerBuilder { | ||
private Client client; | ||
private DataFrameTransformsConfigManager transformsConfigManager; | ||
private DataFrameTransformsCheckpointService transformsCheckpointService; | ||
private DataFrameAuditor auditor; | ||
private Map<String, String> fieldMappings; | ||
private DataFrameTransformConfig transformConfig; | ||
private DataFrameIndexerTransformStats initialStats; | ||
private IndexerState indexerState = IndexerState.STOPPED; | ||
private DataFrameIndexerPosition initialPosition; | ||
private DataFrameTransformProgress progress; | ||
private DataFrameTransformCheckpoint lastCheckpoint; | ||
private DataFrameTransformCheckpoint nextCheckpoint; | ||
|
||
ClientDataFrameIndexerBuilder() { | ||
this.initialStats = new DataFrameIndexerTransformStats(); | ||
} | ||
|
||
ClientDataFrameIndexer build(DataFrameTransformTask parentTask) { | ||
CheckpointProvider checkpointProvider = transformsCheckpointService.getCheckpointProvider(transformConfig); | ||
|
||
return new ClientDataFrameIndexer(this.transformsConfigManager, | ||
checkpointProvider, | ||
new AtomicReference<>(this.indexerState), | ||
this.initialPosition, | ||
this.client, | ||
this.auditor, | ||
this.initialStats, | ||
this.transformConfig, | ||
this.fieldMappings, | ||
this.progress, | ||
this.lastCheckpoint, | ||
this.nextCheckpoint, | ||
parentTask); | ||
} | ||
|
||
ClientDataFrameIndexerBuilder setClient(Client client) { | ||
this.client = client; | ||
return this; | ||
} | ||
|
||
ClientDataFrameIndexerBuilder setTransformsConfigManager(DataFrameTransformsConfigManager transformsConfigManager) { | ||
this.transformsConfigManager = transformsConfigManager; | ||
return this; | ||
} | ||
|
||
ClientDataFrameIndexerBuilder setTransformsCheckpointService(DataFrameTransformsCheckpointService transformsCheckpointService) { | ||
this.transformsCheckpointService = transformsCheckpointService; | ||
return this; | ||
} | ||
|
||
ClientDataFrameIndexerBuilder setAuditor(DataFrameAuditor auditor) { | ||
this.auditor = auditor; | ||
return this; | ||
} | ||
|
||
ClientDataFrameIndexerBuilder setFieldMappings(Map<String, String> fieldMappings) { | ||
this.fieldMappings = fieldMappings; | ||
return this; | ||
} | ||
|
||
ClientDataFrameIndexerBuilder setTransformConfig(DataFrameTransformConfig transformConfig) { | ||
this.transformConfig = transformConfig; | ||
return this; | ||
} | ||
|
||
DataFrameTransformConfig getTransformConfig() { | ||
return this.transformConfig; | ||
} | ||
|
||
ClientDataFrameIndexerBuilder setInitialStats(DataFrameIndexerTransformStats initialStats) { | ||
this.initialStats = initialStats; | ||
return this; | ||
} | ||
|
||
ClientDataFrameIndexerBuilder setIndexerState(IndexerState indexerState) { | ||
this.indexerState = indexerState; | ||
return this; | ||
} | ||
|
||
ClientDataFrameIndexerBuilder setInitialPosition(DataFrameIndexerPosition initialPosition) { | ||
this.initialPosition = initialPosition; | ||
return this; | ||
} | ||
|
||
ClientDataFrameIndexerBuilder setProgress(DataFrameTransformProgress progress) { | ||
this.progress = progress; | ||
return this; | ||
} | ||
|
||
ClientDataFrameIndexerBuilder setLastCheckpoint(DataFrameTransformCheckpoint lastCheckpoint) { | ||
this.lastCheckpoint = lastCheckpoint; | ||
return this; | ||
} | ||
|
||
ClientDataFrameIndexerBuilder setNextCheckpoint(DataFrameTransformCheckpoint nextCheckpoint) { | ||
this.nextCheckpoint = nextCheckpoint; | ||
return this; | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.