Skip to content

Commit df908fa

Browse files
committed
rename FeatureIO.Write to FeatureStoreWrite
1 parent 7c217e2 commit df908fa

22 files changed

Lines changed: 47 additions & 96 deletions

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

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -23,8 +23,8 @@
2323
import com.google.inject.Inject;
2424
import feast.ingestion.model.Specs;
2525
import feast.ingestion.options.ImportJobPipelineOptions;
26-
import feast.ingestion.transform.FeatureIO.Write;
2726
import feast.specs.StorageSpecProto.StorageSpec;
27+
import feast.store.FeatureStoreWrite;
2828
import feast.store.NoOpIO;
2929
import feast.store.errors.FeatureErrorsFactory;
3030
import feast.types.FeatureRowExtendedProto.FeatureRowExtended;
@@ -35,7 +35,7 @@
3535
import org.apache.hadoop.hbase.util.Strings;
3636

3737
@Slf4j
38-
public class ErrorsStoreTransform extends FeatureIO.Write {
38+
public class ErrorsStoreTransform extends FeatureStoreWrite {
3939

4040
private String errorsStoreType;
4141
private StorageSpec errorsStoreSpec;
@@ -61,7 +61,7 @@ public ErrorsStoreTransform(
6161

6262
@Override
6363
public PDone expand(PCollection<FeatureRowExtended> input) {
64-
Write write;
64+
FeatureStoreWrite write;
6565
if (Strings.isEmpty(errorsStoreType)) {
6666
write = new NoOpIO.Write();
6767
} else {

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

Lines changed: 0 additions & 44 deletions
This file was deleted.

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

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -20,11 +20,11 @@
2020
import com.google.common.base.Preconditions;
2121
import com.google.common.collect.Lists;
2222
import feast.ingestion.model.Specs;
23-
import feast.ingestion.transform.FeatureIO.Write;
2423
import feast.ingestion.transform.SplitFeatures.MultiOutputSplit;
2524
import feast.ingestion.values.PFeatureRows;
2625
import feast.specs.FeatureSpecProto.FeatureSpec;
2726
import feast.specs.StorageSpecProto.StorageSpec;
27+
import feast.store.FeatureStoreWrite;
2828
import feast.store.FeatureStoreFactory;
2929
import feast.types.FeatureRowExtendedProto.FeatureRowExtended;
3030
import java.util.Collection;
@@ -52,14 +52,14 @@ public class SplitOutputByStore extends PTransform<PFeatureRows, PFeatureRows> {
5252

5353
@Override
5454
public PFeatureRows expand(PFeatureRows input) {
55-
Map<String, Write> transforms = getFeatureStoreTransforms();
55+
Map<String, FeatureStoreWrite> transforms = getFeatureStoreTransforms();
5656
Set<String> keys = transforms.keySet();
5757

5858
log.info(String.format("Splitting on keys = [%s]", String.join(",", keys)));
5959
MultiOutputSplit<String> splitter = new MultiOutputSplit<>(selector, keys, specs);
6060
PCollectionTuple splits = input.getMain().apply(splitter);
6161

62-
Map<TupleTag<FeatureRowExtended>, Write> taggedTransforms = new HashMap<>();
62+
Map<TupleTag<FeatureRowExtended>, FeatureStoreWrite> taggedTransforms = new HashMap<>();
6363
for (String key : transforms.keySet()) {
6464
TupleTag<FeatureRowExtended> tag = splitter.getStrategy().getTag(key);
6565
taggedTransforms.put(tag, transforms.get(key));
@@ -80,9 +80,9 @@ private Map<String, FeatureStoreFactory> getStoresMap() {
8080
return storesMap;
8181
}
8282

83-
private Map<String, Write> getFeatureStoreTransforms() {
83+
private Map<String, FeatureStoreWrite> getFeatureStoreTransforms() {
8484
Map<String, FeatureStoreFactory> storesMap = getStoresMap();
85-
Map<String, Write> transforms = new HashMap<>();
85+
Map<String, FeatureStoreWrite> transforms = new HashMap<>();
8686
Map<String, StorageSpec> storageSpecs = specs.getStorageSpecs();
8787
for (String storeId : storageSpecs.keySet()) {
8888
StorageSpec storageSpec = storageSpecs.get(storeId);
@@ -109,14 +109,14 @@ private Map<String, Write> getFeatureStoreTransforms() {
109109
public static class WriteTags extends
110110
PTransform<PCollectionTuple, PCollection<FeatureRowExtended>> {
111111

112-
private Map<TupleTag<FeatureRowExtended>, Write> transforms;
112+
private Map<TupleTag<FeatureRowExtended>, FeatureStoreWrite> transforms;
113113
private TupleTag<FeatureRowExtended> mainTag;
114114

115115
@Override
116116
public PCollection<FeatureRowExtended> expand(PCollectionTuple tuple) {
117117
List<PCollection<FeatureRowExtended>> outputList = Lists.newArrayList();
118118
for (TupleTag<FeatureRowExtended> tag : transforms.keySet()) {
119-
Write write = transforms.get(tag);
119+
FeatureStoreWrite write = transforms.get(tag);
120120
Preconditions.checkNotNull(write, String.format("Null transform for tag=%s", tag.getId()));
121121
PCollection<FeatureRowExtended> input = tuple.get(tag);
122122
input.apply(String.format("Write to %s", tag.getId()), write);

ingestion/src/main/java/feast/store/FeatureStoreFactory.java

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -18,11 +18,10 @@
1818
package feast.store;
1919

2020
import feast.ingestion.model.Specs;
21-
import feast.ingestion.transform.FeatureIO;
2221
import feast.specs.StorageSpecProto.StorageSpec;
2322

2423
public interface FeatureStoreFactory {
25-
FeatureIO.Write create(StorageSpec storageSpec, Specs specs);
24+
FeatureStoreWrite create(StorageSpec storageSpec, Specs specs);
2625

2726
String getType();
2827
}

ingestion/src/main/java/feast/store/FeatureStoreWrite.java

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -17,11 +17,9 @@
1717

1818
package feast.store;
1919

20-
import feast.types.FeatureRowProto.FeatureRow;
20+
import feast.types.FeatureRowExtendedProto.FeatureRowExtended;
2121
import org.apache.beam.sdk.transforms.PTransform;
2222
import org.apache.beam.sdk.values.PCollection;
2323
import org.apache.beam.sdk.values.PDone;
2424

25-
public abstract class FeatureStoreWrite extends PTransform<PCollection<FeatureRow>, PDone> {
26-
27-
}
25+
public abstract class FeatureStoreWrite extends PTransform<PCollection<FeatureRowExtended>, PDone> {}

ingestion/src/main/java/feast/store/NoOpIO.java

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,6 @@
1717

1818
package feast.store;
1919

20-
import feast.ingestion.transform.FeatureIO;
2120
import feast.ingestion.transform.fn.Identity;
2221
import feast.types.FeatureRowExtendedProto.FeatureRowExtended;
2322
import org.apache.beam.sdk.transforms.ParDo;
@@ -26,7 +25,7 @@
2625

2726
public class NoOpIO {
2827

29-
public static class Write extends FeatureIO.Write {
28+
public static class Write extends FeatureStoreWrite {
3029

3130
@Override
3231
public PDone expand(PCollection<FeatureRowExtended> input) {

ingestion/src/main/java/feast/store/errors/json/JsonFileErrorsFactory.java

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

2020
import com.google.auto.service.AutoService;
2121
import feast.ingestion.model.Specs;
22-
import feast.ingestion.transform.FeatureIO.Write;
22+
import feast.store.FeatureStoreWrite;
2323
import feast.options.OptionsParser;
2424
import feast.specs.StorageSpecProto.StorageSpec;
2525
import feast.store.FileStoreOptions;
@@ -34,7 +34,7 @@ public class JsonFileErrorsFactory implements FeatureErrorsFactory {
3434
private static final String JSON_FILES_TYPE = "file.json";
3535

3636
@Override
37-
public Write create(StorageSpec storageSpec, Specs specs) {
37+
public FeatureStoreWrite create(StorageSpec storageSpec, Specs specs) {
3838
FileStoreOptions options =
3939
OptionsParser.parse(storageSpec.getOptionsMap(), FileStoreOptions.class);
4040
options.jobName = specs.getJobName();

ingestion/src/main/java/feast/store/errors/json/JsonFileErrorsWrite.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,7 @@
2222

2323
import com.google.protobuf.InvalidProtocolBufferException;
2424
import com.google.protobuf.util.JsonFormat;
25-
import feast.ingestion.transform.FeatureIO;
25+
import feast.store.FeatureStoreWrite;
2626
import feast.store.FileStoreOptions;
2727
import feast.store.TextFileDynamicIO;
2828
import feast.types.FeatureRowExtendedProto.FeatureRowExtended;
@@ -33,7 +33,7 @@
3333
import org.apache.beam.sdk.values.PDone;
3434

3535
@AllArgsConstructor
36-
public class JsonFileErrorsWrite extends FeatureIO.Write {
36+
public class JsonFileErrorsWrite extends FeatureStoreWrite {
3737

3838
private FileStoreOptions options;
3939

ingestion/src/main/java/feast/store/errors/logging/LogIO.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,7 @@
1717

1818
package feast.store.errors.logging;
1919

20-
import feast.ingestion.transform.FeatureIO;
20+
import feast.store.FeatureStoreWrite;
2121
import feast.ingestion.transform.fn.LoggerDoFn;
2222
import feast.types.FeatureRowExtendedProto.FeatureRowExtended;
2323
import lombok.AllArgsConstructor;
@@ -29,7 +29,7 @@
2929
public class LogIO {
3030

3131
@AllArgsConstructor
32-
public static class Write extends FeatureIO.Write {
32+
public static class Write extends FeatureStoreWrite {
3333

3434
private Level level;
3535

ingestion/src/main/java/feast/store/errors/logging/StderrFeatureErrorsFactory.java

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

2020
import com.google.auto.service.AutoService;
2121
import feast.ingestion.model.Specs;
22-
import feast.ingestion.transform.FeatureIO.Write;
22+
import feast.store.FeatureStoreWrite;
2323
import feast.specs.StorageSpecProto.StorageSpec;
2424
import feast.store.errors.FeatureErrorsFactory;
2525
import org.slf4j.event.Level;
@@ -30,7 +30,7 @@ public class StderrFeatureErrorsFactory implements FeatureErrorsFactory {
3030
public static final String TYPE_STDERR = "stderr";
3131

3232
@Override
33-
public Write create(StorageSpec storageSpec, Specs specs) {
33+
public FeatureStoreWrite create(StorageSpec storageSpec, Specs specs) {
3434
return new LogIO.Write(Level.ERROR);
3535
}
3636

0 commit comments

Comments
 (0)