Skip to content

Commit 4160821

Browse files
authored
Provide stable jobName in RowMetrics labels (#1028)
* stable job name in row metrics Signed-off-by: Oleksii Moskalenko <moskalenko.alexey@gmail.com> * fix test Signed-off-by: Oleksii Moskalenko <moskalenko.alexey@gmail.com> * fix test Signed-off-by: Oleksii Moskalenko <moskalenko.alexey@gmail.com>
1 parent e47903f commit 4160821

4 files changed

Lines changed: 16 additions & 9 deletions

File tree

ingestion/src/main/java/feast/ingestion/transform/metrics/WriteFeatureValueMetricsDoFn.java

Lines changed: 7 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -24,12 +24,9 @@
2424
import feast.proto.types.FeatureRowProto.FeatureRow;
2525
import feast.proto.types.FieldProto.Field;
2626
import feast.proto.types.ValueProto.Value;
27-
import java.util.ArrayList;
28-
import java.util.DoubleSummaryStatistics;
29-
import java.util.HashMap;
30-
import java.util.List;
31-
import java.util.Map;
27+
import java.util.*;
3228
import java.util.Map.Entry;
29+
import java.util.stream.Collectors;
3330
import org.apache.beam.sdk.transforms.DoFn;
3431
import org.apache.beam.sdk.values.KV;
3532
import org.apache.commons.math3.stat.descriptive.rank.Percentile;
@@ -149,6 +146,10 @@ public void processElement(
149146
}
150147
}
151148

149+
String[] split = context.getPipelineOptions().getJobName().split("-");
150+
String jobNameWithoutTimestamp =
151+
Arrays.stream(split).limit(split.length - 1).collect(Collectors.joining("-"));
152+
152153
for (Entry<String, DoubleSummaryStatistics> entry : featureNameToStats.entrySet()) {
153154
String featureName = entry.getKey();
154155
DoubleSummaryStatistics stats = entry.getValue();
@@ -157,7 +158,7 @@ public void processElement(
157158
FEATURE_SET_PROJECT_TAG_KEY + ":" + projectName,
158159
FEATURE_SET_NAME_TAG_KEY + ":" + featureSetName,
159160
FEATURE_TAG_KEY + ":" + featureName,
160-
INGESTION_JOB_NAME_KEY + ":" + context.getPipelineOptions().getJobName(),
161+
INGESTION_JOB_NAME_KEY + ":" + jobNameWithoutTimestamp,
161162
METRICS_NAMESPACE_KEY + ":" + getMetricsNamespace(),
162163
};
163164

ingestion/src/main/java/feast/ingestion/transform/metrics/WriteRowMetricsDoFn.java

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -25,9 +25,11 @@
2525
import feast.proto.types.ValueProto.Value;
2626
import feast.proto.types.ValueProto.Value.ValCase;
2727
import java.time.Clock;
28+
import java.util.Arrays;
2829
import java.util.HashMap;
2930
import java.util.Map;
3031
import java.util.Map.Entry;
32+
import java.util.stream.Collectors;
3133
import javax.annotation.Nullable;
3234
import org.apache.beam.sdk.transforms.DoFn;
3335
import org.apache.beam.sdk.values.KV;
@@ -192,11 +194,15 @@ public void processElement(
192194
}
193195
}
194196

197+
String[] split = c.getPipelineOptions().getJobName().split("-");
198+
String jobNameWithoutTimestamp =
199+
Arrays.stream(split).limit(split.length - 1).collect(Collectors.joining("-"));
200+
195201
String[] tags = {
196202
STORE_TAG_KEY + ":" + getStoreName(),
197203
FEATURE_SET_PROJECT_TAG_KEY + ":" + featureSetProject,
198204
FEATURE_SET_NAME_TAG_KEY + ":" + featureSetName,
199-
INGESTION_JOB_NAME_KEY + ":" + c.getPipelineOptions().getJobName(),
205+
INGESTION_JOB_NAME_KEY + ":" + jobNameWithoutTimestamp,
200206
METRICS_NAMESPACE_KEY + ":" + getMetricsNamespace(),
201207
};
202208

ingestion/src/test/java/feast/ingestion/transform/metrics/WriteFeatureValueMetricsDoFnTest.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -62,7 +62,7 @@ public class WriteFeatureValueMetricsDoFnTest {
6262
@Test
6363
public void shouldSendCorrectStatsDMetrics() throws IOException, InterruptedException {
6464
PipelineOptions pipelineOptions = PipelineOptionsFactory.create();
65-
pipelineOptions.setJobName("job");
65+
pipelineOptions.setJobName("job-12345678");
6666

6767
Map<String, Iterable<FeatureRow>> input =
6868
readTestInput("feast/ingestion/transform/WriteFeatureValueMetricsDoFnTest.input");

ingestion/src/test/java/feast/ingestion/transform/metrics/WriteRowMetricsDoFnTest.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -45,7 +45,7 @@ public class WriteRowMetricsDoFnTest {
4545
@Test
4646
public void shouldSendCorrectStatsDMetrics() throws IOException, InterruptedException {
4747
PipelineOptions pipelineOptions = PipelineOptionsFactory.create();
48-
pipelineOptions.setJobName("job");
48+
pipelineOptions.setJobName("job-12345678");
4949
Map<String, Iterable<FeatureRow>> input =
5050
readTestInput("feast/ingestion/transform/WriteRowMetricsDoFnTest.input");
5151
List<String> expectedLines =

0 commit comments

Comments
 (0)