Skip to content

Commit 986381d

Browse files
zhilingcfeast-ci-bot
authored andcommitted
Allow submission of kafka jobs (#94)
1 parent c9e20d4 commit 986381d

2 files changed

Lines changed: 96 additions & 25 deletions

File tree

core/src/main/java/feast/core/validators/SpecValidator.java

Lines changed: 30 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,10 @@
1717

1818
package feast.core.validators;
1919

20+
import static com.google.common.base.Preconditions.checkArgument;
21+
import static com.google.common.base.Preconditions.checkNotNull;
22+
import static feast.core.validators.Matchers.checkLowerSnakeCase;
23+
2024
import com.google.common.base.Strings;
2125
import com.google.common.collect.Lists;
2226
import feast.core.dao.EntityInfoRepository;
@@ -35,36 +39,28 @@
3539
import feast.specs.ImportSpecProto.Field;
3640
import feast.specs.ImportSpecProto.ImportSpec;
3741
import feast.specs.StorageSpecProto.StorageSpec;
38-
import org.springframework.beans.factory.annotation.Autowired;
39-
4042
import java.util.Arrays;
4143
import java.util.Optional;
4244
import java.util.stream.Collectors;
4345
import java.util.stream.Stream;
44-
45-
import static com.google.common.base.Preconditions.checkArgument;
46-
import static com.google.common.base.Preconditions.checkNotNull;
47-
import static feast.core.validators.Matchers.checkLowerSnakeCase;
46+
import org.springframework.beans.factory.annotation.Autowired;
4847

4948
public class SpecValidator {
5049

51-
private StorageInfoRepository storageInfoRepository;
52-
private EntityInfoRepository entityInfoRepository;
53-
private FeatureGroupInfoRepository featureGroupInfoRepository;
54-
private FeatureInfoRepository featureInfoRepository;
5550
private static final String FILE_ERROR_STORE_TYPE = "file.json";
56-
5751
private static final String NO_STORE = "";
58-
5952
private static String[] SUPPORTED_WAREHOUSE_STORES =
60-
new String[] {
53+
new String[]{
6154
BigQueryStorageManager.TYPE, FILE_ERROR_STORE_TYPE,
6255
};
63-
6456
private static String[] SUPPORTED_SERVING_STORES =
65-
new String[] {
57+
new String[]{
6658
BigTableStorageManager.TYPE, PostgresStorageManager.TYPE, RedisStorageManager.TYPE,
6759
};
60+
private StorageInfoRepository storageInfoRepository;
61+
private EntityInfoRepository entityInfoRepository;
62+
private FeatureGroupInfoRepository featureGroupInfoRepository;
63+
private FeatureInfoRepository featureInfoRepository;
6864

6965
@Autowired
7066
public SpecValidator(
@@ -130,7 +126,8 @@ public void validateFeatureSpec(FeatureSpec spec) throws IllegalArgumentExceptio
130126
servingStoreId =
131127
servingStoreId.equals(NO_STORE) ? group.getServingStore().getId() : servingStoreId;
132128
warehouseStoreId =
133-
warehouseStoreId.equals(NO_STORE) ? group.getWarehouseStore().getId() : warehouseStoreId;
129+
warehouseStoreId.equals(NO_STORE) ? group.getWarehouseStore().getId()
130+
: warehouseStoreId;
134131
}
135132
Optional<StorageInfo> servingStore = storageInfoRepository.findById(servingStoreId);
136133
Optional<StorageInfo> warehouseStore = storageInfoRepository.findById(warehouseStoreId);
@@ -221,6 +218,9 @@ public void validateStorageSpec(StorageSpec spec) throws IllegalArgumentExceptio
221218
public void validateImportSpec(ImportSpec spec) throws IllegalArgumentException {
222219
try {
223220
switch (spec.getType()) {
221+
case "kafka":
222+
checkKafkaImportSpecOption(spec);
223+
break;
224224
case "pubsub":
225225
checkPubSubImportSpecOption(spec);
226226
break;
@@ -262,6 +262,20 @@ public void validateImportSpec(ImportSpec spec) throws IllegalArgumentException
262262
}
263263
}
264264

265+
private void checkKafkaImportSpecOption(ImportSpec spec) {
266+
try {
267+
String topics = spec.getOptionsOrDefault("topics", "");
268+
String server = spec.getOptionsOrDefault("server", "");
269+
if (topics.equals("") && server.equals("")) {
270+
throw new IllegalArgumentException(
271+
"Kafka ingestion requires either topics or servers");
272+
}
273+
} catch (NullPointerException | IllegalArgumentException e) {
274+
throw new IllegalArgumentException(
275+
Strings.lenientFormat("Invalid options: %s", e.getMessage()));
276+
}
277+
}
278+
265279
private void checkFileImportSpecOption(ImportSpec spec) throws IllegalArgumentException {
266280
try {
267281
checkArgument(!spec.getOptionsOrDefault("path", "").equals(""), "File path cannot be empty");

core/src/test/java/feast/core/validators/SpecValidatorTest.java

Lines changed: 66 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -44,13 +44,14 @@
4444
import org.mockito.Mockito;
4545

4646
public class SpecValidatorTest {
47+
48+
@Rule
49+
public final ExpectedException exception = ExpectedException.none();
4750
private FeatureInfoRepository featureInfoRepository;
4851
private FeatureGroupInfoRepository featureGroupInfoRepository;
4952
private EntityInfoRepository entityInfoRepository;
5053
private StorageInfoRepository storageInfoRepository;
5154

52-
@Rule public final ExpectedException exception = ExpectedException.none();
53-
5455
@Before
5556
public void setUp() {
5657
featureInfoRepository = Mockito.mock(FeatureInfoRepository.class);
@@ -369,7 +370,7 @@ public void featureSpecWithoutExistingWarehouseStoreShouldThrowIllegalArgumentEx
369370
StorageInfo redis1 = new StorageInfo();
370371
redis1.setId(servingStoreId);
371372
redis1.setType("redis");
372-
when(storageInfoRepository.findById( servingStoreId)).thenReturn(Optional.of(redis1));
373+
when(storageInfoRepository.findById(servingStoreId)).thenReturn(Optional.of(redis1));
373374

374375
SpecValidator validator =
375376
new SpecValidator(
@@ -392,7 +393,8 @@ public void featureSpecWithoutExistingWarehouseStoreShouldThrowIllegalArgumentEx
392393
.setDataStores(dataStores)
393394
.build();
394395
exception.expect(IllegalArgumentException.class);
395-
exception.expectMessage(String.format("Warehouse store with id %s does not exist", warehouseStoreId));
396+
exception.expectMessage(
397+
String.format("Warehouse store with id %s does not exist", warehouseStoreId));
396398
validator.validateFeatureSpec(input);
397399
}
398400

@@ -405,7 +407,7 @@ public void featureSpecWithoutWarehouseStoreShouldBeAllowed() {
405407
StorageInfo redis1 = new StorageInfo();
406408
redis1.setId(servingStoreId);
407409
redis1.setType("redis");
408-
when(storageInfoRepository.findById( servingStoreId)).thenReturn(Optional.of(redis1));
410+
when(storageInfoRepository.findById(servingStoreId)).thenReturn(Optional.of(redis1));
409411

410412
SpecValidator validator =
411413
new SpecValidator(
@@ -432,18 +434,21 @@ public void featureSpecWithoutWarehouseStoreShouldBeAllowed() {
432434
@Test
433435
public void featureSpecWithUnsupportedWarehouseStoreShouldThrowIllegalArgumentException() {
434436
String servingStoreId = "REDIS1";
435-
StorageSpec servingStoreSpec = StorageSpec.newBuilder().setId(servingStoreId).setType("redis").build();
437+
StorageSpec servingStoreSpec = StorageSpec.newBuilder().setId(servingStoreId).setType("redis")
438+
.build();
436439
StorageInfo servingStoreInfo = new StorageInfo(servingStoreSpec);
437440

438441
String warehouseStoreId = "REDIS2";
439-
StorageSpec warehouseStoreSpec = StorageSpec.newBuilder().setId(warehouseStoreId).setType("redis").build();
442+
StorageSpec warehouseStoreSpec = StorageSpec.newBuilder().setId(warehouseStoreId)
443+
.setType("redis").build();
440444
StorageInfo warehouseStoreInfo = new StorageInfo(warehouseStoreSpec);
441445

442446
when(entityInfoRepository.existsById("entity")).thenReturn(true);
443447
when(storageInfoRepository.existsById(servingStoreId)).thenReturn(true);
444448
when(storageInfoRepository.existsById(warehouseStoreId)).thenReturn(true);
445449
when(storageInfoRepository.findById(servingStoreId)).thenReturn(Optional.of(servingStoreInfo));
446-
when(storageInfoRepository.findById(warehouseStoreId)).thenReturn(Optional.of(warehouseStoreInfo));
450+
when(storageInfoRepository.findById(warehouseStoreId))
451+
.thenReturn(Optional.of(warehouseStoreInfo));
447452
SpecValidator validator =
448453
new SpecValidator(
449454
storageInfoRepository,
@@ -488,7 +493,8 @@ public void featureSpecWithUnsupportedServingStoreShouldThrowIllegalArgumentExce
488493
when(entityInfoRepository.existsById("entity")).thenReturn(true);
489494
when(storageInfoRepository.existsById(servingStoreName)).thenReturn(true);
490495
when(storageInfoRepository.existsById(warehouseStorageName)).thenReturn(true);
491-
when(storageInfoRepository.findById(servingStoreName)).thenReturn(Optional.of(redis1StorageInfo));
496+
when(storageInfoRepository.findById(servingStoreName))
497+
.thenReturn(Optional.of(redis1StorageInfo));
492498
when(storageInfoRepository.findById(warehouseStorageName)).thenReturn(Optional.of(bqInfo));
493499
SpecValidator validator =
494500
new SpecValidator(
@@ -778,4 +784,55 @@ public void importSpecWithUnregisteredFeaturesShouldThrowIllegalArgumentExceptio
778784
"Validation for import spec failed: Feature some_nonexistent_feature not registered");
779785
validator.validateImportSpec(input);
780786
}
787+
788+
@Test
789+
public void importSpecWithKafkaSourceAndCorrectOptionsShouldPassValidation() {
790+
SpecValidator validator =
791+
new SpecValidator(
792+
storageInfoRepository,
793+
entityInfoRepository,
794+
featureGroupInfoRepository,
795+
featureInfoRepository);
796+
when(featureInfoRepository.existsById("some_existing_feature")).thenReturn(true);
797+
when(entityInfoRepository.existsById("someEntity")).thenReturn(true);
798+
Schema schema =
799+
Schema.newBuilder()
800+
.addFields(Field.newBuilder().setFeatureId("some_existing_feature").build())
801+
.build();
802+
ImportSpec input =
803+
ImportSpec.newBuilder()
804+
.setType("kafka")
805+
.putOptions("topics", "my-kafka-topic")
806+
.putOptions("server", "localhost:54321")
807+
.setSchema(schema)
808+
.addEntities("someEntity")
809+
.build();
810+
validator.validateImportSpec(input);
811+
}
812+
813+
@Test
814+
public void importSpecWithKafkaSourceWithoutOptionsShouldThrowIllegalArgumentException() {
815+
SpecValidator validator =
816+
new SpecValidator(
817+
storageInfoRepository,
818+
entityInfoRepository,
819+
featureGroupInfoRepository,
820+
featureInfoRepository);
821+
when(featureInfoRepository.existsById("some_existing_feature")).thenReturn(true);
822+
when(entityInfoRepository.existsById("someEntity")).thenReturn(true);
823+
Schema schema =
824+
Schema.newBuilder()
825+
.addFields(Field.newBuilder().setFeatureId("some_existing_feature").build())
826+
.build();
827+
ImportSpec input =
828+
ImportSpec.newBuilder()
829+
.setType("kafka")
830+
.setSchema(schema)
831+
.addEntities("someEntity")
832+
.build();
833+
exception.expect(IllegalArgumentException.class);
834+
exception.expectMessage(
835+
"Validation for import spec failed: Invalid options: Kafka ingestion requires either topics or servers");
836+
validator.validateImportSpec(input);
837+
}
781838
}

0 commit comments

Comments
 (0)