4040import feast .core .model .Source ;
4141import feast .core .model .Store ;
4242import java .util .Arrays ;
43+ import java .util .Collections ;
44+ import java .util .List ;
4345import java .util .Optional ;
4446import org .hamcrest .core .IsNull ;
4547import org .junit .Before ;
4850
4951public class JobUpdateTaskTest {
5052
53+ private static final FeatureSetProto .FeatureSet .Builder fsBuilder =
54+ FeatureSetProto .FeatureSet .newBuilder ().setMeta (FeatureSetMeta .newBuilder ());
55+ private static final FeatureSetSpec .Builder specBuilder =
56+ FeatureSetSpec .newBuilder ().setProject ("project1" ).setVersion (1 );
57+
5158 @ Mock private JobManager jobManager ;
5259
5360 private StoreProto .Store store ;
5461 private SourceProto .Source source ;
62+ private FeatureSetProto .FeatureSet featureSet1 ;
5563
5664 @ Before
5765 public void setUp () {
@@ -76,38 +84,27 @@ public void setUp() {
7684 .setBootstrapServers ("servers:9092" )
7785 .build ())
7886 .build ();
87+
88+ featureSet1 = fsBuilder .setSpec (specBuilder .setName ("featureSet1" ).setSource (source )).build ();
7989 }
8090
8191 @ Test
8292 public void shouldUpdateJobIfPresent () {
83- FeatureSetProto .FeatureSet featureSet1 =
84- FeatureSetProto .FeatureSet .newBuilder ()
85- .setSpec (
86- FeatureSetSpec .newBuilder ()
87- .setSource (source )
88- .setProject ("project1" )
89- .setName ("featureSet1" )
90- .setVersion (1 ))
91- .setMeta (FeatureSetMeta .newBuilder ())
92- .build ();
9393 FeatureSetProto .FeatureSet featureSet2 =
94- FeatureSetProto .FeatureSet .newBuilder ()
95- .setSpec (
96- FeatureSetSpec .newBuilder ()
97- .setSource (source )
98- .setProject ("project1" )
99- .setName ("featureSet2" )
100- .setVersion (1 ))
101- .setMeta (FeatureSetMeta .newBuilder ())
102- .build ();
94+ fsBuilder .setSpec (specBuilder .setName ("featureSet2" )).build ();
95+ List <FeatureSet > existingFeatureSetsPopulatedByJob =
96+ Collections .singletonList (FeatureSet .fromProto (featureSet1 ));
97+ List <FeatureSet > newFeatureSetsPopulatedByJob =
98+ Arrays .asList (FeatureSet .fromProto (featureSet1 ), FeatureSet .fromProto (featureSet2 ));
99+
103100 Job originalJob =
104101 new Job (
105102 "job" ,
106103 "old_ext" ,
107104 Runner .DATAFLOW .name (),
108- feast . core . model . Source .fromProto (source ),
109- feast . core . model . Store .fromProto (store ),
110- Arrays . asList ( FeatureSet . fromProto ( featureSet1 )) ,
105+ Source .fromProto (source ),
106+ Store .fromProto (store ),
107+ existingFeatureSetsPopulatedByJob ,
111108 JobStatus .RUNNING );
112109 JobUpdateTask jobUpdateTask =
113110 new JobUpdateTask (
@@ -122,9 +119,9 @@ public void shouldUpdateJobIfPresent() {
122119 "job" ,
123120 "old_ext" ,
124121 Runner .DATAFLOW .name (),
125- feast . core . model . Source .fromProto (source ),
126- feast . core . model . Store .fromProto (store ),
127- Arrays . asList ( FeatureSet . fromProto ( featureSet1 ), FeatureSet . fromProto ( featureSet2 )) ,
122+ Source .fromProto (source ),
123+ Store .fromProto (store ),
124+ newFeatureSetsPopulatedByJob ,
128125 JobStatus .RUNNING );
129126
130127 Job expected =
@@ -134,7 +131,7 @@ public void shouldUpdateJobIfPresent() {
134131 Runner .DATAFLOW .name (),
135132 Source .fromProto (source ),
136133 Store .fromProto (store ),
137- Arrays . asList ( FeatureSet . fromProto ( featureSet1 ), FeatureSet . fromProto ( featureSet2 )) ,
134+ newFeatureSetsPopulatedByJob ,
138135 JobStatus .PENDING );
139136 when (jobManager .updateJob (submittedJob )).thenReturn (expected );
140137 Job actual = jobUpdateTask .call ();
@@ -144,16 +141,6 @@ public void shouldUpdateJobIfPresent() {
144141
145142 @ Test
146143 public void shouldCreateJobIfNotPresent () {
147- FeatureSetProto .FeatureSet featureSet1 =
148- FeatureSetProto .FeatureSet .newBuilder ()
149- .setSpec (
150- FeatureSetSpec .newBuilder ()
151- .setSource (source )
152- .setProject ("project1" )
153- .setName ("featureSet1" )
154- .setVersion (1 ))
155- .setMeta (FeatureSetMeta .newBuilder ())
156- .build ();
157144 JobUpdateTask jobUpdateTask =
158145 spy (
159146 new JobUpdateTask (
@@ -165,8 +152,8 @@ public void shouldCreateJobIfNotPresent() {
165152 "job" ,
166153 "" ,
167154 Runner .DATAFLOW .name (),
168- feast . core . model . Source .fromProto (source ),
169- feast . core . model . Store .fromProto (store ),
155+ Source .fromProto (source ),
156+ Store .fromProto (store ),
170157 Arrays .asList (FeatureSet .fromProto (featureSet1 )),
171158 JobStatus .PENDING );
172159
@@ -175,8 +162,8 @@ public void shouldCreateJobIfNotPresent() {
175162 "job" ,
176163 "ext" ,
177164 Runner .DATAFLOW .name (),
178- feast . core . model . Source .fromProto (source ),
179- feast . core . model . Store .fromProto (store ),
165+ Source .fromProto (source ),
166+ Store .fromProto (store ),
180167 Arrays .asList (FeatureSet .fromProto (featureSet1 )),
181168 JobStatus .RUNNING );
182169
@@ -188,56 +175,27 @@ public void shouldCreateJobIfNotPresent() {
188175
189176 @ Test
190177 public void shouldUpdateJobStatusIfNotCreateOrUpdate () {
191- FeatureSetProto .FeatureSet featureSet1 =
192- FeatureSetProto .FeatureSet .newBuilder ()
193- .setSpec (
194- FeatureSetSpec .newBuilder ()
195- .setSource (source )
196- .setProject ("project1" )
197- .setName ("featureSet1" )
198- .setVersion (1 ))
199- .setMeta (FeatureSetMeta .newBuilder ())
200- .build ();
201178 Job originalJob =
202179 new Job (
203180 "job" ,
204181 "ext" ,
205182 Runner .DATAFLOW .name (),
206- feast . core . model . Source .fromProto (source ),
207- feast . core . model . Store .fromProto (store ),
183+ Source .fromProto (source ),
184+ Store .fromProto (store ),
208185 Arrays .asList (FeatureSet .fromProto (featureSet1 )),
209186 JobStatus .RUNNING );
210187 JobUpdateTask jobUpdateTask =
211188 new JobUpdateTask (
212189 Arrays .asList (featureSet1 ), source , store , Optional .of (originalJob ), jobManager , 100L );
213190
214191 when (jobManager .getJobStatus (originalJob )).thenReturn (JobStatus .ABORTING );
215- Job expected =
216- new Job (
217- "job" ,
218- "ext" ,
219- Runner .DATAFLOW .name (),
220- Source .fromProto (source ),
221- Store .fromProto (store ),
222- Arrays .asList (FeatureSet .fromProto (featureSet1 )),
223- JobStatus .ABORTING );
224- Job actual = jobUpdateTask .call ();
192+ Job updated = jobUpdateTask .call ();
225193
226- assertThat (actual , equalTo (expected ));
194+ assertThat (updated . getStatus () , equalTo (JobStatus . ABORTING ));
227195 }
228196
229197 @ Test
230198 public void shouldReturnJobWithErrorStatusIfFailedToSubmit () {
231- FeatureSetProto .FeatureSet featureSet1 =
232- FeatureSetProto .FeatureSet .newBuilder ()
233- .setSpec (
234- FeatureSetSpec .newBuilder ()
235- .setSource (source )
236- .setProject ("project1" )
237- .setName ("featureSet1" )
238- .setVersion (1 ))
239- .setMeta (FeatureSetMeta .newBuilder ())
240- .build ();
241199 JobUpdateTask jobUpdateTask =
242200 spy (
243201 new JobUpdateTask (
@@ -249,8 +207,8 @@ public void shouldReturnJobWithErrorStatusIfFailedToSubmit() {
249207 "job" ,
250208 "" ,
251209 Runner .DATAFLOW .name (),
252- feast . core . model . Source .fromProto (source ),
253- feast . core . model . Store .fromProto (store ),
210+ Source .fromProto (source ),
211+ Store .fromProto (store ),
254212 Arrays .asList (FeatureSet .fromProto (featureSet1 )),
255213 JobStatus .PENDING );
256214
@@ -259,8 +217,8 @@ public void shouldReturnJobWithErrorStatusIfFailedToSubmit() {
259217 "job" ,
260218 "" ,
261219 Runner .DATAFLOW .name (),
262- feast . core . model . Source .fromProto (source ),
263- feast . core . model . Store .fromProto (store ),
220+ Source .fromProto (source ),
221+ Store .fromProto (store ),
264222 Arrays .asList (FeatureSet .fromProto (featureSet1 )),
265223 JobStatus .ERROR );
266224
@@ -273,17 +231,6 @@ public void shouldReturnJobWithErrorStatusIfFailedToSubmit() {
273231
274232 @ Test
275233 public void shouldTimeout () {
276- FeatureSetProto .FeatureSet featureSet1 =
277- FeatureSetProto .FeatureSet .newBuilder ()
278- .setSpec (
279- FeatureSetSpec .newBuilder ()
280- .setSource (source )
281- .setProject ("project1" )
282- .setName ("featureSet1" )
283- .setVersion (1 ))
284- .setMeta (FeatureSetMeta .newBuilder ())
285- .build ();
286-
287234 JobUpdateTask jobUpdateTask =
288235 spy (
289236 new JobUpdateTask (
0 commit comments