Skip to content

Commit a35fe05

Browse files
author
zhilingc
committed
Re-implement metrics for ingestion, remove more redundant code
1 parent 225582f commit a35fe05

18 files changed

Lines changed: 205 additions & 2019 deletions

ingestion/src/main/java/feast/ingestion/ImportJob.java

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@
77
import feast.ingestion.transform.ReadFromSource;
88
import feast.ingestion.transform.WriteFailedElementToBigQuery;
99
import feast.ingestion.transform.WriteToStore;
10+
import feast.ingestion.transform.metrics.WriteMetricsTransform;
1011
import feast.ingestion.utils.ResourceUtil;
1112
import feast.ingestion.utils.SpecUtil;
1213
import feast.ingestion.utils.StoreUtil;
@@ -47,6 +48,7 @@ public static PipelineResult runPipeline(ImportOptions options)
4748
* 1. Read messages from Feast Source as FeatureRow
4849
* 2. Write FeatureRow to the corresponding Store
4950
* 3. Write elements that failed to be processed to a dead letter queue.
51+
* 4. Write metrics to a metrics sink
5052
*/
5153

5254
PipelineOptionsValidator.validate(ImportOptions.class, options);
@@ -95,6 +97,15 @@ public static PipelineResult runPipeline(ImportOptions options)
9597
.setTableSpec(options.getDeadLetterTableSpec())
9698
.build());
9799
}
100+
101+
// Step 4. Write metrics to a metrics sink.
102+
convertedFeatureRows
103+
.apply("WriteMetrics", WriteMetricsTransform.newBuilder()
104+
.setFeatureSetSpec(featureSet)
105+
.setStoreName(store.getName())
106+
.setSuccessTag(FEATURE_ROW_OUT)
107+
.setFailureTag(DEADLETTER_OUT)
108+
.build());
98109
}
99110
}
100111

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

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

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

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

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

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

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

Lines changed: 0 additions & 85 deletions
This file was deleted.
Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
1+
package feast.ingestion.transform.metrics;
2+
3+
import org.apache.beam.sdk.transforms.DoFn;
4+
import org.apache.beam.sdk.transforms.GroupByKey;
5+
import org.apache.beam.sdk.transforms.PTransform;
6+
import org.apache.beam.sdk.transforms.ParDo;
7+
import org.apache.beam.sdk.transforms.windowing.FixedWindows;
8+
import org.apache.beam.sdk.transforms.windowing.Window;
9+
import org.apache.beam.sdk.values.KV;
10+
import org.apache.beam.sdk.values.PCollection;
11+
import org.joda.time.Duration;
12+
13+
public class WindowRecords<T> extends
14+
PTransform<PCollection<T>, PCollection<KV<Integer, Iterable<T>>>> {
15+
16+
private final long windowSize;
17+
18+
public WindowRecords(long windowSize) {
19+
this.windowSize = windowSize;
20+
}
21+
22+
@Override
23+
public PCollection<KV<Integer, Iterable<T>>> expand(PCollection<T> input) {
24+
return input
25+
.apply("Window records",
26+
Window.into(FixedWindows.of(Duration.standardSeconds(windowSize))))
27+
.apply("Add key", ParDo.of(new DoFn<T, KV<Integer, T>>() {
28+
@ProcessElement
29+
public void processElement(ProcessContext c) {
30+
c.output(KV.of(1, c.element()));
31+
}
32+
}))
33+
.apply("Collect", GroupByKey.create());
34+
}
35+
}

0 commit comments

Comments
 (0)