From a0102b241d8dc48c06eff0a6846c16ddeab9a1cc Mon Sep 17 00:00:00 2001 From: Oleksii Moskalenko Date: Tue, 1 Dec 2020 17:58:51 +0800 Subject: [PATCH 1/2] fix max age Signed-off-by: Oleksii Moskalenko --- .../java/feast/core/model/FeatureTable.java | 2 +- .../feast/core/service/SpecServiceIT.java | 212 +++++++----------- 2 files changed, 85 insertions(+), 129 deletions(-) diff --git a/core/src/main/java/feast/core/model/FeatureTable.java b/core/src/main/java/feast/core/model/FeatureTable.java index a1a89f0ae4f..b442d57b583 100644 --- a/core/src/main/java/feast/core/model/FeatureTable.java +++ b/core/src/main/java/feast/core/model/FeatureTable.java @@ -411,7 +411,7 @@ && getProject().equals(other.getProject()) && getLabelsJSON().equals(other.getLabelsJSON()) && getFeatures().equals(other.getFeatures()) && getEntities().equals(other.getEntities()) - && getMaxAgeSecs() == getMaxAgeSecs() + && getMaxAgeSecs() == other.getMaxAgeSecs() && Optional.ofNullable(getBatchSource()).equals(Optional.ofNullable(other.getBatchSource())) && Optional.ofNullable(getStreamSource()) .equals(Optional.ofNullable(other.getStreamSource())); diff --git a/core/src/test/java/feast/core/service/SpecServiceIT.java b/core/src/test/java/feast/core/service/SpecServiceIT.java index 61a3195a55f..b0a0bfc9304 100644 --- a/core/src/test/java/feast/core/service/SpecServiceIT.java +++ b/core/src/test/java/feast/core/service/SpecServiceIT.java @@ -27,6 +27,7 @@ import static org.junit.jupiter.api.Assertions.assertThrows; import avro.shaded.com.google.common.collect.ImmutableMap; +import com.google.protobuf.Duration; import feast.common.it.BaseIT; import feast.common.it.DataGenerator; import feast.common.it.SimpleCoreClient; @@ -38,6 +39,9 @@ import io.grpc.ManagedChannelBuilder; import io.grpc.StatusRuntimeException; import java.util.*; +import java.util.stream.IntStream; + +import org.apache.commons.lang3.ArrayUtils; import org.apache.commons.lang3.tuple.Triple; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.BeforeEach; @@ -61,6 +65,9 @@ public static void globalSetUp(@Value("${grpc.server.port}") int port) { apiClient = new SimpleCoreClient(stub); } + private FeatureTableProto.FeatureTableSpec example1; + private FeatureTableProto.FeatureTableSpec example2; + @BeforeEach public void initState() { @@ -78,40 +85,41 @@ public void initState() { ImmutableMap.of("label_key2", "label_value2")); apiClient.simpleApplyEntity("default", entitySpec1); apiClient.simpleApplyEntity("default", entitySpec2); - apiClient.applyFeatureTable( - "default", - DataGenerator.createFeatureTableSpec( - "featuretable1", - Arrays.asList("entity1", "entity2"), - new HashMap<>() { - { - put("feature1", ValueProto.ValueType.Enum.STRING); - put("feature2", ValueProto.ValueType.Enum.FLOAT); - } - }, - 7200, - ImmutableMap.of("feat_key2", "feat_value2")) - .toBuilder() - .setBatchSource( - DataGenerator.createFileDataSourceSpec("file:///path/to/file", "ts_col", "")) - .build()); - apiClient.applyFeatureTable( - "default", - DataGenerator.createFeatureTableSpec( - "featuretable2", - Arrays.asList("entity1", "entity2"), - new HashMap<>() { - { - put("feature3", ValueProto.ValueType.Enum.STRING); - put("feature4", ValueProto.ValueType.Enum.FLOAT); - } - }, - 7200, - ImmutableMap.of("feat_key4", "feat_value4")) - .toBuilder() - .setBatchSource( - DataGenerator.createFileDataSourceSpec("file:///path/to/file", "ts_col", "")) - .build()); + + example1 = DataGenerator.createFeatureTableSpec( + "featuretable1", + Arrays.asList("entity1", "entity2"), + new HashMap<>() { + { + put("feature1", ValueProto.ValueType.Enum.STRING); + put("feature2", ValueProto.ValueType.Enum.FLOAT); + } + }, + 7200, + ImmutableMap.of("feat_key2", "feat_value2")) + .toBuilder() + .setBatchSource( + DataGenerator.createFileDataSourceSpec("file:///path/to/file", "ts_col", "")) + .build(); + + example2 = DataGenerator.createFeatureTableSpec( + "featuretable2", + Arrays.asList("entity1", "entity2"), + new HashMap<>() { + { + put("feature3", ValueProto.ValueType.Enum.STRING); + put("feature4", ValueProto.ValueType.Enum.FLOAT); + } + }, + 7200, + ImmutableMap.of("feat_key4", "feat_value4")) + .toBuilder() + .setBatchSource( + DataGenerator.createFileDataSourceSpec("file:///path/to/file", "ts_col", "")) + .build(); + + apiClient.applyFeatureTable("default", example1); + apiClient.applyFeatureTable("default", example2); apiClient.simpleApplyEntity( "project1", DataGenerator.createEntitySpecV2( @@ -422,26 +430,10 @@ public void shouldThrowExceptionGivenNoSuchFeatureTable() { @Test public void shouldReturnFeatureTableIfExists() { - FeatureTableSpec featureTableSpec = - DataGenerator.createFeatureTableSpec( - "featuretable1", - Arrays.asList("entity1", "entity2"), - new HashMap<>() { - { - put("feature1", ValueProto.ValueType.Enum.STRING); - put("feature2", ValueProto.ValueType.Enum.FLOAT); - } - }, - 7200, - ImmutableMap.of("feat_key2", "feat_value2")) - .toBuilder() - .setBatchSource( - DataGenerator.createFileDataSourceSpec("file:///path/to/file", "ts_col", "")) - .build(); FeatureTableProto.FeatureTable featureTable = apiClient.simpleGetFeatureTable("default", "featuretable1"); - assertTrue(TestUtil.compareFeatureTableSpec(featureTable.getSpec(), featureTableSpec)); + assertTrue(TestUtil.compareFeatureTableSpec(featureTable.getSpec(), example1)); } } @@ -544,17 +536,8 @@ public void shouldFilterFeaturesByEntitiesAndLabels() { @Nested public class ApplyFeatureTable { private FeatureTableSpec getTestSpec() { - return DataGenerator.createFeatureTableSpec( - "ft", - List.of("entity1", "entity2"), - Map.of( - "feature1", ValueProto.ValueType.Enum.INT64, - "feature2", ValueProto.ValueType.Enum.FLOAT), - 3600, - Map.of()) - .toBuilder() - .setBatchSource( - DataGenerator.createFileDataSourceSpec("file:///path/to/file", "ts_col", "")) + return example1.toBuilder() + .setName("apply_test") .setStreamSource( DataGenerator.createKafkaDataSourceSpec( "localhost:9092", "topic", "class.path", "ts_col")) @@ -573,23 +556,17 @@ public void shouldApplyNewValidTable() { public void shouldUpdateExistingTableWithValidSpec() { FeatureTableProto.FeatureTable table = apiClient.applyFeatureTable("default", getTestSpec()); - FeatureTableSpec updatedSpec = - DataGenerator.createFeatureTableSpec( - "ft", - List.of("entity1", "entity2"), - Map.of( - "feature2", ValueProto.ValueType.Enum.FLOAT, - "feature3", ValueProto.ValueType.Enum.INT64, - "feature4", ValueProto.ValueType.Enum.INT64), - 2100, - Map.of("test", "labels")) - .toBuilder() - .setStreamSource( - DataGenerator.createFileDataSourceSpec("file:///path/to/file", "ts_col", "")) - .setBatchSource( - DataGenerator.createKafkaDataSourceSpec( - "localhost:9092", "topic", "class.path", "ts_col")) - .build(); + FeatureTableSpec updatedSpec = getTestSpec().toBuilder() + .clearFeatures() + .addFeatures( + DataGenerator.createFeatureSpecV2("feature5", ValueProto.ValueType.Enum.FLOAT, ImmutableMap.of()) + ) + .setStreamSource( + DataGenerator.createKafkaDataSourceSpec( + "localhost:9092", "new_topic", "new.class", "ts_col") + ) + .build(); + FeatureTableProto.FeatureTable updatedTable = apiClient.applyFeatureTable("default", updatedSpec); @@ -599,22 +576,9 @@ public void shouldUpdateExistingTableWithValidSpec() { @Test public void shouldUpdateFeatureTableOnEntityChange() { - List entities = Arrays.asList("entity1", "entity2"); - FeatureTableProto.FeatureTableSpec updatedSpec = - DataGenerator.createFeatureTableSpec( - "featuretable1", - Arrays.asList("entity1"), - new HashMap<>() { - { - put("feature1", ValueProto.ValueType.Enum.STRING); - put("feature2", ValueProto.ValueType.Enum.FLOAT); - } - }, - 7200, - ImmutableMap.of("feat_key2", "feat_value2")) - .toBuilder() - .setBatchSource( - DataGenerator.createFileDataSourceSpec("file:///path/to/file", "ts_col", "")) + FeatureTableProto.FeatureTableSpec updatedSpec = getTestSpec().toBuilder() + .clearEntities() + .addEntities("entity1") .build(); FeatureTableProto.FeatureTable updatedTable = @@ -623,24 +587,28 @@ public void shouldUpdateFeatureTableOnEntityChange() { assertTrue(TestUtil.compareFeatureTableSpec(updatedTable.getSpec(), updatedSpec)); } + @Test + public void shouldUpdateFeatureTableOnMaxAgeChange() { + FeatureTableProto.FeatureTableSpec updatedSpec = getTestSpec().toBuilder() + .setMaxAge(Duration.newBuilder().setSeconds(600).build()) + .build(); + + FeatureTableProto.FeatureTable updatedTable = + apiClient.applyFeatureTable("default", updatedSpec); + + assertTrue(TestUtil.compareFeatureTableSpec(updatedTable.getSpec(), updatedSpec)); + } + @Test public void shouldUpdateFeatureTableOnFeatureTypeChange() { - FeatureTableProto.FeatureTableSpec updatedSpec = - DataGenerator.createFeatureTableSpec( - "featuretable1", - Arrays.asList("entity1", "entity2"), - new HashMap<>() { - { - put("feature1", ValueProto.ValueType.Enum.STRING); - put("feature2", ValueProto.ValueType.Enum.STRING_LIST); - } - }, - 7200, - ImmutableMap.of("feat_key2", "feat_value2")) - .toBuilder() - .setBatchSource( - DataGenerator.createFileDataSourceSpec("file:///path/to/file", "ts_col", "")) - .build(); + int featureIdx = IntStream.range(0, getTestSpec().getFeaturesCount()) + .filter(i -> getTestSpec().getFeatures(i).getName().equals("feature2")) + .findFirst().orElse(-1); + + FeatureTableProto.FeatureTableSpec updatedSpec = getTestSpec().toBuilder() + .setFeatures(featureIdx, + DataGenerator.createFeatureSpecV2("feature2", ValueProto.ValueType.Enum.STRING_LIST, ImmutableMap.of())) + .build(); FeatureTableProto.FeatureTable updatedTable = apiClient.applyFeatureTable("default", updatedSpec); @@ -650,23 +618,11 @@ public void shouldUpdateFeatureTableOnFeatureTypeChange() { @Test public void shouldUpdateFeatureTableOnFeatureAddition() { - FeatureTableProto.FeatureTableSpec updatedSpec = - DataGenerator.createFeatureTableSpec( - "featuretable1", - Arrays.asList("entity1", "entity2"), - new HashMap<>() { - { - put("feature1", ValueProto.ValueType.Enum.STRING); - put("feature2", ValueProto.ValueType.Enum.FLOAT); - put("feature3", ValueProto.ValueType.Enum.FLOAT); - } - }, - 7200, - ImmutableMap.of("feat_key2", "feat_value2")) - .toBuilder() - .setBatchSource( - DataGenerator.createFileDataSourceSpec("file:///path/to/file", "ts_col", "")) - .build(); + FeatureTableProto.FeatureTableSpec updatedSpec = getTestSpec().toBuilder() + .addFeatures( + DataGenerator.createFeatureSpecV2("feature6", ValueProto.ValueType.Enum.FLOAT, ImmutableMap.of()) + ) + .build(); FeatureTableProto.FeatureTable updatedTable = apiClient.applyFeatureTable("default", updatedSpec); From 1a4d38afb7174e80ead4767ee5e3bdda902fdfe7 Mon Sep 17 00:00:00 2001 From: Oleksii Moskalenko Date: Tue, 1 Dec 2020 17:59:28 +0800 Subject: [PATCH 2/2] format Signed-off-by: Oleksii Moskalenko --- .../feast/core/service/SpecServiceIT.java | 138 ++++++++++-------- 1 file changed, 74 insertions(+), 64 deletions(-) diff --git a/core/src/test/java/feast/core/service/SpecServiceIT.java b/core/src/test/java/feast/core/service/SpecServiceIT.java index b0a0bfc9304..e6dc068ef69 100644 --- a/core/src/test/java/feast/core/service/SpecServiceIT.java +++ b/core/src/test/java/feast/core/service/SpecServiceIT.java @@ -40,8 +40,6 @@ import io.grpc.StatusRuntimeException; import java.util.*; import java.util.stream.IntStream; - -import org.apache.commons.lang3.ArrayUtils; import org.apache.commons.lang3.tuple.Triple; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.BeforeEach; @@ -86,37 +84,39 @@ public void initState() { apiClient.simpleApplyEntity("default", entitySpec1); apiClient.simpleApplyEntity("default", entitySpec2); - example1 = DataGenerator.createFeatureTableSpec( - "featuretable1", - Arrays.asList("entity1", "entity2"), - new HashMap<>() { - { - put("feature1", ValueProto.ValueType.Enum.STRING); - put("feature2", ValueProto.ValueType.Enum.FLOAT); - } - }, - 7200, - ImmutableMap.of("feat_key2", "feat_value2")) - .toBuilder() - .setBatchSource( - DataGenerator.createFileDataSourceSpec("file:///path/to/file", "ts_col", "")) - .build(); - - example2 = DataGenerator.createFeatureTableSpec( - "featuretable2", - Arrays.asList("entity1", "entity2"), - new HashMap<>() { - { - put("feature3", ValueProto.ValueType.Enum.STRING); - put("feature4", ValueProto.ValueType.Enum.FLOAT); - } - }, - 7200, - ImmutableMap.of("feat_key4", "feat_value4")) - .toBuilder() - .setBatchSource( - DataGenerator.createFileDataSourceSpec("file:///path/to/file", "ts_col", "")) - .build(); + example1 = + DataGenerator.createFeatureTableSpec( + "featuretable1", + Arrays.asList("entity1", "entity2"), + new HashMap<>() { + { + put("feature1", ValueProto.ValueType.Enum.STRING); + put("feature2", ValueProto.ValueType.Enum.FLOAT); + } + }, + 7200, + ImmutableMap.of("feat_key2", "feat_value2")) + .toBuilder() + .setBatchSource( + DataGenerator.createFileDataSourceSpec("file:///path/to/file", "ts_col", "")) + .build(); + + example2 = + DataGenerator.createFeatureTableSpec( + "featuretable2", + Arrays.asList("entity1", "entity2"), + new HashMap<>() { + { + put("feature3", ValueProto.ValueType.Enum.STRING); + put("feature4", ValueProto.ValueType.Enum.FLOAT); + } + }, + 7200, + ImmutableMap.of("feat_key4", "feat_value4")) + .toBuilder() + .setBatchSource( + DataGenerator.createFileDataSourceSpec("file:///path/to/file", "ts_col", "")) + .build(); apiClient.applyFeatureTable("default", example1); apiClient.applyFeatureTable("default", example2); @@ -536,7 +536,8 @@ public void shouldFilterFeaturesByEntitiesAndLabels() { @Nested public class ApplyFeatureTable { private FeatureTableSpec getTestSpec() { - return example1.toBuilder() + return example1 + .toBuilder() .setName("apply_test") .setStreamSource( DataGenerator.createKafkaDataSourceSpec( @@ -556,16 +557,17 @@ public void shouldApplyNewValidTable() { public void shouldUpdateExistingTableWithValidSpec() { FeatureTableProto.FeatureTable table = apiClient.applyFeatureTable("default", getTestSpec()); - FeatureTableSpec updatedSpec = getTestSpec().toBuilder() - .clearFeatures() - .addFeatures( - DataGenerator.createFeatureSpecV2("feature5", ValueProto.ValueType.Enum.FLOAT, ImmutableMap.of()) - ) - .setStreamSource( - DataGenerator.createKafkaDataSourceSpec( - "localhost:9092", "new_topic", "new.class", "ts_col") - ) - .build(); + FeatureTableSpec updatedSpec = + getTestSpec() + .toBuilder() + .clearFeatures() + .addFeatures( + DataGenerator.createFeatureSpecV2( + "feature5", ValueProto.ValueType.Enum.FLOAT, ImmutableMap.of())) + .setStreamSource( + DataGenerator.createKafkaDataSourceSpec( + "localhost:9092", "new_topic", "new.class", "ts_col")) + .build(); FeatureTableProto.FeatureTable updatedTable = apiClient.applyFeatureTable("default", updatedSpec); @@ -576,10 +578,8 @@ public void shouldUpdateExistingTableWithValidSpec() { @Test public void shouldUpdateFeatureTableOnEntityChange() { - FeatureTableProto.FeatureTableSpec updatedSpec = getTestSpec().toBuilder() - .clearEntities() - .addEntities("entity1") - .build(); + FeatureTableProto.FeatureTableSpec updatedSpec = + getTestSpec().toBuilder().clearEntities().addEntities("entity1").build(); FeatureTableProto.FeatureTable updatedTable = apiClient.applyFeatureTable("default", updatedSpec); @@ -589,9 +589,11 @@ public void shouldUpdateFeatureTableOnEntityChange() { @Test public void shouldUpdateFeatureTableOnMaxAgeChange() { - FeatureTableProto.FeatureTableSpec updatedSpec = getTestSpec().toBuilder() - .setMaxAge(Duration.newBuilder().setSeconds(600).build()) - .build(); + FeatureTableProto.FeatureTableSpec updatedSpec = + getTestSpec() + .toBuilder() + .setMaxAge(Duration.newBuilder().setSeconds(600).build()) + .build(); FeatureTableProto.FeatureTable updatedTable = apiClient.applyFeatureTable("default", updatedSpec); @@ -601,14 +603,20 @@ public void shouldUpdateFeatureTableOnMaxAgeChange() { @Test public void shouldUpdateFeatureTableOnFeatureTypeChange() { - int featureIdx = IntStream.range(0, getTestSpec().getFeaturesCount()) - .filter(i -> getTestSpec().getFeatures(i).getName().equals("feature2")) - .findFirst().orElse(-1); - - FeatureTableProto.FeatureTableSpec updatedSpec = getTestSpec().toBuilder() - .setFeatures(featureIdx, - DataGenerator.createFeatureSpecV2("feature2", ValueProto.ValueType.Enum.STRING_LIST, ImmutableMap.of())) - .build(); + int featureIdx = + IntStream.range(0, getTestSpec().getFeaturesCount()) + .filter(i -> getTestSpec().getFeatures(i).getName().equals("feature2")) + .findFirst() + .orElse(-1); + + FeatureTableProto.FeatureTableSpec updatedSpec = + getTestSpec() + .toBuilder() + .setFeatures( + featureIdx, + DataGenerator.createFeatureSpecV2( + "feature2", ValueProto.ValueType.Enum.STRING_LIST, ImmutableMap.of())) + .build(); FeatureTableProto.FeatureTable updatedTable = apiClient.applyFeatureTable("default", updatedSpec); @@ -618,11 +626,13 @@ public void shouldUpdateFeatureTableOnFeatureTypeChange() { @Test public void shouldUpdateFeatureTableOnFeatureAddition() { - FeatureTableProto.FeatureTableSpec updatedSpec = getTestSpec().toBuilder() - .addFeatures( - DataGenerator.createFeatureSpecV2("feature6", ValueProto.ValueType.Enum.FLOAT, ImmutableMap.of()) - ) - .build(); + FeatureTableProto.FeatureTableSpec updatedSpec = + getTestSpec() + .toBuilder() + .addFeatures( + DataGenerator.createFeatureSpecV2( + "feature6", ValueProto.ValueType.Enum.FLOAT, ImmutableMap.of())) + .build(); FeatureTableProto.FeatureTable updatedTable = apiClient.applyFeatureTable("default", updatedSpec);