From ada1078f1d3948f1da99bf1ef326ef0fdf6f4599 Mon Sep 17 00:00:00 2001 From: Oleksii Moskalenko Date: Wed, 7 Oct 2020 14:54:08 +0800 Subject: [PATCH] Provide stable jobName in RowMetrics labels (#1028) * stable job name in row metrics Signed-off-by: Oleksii Moskalenko * fix test Signed-off-by: Oleksii Moskalenko * fix test Signed-off-by: Oleksii Moskalenko --- .../metrics/WriteFeatureValueMetricsDoFn.java | 13 +++++++------ .../transform/metrics/WriteRowMetricsDoFn.java | 8 +++++++- .../metrics/WriteFeatureValueMetricsDoFnTest.java | 2 +- .../transform/metrics/WriteRowMetricsDoFnTest.java | 2 +- 4 files changed, 16 insertions(+), 9 deletions(-) diff --git a/ingestion/src/main/java/feast/ingestion/transform/metrics/WriteFeatureValueMetricsDoFn.java b/ingestion/src/main/java/feast/ingestion/transform/metrics/WriteFeatureValueMetricsDoFn.java index 45fa4d25cf7..1e40e44d5cd 100644 --- a/ingestion/src/main/java/feast/ingestion/transform/metrics/WriteFeatureValueMetricsDoFn.java +++ b/ingestion/src/main/java/feast/ingestion/transform/metrics/WriteFeatureValueMetricsDoFn.java @@ -24,12 +24,9 @@ import feast.proto.types.FeatureRowProto.FeatureRow; import feast.proto.types.FieldProto.Field; import feast.proto.types.ValueProto.Value; -import java.util.ArrayList; -import java.util.DoubleSummaryStatistics; -import java.util.HashMap; -import java.util.List; -import java.util.Map; +import java.util.*; import java.util.Map.Entry; +import java.util.stream.Collectors; import org.apache.beam.sdk.transforms.DoFn; import org.apache.beam.sdk.values.KV; import org.apache.commons.math3.stat.descriptive.rank.Percentile; @@ -149,6 +146,10 @@ public void processElement( } } + String[] split = context.getPipelineOptions().getJobName().split("-"); + String jobNameWithoutTimestamp = + Arrays.stream(split).limit(split.length - 1).collect(Collectors.joining("-")); + for (Entry entry : featureNameToStats.entrySet()) { String featureName = entry.getKey(); DoubleSummaryStatistics stats = entry.getValue(); @@ -157,7 +158,7 @@ public void processElement( FEATURE_SET_PROJECT_TAG_KEY + ":" + projectName, FEATURE_SET_NAME_TAG_KEY + ":" + featureSetName, FEATURE_TAG_KEY + ":" + featureName, - INGESTION_JOB_NAME_KEY + ":" + context.getPipelineOptions().getJobName(), + INGESTION_JOB_NAME_KEY + ":" + jobNameWithoutTimestamp, METRICS_NAMESPACE_KEY + ":" + getMetricsNamespace(), }; diff --git a/ingestion/src/main/java/feast/ingestion/transform/metrics/WriteRowMetricsDoFn.java b/ingestion/src/main/java/feast/ingestion/transform/metrics/WriteRowMetricsDoFn.java index 8650285445a..29597ae84f4 100644 --- a/ingestion/src/main/java/feast/ingestion/transform/metrics/WriteRowMetricsDoFn.java +++ b/ingestion/src/main/java/feast/ingestion/transform/metrics/WriteRowMetricsDoFn.java @@ -25,9 +25,11 @@ import feast.proto.types.ValueProto.Value; import feast.proto.types.ValueProto.Value.ValCase; import java.time.Clock; +import java.util.Arrays; import java.util.HashMap; import java.util.Map; import java.util.Map.Entry; +import java.util.stream.Collectors; import javax.annotation.Nullable; import org.apache.beam.sdk.transforms.DoFn; import org.apache.beam.sdk.values.KV; @@ -192,11 +194,15 @@ public void processElement( } } + String[] split = c.getPipelineOptions().getJobName().split("-"); + String jobNameWithoutTimestamp = + Arrays.stream(split).limit(split.length - 1).collect(Collectors.joining("-")); + String[] tags = { STORE_TAG_KEY + ":" + getStoreName(), FEATURE_SET_PROJECT_TAG_KEY + ":" + featureSetProject, FEATURE_SET_NAME_TAG_KEY + ":" + featureSetName, - INGESTION_JOB_NAME_KEY + ":" + c.getPipelineOptions().getJobName(), + INGESTION_JOB_NAME_KEY + ":" + jobNameWithoutTimestamp, METRICS_NAMESPACE_KEY + ":" + getMetricsNamespace(), }; diff --git a/ingestion/src/test/java/feast/ingestion/transform/metrics/WriteFeatureValueMetricsDoFnTest.java b/ingestion/src/test/java/feast/ingestion/transform/metrics/WriteFeatureValueMetricsDoFnTest.java index 677f5937902..54b68eb2203 100644 --- a/ingestion/src/test/java/feast/ingestion/transform/metrics/WriteFeatureValueMetricsDoFnTest.java +++ b/ingestion/src/test/java/feast/ingestion/transform/metrics/WriteFeatureValueMetricsDoFnTest.java @@ -62,7 +62,7 @@ public class WriteFeatureValueMetricsDoFnTest { @Test public void shouldSendCorrectStatsDMetrics() throws IOException, InterruptedException { PipelineOptions pipelineOptions = PipelineOptionsFactory.create(); - pipelineOptions.setJobName("job"); + pipelineOptions.setJobName("job-12345678"); Map> input = readTestInput("feast/ingestion/transform/WriteFeatureValueMetricsDoFnTest.input"); diff --git a/ingestion/src/test/java/feast/ingestion/transform/metrics/WriteRowMetricsDoFnTest.java b/ingestion/src/test/java/feast/ingestion/transform/metrics/WriteRowMetricsDoFnTest.java index bc2b8f2dd5e..3bed89e9b13 100644 --- a/ingestion/src/test/java/feast/ingestion/transform/metrics/WriteRowMetricsDoFnTest.java +++ b/ingestion/src/test/java/feast/ingestion/transform/metrics/WriteRowMetricsDoFnTest.java @@ -45,7 +45,7 @@ public class WriteRowMetricsDoFnTest { @Test public void shouldSendCorrectStatsDMetrics() throws IOException, InterruptedException { PipelineOptions pipelineOptions = PipelineOptionsFactory.create(); - pipelineOptions.setJobName("job"); + pipelineOptions.setJobName("job-12345678"); Map> input = readTestInput("feast/ingestion/transform/WriteRowMetricsDoFnTest.input"); List expectedLines =