22
33import com .google .protobuf .InvalidProtocolBufferException ;
44import feast .core .FeatureSetProto .FeatureSetSpec ;
5+ import feast .core .SourceProto .Source ;
56import feast .core .StoreProto .Store ;
67import feast .ingestion .options .ImportOptions ;
78import feast .ingestion .transform .ReadFromSource ;
9+ import feast .ingestion .transform .ValidateFeatureRows ;
810import feast .ingestion .transform .WriteFailedElementToBigQuery ;
911import feast .ingestion .transform .WriteToStore ;
1012import feast .ingestion .transform .metrics .WriteMetricsTransform ;
1416import feast .ingestion .values .FailedElement ;
1517import feast .types .FeatureRowProto .FeatureRow ;
1618import java .util .List ;
19+ import java .util .Map ;
20+ import java .util .stream .Collectors ;
1721import org .apache .beam .sdk .Pipeline ;
1822import org .apache .beam .sdk .PipelineResult ;
1923import org .apache .beam .sdk .options .PipelineOptionsFactory ;
2024import org .apache .beam .sdk .options .PipelineOptionsValidator ;
2125import org .apache .beam .sdk .values .PCollectionTuple ;
2226import org .apache .beam .sdk .values .TupleTag ;
27+ import org .apache .commons .lang3 .tuple .Pair ;
2328import org .slf4j .Logger ;
2429
2530public class ImportJob {
31+
2632 // Tag for main output containing Feature Row that has been successfully processed.
27- private static final TupleTag <FeatureRow > FEATURE_ROW_OUT = new TupleTag <FeatureRow >() {};
33+ private static final TupleTag <FeatureRow > FEATURE_ROW_OUT = new TupleTag <FeatureRow >() {
34+ };
2835
2936 // Tag for deadletter output containing elements and error messages from invalid input/transform.
30- private static final TupleTag <FailedElement > DEADLETTER_OUT = new TupleTag <FailedElement >() {};
37+ private static final TupleTag <FailedElement > DEADLETTER_OUT = new TupleTag <FailedElement >() {
38+ };
3139 private static final Logger log = org .slf4j .LoggerFactory .getLogger (ImportJob .class );
3240
3341 /**
@@ -46,14 +54,17 @@ public static PipelineResult runPipeline(ImportOptions options)
4654 /*
4755 * Steps:
4856 * 1. Read messages from Feast Source as FeatureRow
49- * 2. Write FeatureRow to the corresponding Store
50- * 3. Write elements that failed to be processed to a dead letter queue.
51- * 4. Write metrics to a metrics sink
57+ * 2. Validate the feature rows to ensure the schema matches what is registered to the system
58+ * 3. Write FeatureRow to the corresponding Store
59+ * 4. Write elements that failed to be processed to a dead letter queue.
60+ * 5. Write metrics to a metrics sink
5261 */
5362
5463 PipelineOptionsValidator .validate (ImportOptions .class , options );
5564 Pipeline pipeline = Pipeline .create (options );
5665
66+ log .info ("Starting import job with settings: \n {}" , options .toString ());
67+
5768 List <FeatureSetSpec > featureSetSpecs =
5869 SpecUtil .parseFeatureSetSpecJsonList (options .getFeatureSetSpecJson ());
5970 List <Store > stores = SpecUtil .parseStoreJsonList (options .getStoreJson ());
@@ -62,44 +73,71 @@ public static PipelineResult runPipeline(ImportOptions options)
6273 List <FeatureSetSpec > subscribedFeatureSets =
6374 SpecUtil .getSubscribedFeatureSets (store .getSubscriptionsList (), featureSetSpecs );
6475
76+ // Generate tags by key
77+ Map <String , TupleTag <FeatureRow >> featureSetTagsByKey = subscribedFeatureSets .stream ()
78+ .map (fs -> {
79+ String id = String .format ("%s:%s" , fs .getName (), fs .getVersion ());
80+ return Pair .of (id , new TupleTag <FeatureRow >(id ) {
81+ });
82+ })
83+ .collect (Collectors .toMap (Pair ::getLeft , Pair ::getRight ));
84+
85+ // TODO: make the source part of the job initialisation options
86+ Source source = subscribedFeatureSets .get (0 ).getSource ();
87+
88+ // Step 1. Read messages from Feast Source as FeatureRow.
89+ PCollectionTuple convertedFeatureRows =
90+ pipeline .apply (
91+ "ReadFeatureRowFromSource" ,
92+ ReadFromSource .newBuilder ()
93+ .setSource (source )
94+ .setFeatureSetTagByKey (featureSetTagsByKey )
95+ .setFailureTag (DEADLETTER_OUT )
96+ .build ());
97+
6598 for (FeatureSetSpec featureSet : subscribedFeatureSets ) {
6699 // Ensure Store has valid configuration and Feast can access it.
67100 StoreUtil .setupStore (store , featureSet );
101+ String id = String .format ("%s:%s" , featureSet .getName (), featureSet .getVersion ());
102+
103+ // Step 2. Validate incoming FeatureRows
104+ PCollectionTuple validatedRows = convertedFeatureRows
105+ .get (featureSetTagsByKey .get (id ))
106+ .apply (ValidateFeatureRows .newBuilder ()
107+ .setFeatureSetSpec (featureSet )
108+ .setSuccessTag (FEATURE_ROW_OUT )
109+ .setFailureTag (DEADLETTER_OUT )
110+ .build ());
68111
69- // Step 1. Read messages from Feast Source as FeatureRow.
70- PCollectionTuple convertedFeatureRows =
71- pipeline .apply (
72- "ReadFeatureRowFromSource" ,
73- ReadFromSource .newBuilder ()
74- .setSource (featureSet .getSource ())
75- .setFieldByName (SpecUtil .getFieldByName (featureSet ))
76- .setFeatureSetName (featureSet .getName ())
77- .setFeatureSetVersion (featureSet .getVersion ())
78- .setSuccessTag (FEATURE_ROW_OUT )
79- .setFailureTag (DEADLETTER_OUT )
80- .build ());
81-
82- // Step 2. Write FeatureRow to the corresponding Store.
83- convertedFeatureRows
112+ // Step 3. Write FeatureRow to the corresponding Store.
113+ validatedRows
84114 .get (FEATURE_ROW_OUT )
85115 .apply (
86116 "WriteFeatureRowToStore" ,
87117 WriteToStore .newBuilder ().setFeatureSetSpec (featureSet ).setStore (store ).build ());
88118
89- // Step 3 . Write FailedElements to a dead letter table in BigQuery.
119+ // Step 4 . Write FailedElements to a dead letter table in BigQuery.
90120 if (options .getDeadLetterTableSpec () != null ) {
91121 convertedFeatureRows
92122 .get (DEADLETTER_OUT )
93123 .apply (
94- "WriteFailedElements" ,
124+ "WriteFailedElements_ReadFromSource" ,
125+ WriteFailedElementToBigQuery .newBuilder ()
126+ .setJsonSchema (ResourceUtil .getDeadletterTableSchemaJson ())
127+ .setTableSpec (options .getDeadLetterTableSpec ())
128+ .build ());
129+
130+ validatedRows
131+ .get (DEADLETTER_OUT )
132+ .apply ("WriteFailedElements_ValidateRows" ,
95133 WriteFailedElementToBigQuery .newBuilder ()
96134 .setJsonSchema (ResourceUtil .getDeadletterTableSchemaJson ())
97135 .setTableSpec (options .getDeadLetterTableSpec ())
98136 .build ());
99137 }
100138
101- // Step 4 . Write metrics to a metrics sink.
102- convertedFeatureRows
139+ // Step 5 . Write metrics to a metrics sink.
140+ validatedRows
103141 .apply ("WriteMetrics" , WriteMetricsTransform .newBuilder ()
104142 .setFeatureSetSpec (featureSet )
105143 .setStoreName (store .getName ())
0 commit comments