diff --git a/core/src/main/java/feast/core/job/dataflow/DataflowJobManager.java b/core/src/main/java/feast/core/job/dataflow/DataflowJobManager.java index e76568dfb48..f4df3d352a9 100644 --- a/core/src/main/java/feast/core/job/dataflow/DataflowJobManager.java +++ b/core/src/main/java/feast/core/job/dataflow/DataflowJobManager.java @@ -187,16 +187,7 @@ private Job submitDataflowJob( ImportOptions pipelineOptions = getPipelineOptions(jobName, featureSetProtos, sink, update); DataflowPipelineJob pipelineResult = runPipeline(pipelineOptions); List featureSets = - featureSetProtos.stream() - .map( - fsp -> { - FeatureSet featureSet = new FeatureSet(); - featureSet.setName(fsp.getSpec().toString()); - featureSet.setVersion(fsp.getSpec().getVersion()); - featureSet.setProject(new Project(fsp.getSpec().getProject())); - return featureSet; - }) - .collect(Collectors.toList()); + featureSetProtos.stream().map(FeatureSet::fromProto).collect(Collectors.toList()); String jobId = waitForJobToRun(pipelineResult); return new Job( jobName,