Skip to content

Commit 5ac4b9c

Browse files
zhilingcfeast-ci-bot
authored andcommitted
Change error store to be part of configuration instead (#39)
* Change error store to be part of configuration instead * Revert to provider providing list of error stores * merge master, reformat, align mocks stores to lowercase type * Remove extra parantheses * Format code to follow style * Format code to follow style
1 parent f389983 commit 5ac4b9c

19 files changed

Lines changed: 286 additions & 116 deletions

File tree

core/src/main/java/feast/core/config/AppConfig.java

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,9 @@ public ImportJobDefaults getImportJobDefaults(
3232
@Value("${feast.jobs.runner}") String runner,
3333
@Value("${feast.jobs.options}") String options,
3434
@Value("${feast.jobs.executable}") String executable,
35-
@Value("${feast.jobs.errorsStoreId}") String errorsStoreId) {
36-
return new ImportJobDefaults(coreApiUri, runner, options, executable, errorsStoreId);
35+
@Value("${feast.jobs.errorsStoreType}") String errorsStoreType,
36+
@Value("${feast.jobs.errorsStoreOptions}") String errorsStoreOptions) {
37+
return new ImportJobDefaults(
38+
coreApiUri, runner, options, executable, errorsStoreType, errorsStoreOptions);
3739
}
3840
}

core/src/main/java/feast/core/config/ImportJobDefaults.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@ public class ImportJobDefaults {
3232
private String runner;
3333
private String importJobOptions;
3434
private String executable;
35-
private String errorsStoreId;
35+
private String errorsStoreType;
36+
private String errorsStoreOptions;
3637
}
3738

core/src/main/java/feast/core/service/JobExecutionService.java

Lines changed: 13 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -28,24 +28,26 @@
2828
import feast.core.model.JobStatus;
2929
import feast.core.util.TypeConversion;
3030
import feast.specs.ImportSpecProto.ImportSpec;
31-
import lombok.extern.slf4j.Slf4j;
32-
import org.springframework.beans.factory.annotation.Autowired;
33-
import org.springframework.stereotype.Service;
34-
3531
import java.io.BufferedReader;
3632
import java.io.InputStreamReader;
3733
import java.time.Instant;
38-
import java.util.*;
34+
import java.util.ArrayList;
35+
import java.util.Base64;
36+
import java.util.List;
37+
import java.util.Map;
38+
import java.util.Optional;
3939
import java.util.regex.Pattern;
40+
import lombok.extern.slf4j.Slf4j;
41+
import org.springframework.beans.factory.annotation.Autowired;
42+
import org.springframework.stereotype.Service;
4043

4144
@Slf4j
4245
@Service
4346
public class JobExecutionService {
44-
private static final int SLEEP_MS = 10;
45-
private static final Pattern JOB_EXT_ID_PREFIX_REGEX = Pattern.compile(".*FeastImportJobId:.*");
4647

4748
public static final String JOB_PREFIX_DEFAULT = "feastimport";
48-
49+
private static final int SLEEP_MS = 10;
50+
private static final Pattern JOB_EXT_ID_PREFIX_REGEX = Pattern.compile(".*FeastImportJobId:.*");
4951
private JobInfoRepository jobInfoRepository;
5052
private ImportJobDefaults defaults;
5153

@@ -106,9 +108,6 @@ public SubmitImportJobResponse submitJob(ImportSpec importSpec, String jobPrefix
106108

107109
/**
108110
* Update a given job's status
109-
*
110-
* @param jobId
111-
* @param status
112111
*/
113112
public void updateJobStatus(String jobId, JobStatus status) {
114113
Optional<JobInfo> jobRecordOptional = jobInfoRepository.findById(jobId);
@@ -121,9 +120,6 @@ public void updateJobStatus(String jobId, JobStatus status) {
121120

122121
/**
123122
* Update a given job's external id
124-
*
125-
* @param jobId
126-
* @param jobExtId
127123
*/
128124
public void updateJobExtId(String jobId, String jobExtId) {
129125
Optional<JobInfo> jobRecordOptional = jobInfoRepository.findById(jobId);
@@ -137,8 +133,6 @@ public void updateJobExtId(String jobId, String jobExtId) {
137133
/**
138134
* Builds the command to execute the ingestion job
139135
*
140-
* @param importSpec
141-
* @param jobId
142136
* @return configured ProcessBuilder
143137
*/
144138
public ProcessBuilder getProcessBuilder(ImportSpec importSpec, String jobId) {
@@ -153,7 +147,8 @@ public ProcessBuilder getProcessBuilder(ImportSpec importSpec, String jobId) {
153147
commands.add(
154148
option("importSpecBase64", Base64.getEncoder().encodeToString(importSpec.toByteArray())));
155149
commands.add(option("coreApiUri", defaults.getCoreApiUri()));
156-
commands.add(option("errorsStoreId", defaults.getErrorsStoreId()));
150+
commands.add(option("errorsStoreType", defaults.getErrorsStoreType()));
151+
commands.add(option("errorsStoreOptions", defaults.getErrorsStoreOptions()));
157152
options.forEach((k, v) -> commands.add(option(k, v)));
158153
return new ProcessBuilder(commands);
159154
}
@@ -170,7 +165,7 @@ private String option(String key, String value) {
170165
*/
171166
public String runProcess(Process p) {
172167
try (BufferedReader outputStream =
173-
new BufferedReader(new InputStreamReader(p.getInputStream()));
168+
new BufferedReader(new InputStreamReader(p.getInputStream()));
174169
BufferedReader errorsStream =
175170
new BufferedReader(new InputStreamReader(p.getErrorStream()))) {
176171
String extId = "";

core/src/main/resources/application.properties

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,8 @@ feast.jobs.coreUri=${CORE_API_URI:localhost:8433}
2020
feast.jobs.runner=${JOB_RUNNER:DirectRunner}
2121
feast.jobs.options=${JOB_OPTIONS:{}}
2222
feast.jobs.executable=${JOB_EXECUTABLE:feast-ingestion.jar}
23-
feast.jobs.errorsStoreId=${JOB_ERRORS_STORE_ID:STDOUT}
23+
feast.jobs.errorsStoreType=${JOB_ERRORS_STORE_TYPE:stdout}
24+
feast.jobs.errorsStoreOptions=${JOB_ERRORS_STORE_OPTIONS:{}}
2425

2526
feast.jobs.dataflow.projectId = ${DATAFLOW_PROJECT_ID:}
2627
feast.jobs.dataflow.location = ${DATAFLOW_LOCATION:}

core/src/test/java/feast/core/service/JobExecutionServiceTest.java

Lines changed: 22 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -17,12 +17,25 @@
1717

1818
package feast.core.service;
1919

20+
import static org.hamcrest.MatcherAssert.assertThat;
21+
import static org.hamcrest.Matchers.equalTo;
22+
import static org.mockito.Mockito.verify;
23+
import static org.mockito.Mockito.when;
24+
import static org.mockito.MockitoAnnotations.initMocks;
25+
import static org.mockito.internal.verification.VerificationModeFactory.times;
26+
2027
import com.google.common.collect.Lists;
2128
import feast.core.config.ImportJobDefaults;
2229
import feast.core.dao.JobInfoRepository;
2330
import feast.core.model.JobInfo;
2431
import feast.core.model.JobStatus;
2532
import feast.specs.ImportSpecProto.ImportSpec;
33+
import java.io.ByteArrayInputStream;
34+
import java.io.IOException;
35+
import java.io.InputStream;
36+
import java.nio.charset.StandardCharsets;
37+
import java.util.List;
38+
import java.util.Optional;
2639
import org.junit.Before;
2740
import org.junit.Rule;
2841
import org.junit.Test;
@@ -31,25 +44,13 @@
3144
import org.mockito.Mock;
3245
import org.mockito.Mockito;
3346

34-
import java.io.ByteArrayInputStream;
35-
import java.io.IOException;
36-
import java.io.InputStream;
37-
import java.nio.charset.StandardCharsets;
38-
import java.util.List;
39-
import java.util.Optional;
40-
41-
import static org.hamcrest.MatcherAssert.assertThat;
42-
import static org.hamcrest.Matchers.equalTo;
43-
import static org.mockito.Mockito.verify;
44-
import static org.mockito.Mockito.when;
45-
import static org.mockito.MockitoAnnotations.initMocks;
46-
import static org.mockito.internal.verification.VerificationModeFactory.times;
47-
4847
public class JobExecutionServiceTest {
49-
private ImportJobDefaults defaults;
50-
@Mock JobInfoRepository jobInfoRepository;
5148

52-
@Rule public final ExpectedException expectedException = ExpectedException.none();
49+
@Rule
50+
public final ExpectedException expectedException = ExpectedException.none();
51+
@Mock
52+
JobInfoRepository jobInfoRepository;
53+
private ImportJobDefaults defaults;
5354

5455
@Before
5556
public void setUp() {
@@ -60,7 +61,8 @@ public void setUp() {
6061
"DirectRunner",
6162
"{\"key\":\"value\"}",
6263
"ingestion.jar",
63-
"STDOUT");
64+
"STDOUT",
65+
"{}");
6466
}
6567

6668
@Test
@@ -77,7 +79,8 @@ public void shouldBuildProcessBuilderWithCorrectOptions() {
7779
"--runner=DirectRunner",
7880
"--importSpecBase64=CgRmaWxl",
7981
"--coreApiUri=localhost:8080",
80-
"--errorsStoreId=STDOUT",
82+
"--errorsStoreType=STDOUT",
83+
"--errorsStoreOptions={}",
8184
"--key=value");
8285
assertThat(pb.command(), equalTo(expected));
8386
}

ingestion/src/main/java/feast/ingestion/boot/ImportJobModule.java

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -37,7 +37,9 @@
3737
import java.util.List;
3838
import org.apache.beam.sdk.options.PipelineOptions;
3939

40-
/** An ImportJobModule is a Guice module for creating dependency injection bindings. */
40+
/**
41+
* An ImportJobModule is a Guice module for creating dependency injection bindings.
42+
*/
4143
public class ImportJobModule extends AbstractModule {
4244

4345
private final ImportJobOptions options;

ingestion/src/main/java/feast/ingestion/options/ImportJobOptions.java

Lines changed: 14 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@
2727
import org.apache.beam.sdk.options.Validation.Required;
2828

2929
public interface ImportJobOptions extends PipelineOptions {
30+
3031
@Description("Import spec yaml file path")
3132
@Required(groups = {"importSpec"})
3233
String getImportSpecYamlFile();
@@ -72,14 +73,23 @@ public interface ImportJobOptions extends PipelineOptions {
7273
void setLimit(Long value);
7374

7475
@Description(
75-
"Set a store id to store errors in, if your data input is **very** small, you can use STDOUT"
76-
+ " or STDERR as the store id, otherwise it must match an associated storage spec")
77-
String getErrorsStoreId();
76+
"Set an errors store type. One of: [stderr, stdout, file.json]. Note that you should not use "
77+
+ "stderr/stdout in production unless your data volume is extremely small.")
78+
String getErrorsStoreType();
79+
80+
void setErrorsStoreType(String value);
7881

79-
void setErrorsStoreId(String value);
82+
@Description(
83+
"Provide errors store options as a json string containing key-values. Options required"
84+
+ "depend on the type of store set.")
85+
@Default.String("{}")
86+
String getErrorsStoreOptions();
87+
88+
void setErrorsStoreOptions(String value);
8089

8190
@AutoService(PipelineOptionsRegistrar.class)
8291
class ImportJobOptionsRegistrar implements PipelineOptionsRegistrar {
92+
8393
@Override
8494
public Iterable<Class<? extends PipelineOptions>> getPipelineOptions() {
8595
return Collections.singleton(ImportJobOptions.class);

ingestion/src/main/java/feast/ingestion/transform/ErrorsStoreTransform.java

Lines changed: 41 additions & 40 deletions
Original file line numberDiff line numberDiff line change
@@ -17,14 +17,9 @@
1717

1818
package feast.ingestion.transform;
1919

20-
import static com.google.common.base.Preconditions.checkNotNull;
20+
import static feast.ingestion.util.JsonUtil.convertJsonStringToMap;
2121

2222
import com.google.inject.Inject;
23-
import java.util.List;
24-
import lombok.extern.slf4j.Slf4j;
25-
import org.apache.beam.sdk.transforms.ParDo;
26-
import org.apache.beam.sdk.values.PCollection;
27-
import org.apache.beam.sdk.values.PDone;
2823
import feast.ingestion.model.Specs;
2924
import feast.ingestion.options.ImportJobOptions;
3025
import feast.ingestion.transform.FeatureIO.Write;
@@ -33,54 +28,60 @@
3328
import feast.storage.ErrorsStore;
3429
import feast.storage.noop.NoOpIO;
3530
import feast.types.FeatureRowExtendedProto.FeatureRowExtended;
31+
import java.util.List;
32+
import lombok.extern.slf4j.Slf4j;
33+
import org.apache.beam.sdk.transforms.ParDo;
34+
import org.apache.beam.sdk.values.PCollection;
35+
import org.apache.beam.sdk.values.PDone;
3636
import org.slf4j.event.Level;
3737

3838
@Slf4j
3939
public class ErrorsStoreTransform extends FeatureIO.Write {
40-
public static final String STDERR_STORE_ID = "STDERR";
41-
public static final String STDOUT_STORE_ID = "STDOUT";
4240

43-
private String errorsStoreId;
44-
private List<ErrorsStore> stores;
41+
public static final String ERRORS_STORE_STDERR = "stderr";
42+
public static final String ERRORS_STORE_STDOUT = "stdout";
43+
public static final String ERRORS_STORE_JSON = "file.json";
44+
45+
private String errorsStoreType;
46+
private StorageSpec errorsStoreSpec;
47+
private ErrorsStore errorsStore;
4548
private Specs specs;
4649

4750
@Inject
48-
public ErrorsStoreTransform(ImportJobOptions options, List<ErrorsStore> stores, Specs specs) {
49-
this.errorsStoreId = options.getErrorsStoreId();
50-
this.stores = stores;
51+
public ErrorsStoreTransform(
52+
ImportJobOptions options, Specs specs, List<ErrorsStore> errorsStores) {
5153
this.specs = specs;
54+
this.errorsStoreType = options.getErrorsStoreType();
55+
56+
for (ErrorsStore errorsStore : errorsStores) {
57+
if (errorsStore.getType().equals(errorsStoreType)) {
58+
this.errorsStore = errorsStore;
59+
}
60+
}
61+
62+
this.errorsStoreSpec =
63+
StorageSpec.newBuilder()
64+
.setType(errorsStoreType)
65+
.putAllOptions(convertJsonStringToMap(options.getErrorsStoreOptions()))
66+
.build();
5267
}
5368

5469
@Override
5570
public PDone expand(PCollection<FeatureRowExtended> input) {
56-
if (errorsStoreId == null) {
57-
log.warn("No errorsStoreId specified, errors will be discarded");
58-
return input.apply(new NoOpIO.Write());
59-
}
60-
61-
if (errorsStoreId.equals(STDOUT_STORE_ID)) {
62-
input.apply("Log errors to STDOUT", ParDo.of(new LoggerDoFn(Level.INFO)));
63-
} else if (errorsStoreId.equals(STDERR_STORE_ID)) {
64-
input.apply("Log errors to STDERR", ParDo.of(new LoggerDoFn(Level.ERROR)));
65-
} else {
66-
StorageSpec storageSpec = specs.getStorageSpec(errorsStoreId);
67-
storageSpec =
68-
checkNotNull(
69-
storageSpec,
70-
String.format("errorsStoreId=%s not found in storage specs", errorsStoreId));
71-
Write write = null;
72-
for (ErrorsStore errorsStore : stores) {
73-
if (errorsStore.getType().equals(storageSpec.getType())) {
74-
write = errorsStore.create(storageSpec, specs);
71+
switch (errorsStoreType) {
72+
case ERRORS_STORE_STDOUT:
73+
input.apply("Log errors to STDOUT", ParDo.of(new LoggerDoFn(Level.INFO)));
74+
break;
75+
case ERRORS_STORE_STDERR:
76+
input.apply("Log errors to STDERR", ParDo.of(new LoggerDoFn(Level.ERROR)));
77+
break;
78+
default:
79+
if (errorsStore == null) {
80+
log.warn("No valid errors store specified, errors will be discarded");
81+
return input.apply(new NoOpIO.Write());
7582
}
76-
}
77-
write =
78-
checkNotNull(
79-
write,
80-
"No errors storage factory found for errorsStoreId=%s with type=%s",
81-
errorsStoreId,
82-
storageSpec.getType());
83-
return input.apply(write);
83+
Write write = errorsStore.create(this.errorsStoreSpec, specs);
84+
return input.apply(write);
8485
}
8586
return PDone.in(input.getPipeline());
8687
}
Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,27 @@
1+
package feast.ingestion.util;
2+
3+
import com.google.gson.Gson;
4+
import com.google.gson.reflect.TypeToken;
5+
import java.lang.reflect.Type;
6+
import java.util.Collections;
7+
import java.util.Map;
8+
9+
public class JsonUtil {
10+
11+
private static Gson gson = new Gson();
12+
13+
/**
14+
* Unmarshals a given json string to map
15+
*
16+
* @param jsonString valid json formatted string
17+
* @return map of keys to values in json
18+
*/
19+
public static Map<String, String> convertJsonStringToMap(String jsonString) {
20+
if (jsonString == null || jsonString.equals("") || jsonString.equals("{}")) {
21+
return Collections.emptyMap();
22+
}
23+
Type stringMapType = new TypeToken<Map<String, String>>() {
24+
}.getType();
25+
return gson.fromJson(jsonString, stringMapType);
26+
}
27+
}

ingestion/src/main/java/feast/storage/service/ErrorsStoreService.java

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -19,14 +19,16 @@
1919

2020
import com.google.common.collect.Iterators;
2121
import com.google.common.collect.Lists;
22+
import feast.storage.ErrorsStore;
23+
import lombok.extern.slf4j.Slf4j;
24+
2225
import java.util.ArrayList;
2326
import java.util.List;
2427
import java.util.ServiceLoader;
25-
import lombok.extern.slf4j.Slf4j;
26-
import feast.storage.ErrorsStore;
2728

2829
@Slf4j
2930
public class ErrorsStoreService {
31+
3032
private static ServiceLoader<ErrorsStore> serviceLoader = ServiceLoader.load(ErrorsStore.class);
3133
private static List<ErrorsStore> manuallyRegistered = new ArrayList<>();
3234

0 commit comments

Comments
 (0)