[
https://issues.apache.org/jira/browse/FLINK-2423?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=14663270#comment-14663270
]
ASF GitHub Bot commented on FLINK-2423:
---------------------------------------
Github user gyfora commented on a diff in the pull request:
https://github.com/apache/flink/pull/980#discussion_r36584713
--- Diff:
flink-tests/src/test/java/org/apache/flink/test/checkpointing/StreamCheckpointNotifierITCase.java
---
@@ -0,0 +1,340 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.flink.test.checkpointing;
+
+import org.apache.flink.api.common.functions.RichFilterFunction;
+import org.apache.flink.api.common.functions.RichMapFunction;
+import org.apache.flink.api.common.functions.RichReduceFunction;
+import org.apache.flink.api.common.state.OperatorState;
+import org.apache.flink.api.java.tuple.Tuple1;
+import org.apache.flink.configuration.Configuration;
+import org.apache.flink.streaming.api.checkpoint.CheckpointNotifier;
+import org.apache.flink.streaming.api.checkpoint.Checkpointed;
+import org.apache.flink.streaming.api.datastream.DataStream;
+import
org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
+import org.apache.flink.streaming.api.functions.co.RichCoFlatMapFunction;
+import org.apache.flink.streaming.api.functions.sink.SinkFunction;
+import
org.apache.flink.streaming.api.functions.source.ParallelSourceFunction;
+import org.apache.flink.streaming.api.functions.source.RichSourceFunction;
+import org.apache.flink.streaming.api.operators.OneInputStreamOperator;
+import org.apache.flink.streaming.api.operators.TwoInputStreamOperator;
+import org.apache.flink.util.Collector;
+
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Random;
+
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertTrue;
+
+/**
+ * Integration test for the {@link CheckpointNotifier} interface. The test
ensures that
+ * {@link CheckpointNotifier#notifyCheckpointComplete(long)} is called for
some completed
+ * checkpoints, that it is called at most once for any checkpoint id and
that it is not
+ * called for a deliberately failed checkpoint.
+ *
+ * <p>
+ * The topology tested here includes a number of {@link
OneInputStreamOperator}s and a
+ * {@link TwoInputStreamOperator}.
+ *
+ * <p>
+ * Note that as a result of doing the checks on the task level there is no
way to verify
+ * that the {@link CheckpointNotifier#notifyCheckpointComplete(long)} is
called for every
+ * successfully completed checkpoint.
+ */
+@SuppressWarnings("serial")
+public class StreamCheckpointNotifierITCase extends
StreamFaultToleranceTestBase {
+
+ final long NUM_LONGS = 10_000_000L;
+
+ /**
+ * Runs the following program:
+ *
+ * <pre>
+ * [ (source)->(filter) ] -> [ (co-map) ] -> [ (map) ] -> [
(groupBy/reduce)->(sink) ]
+ * </pre>
+ */
+ @Override
+ public void testProgram(StreamExecutionEnvironment env) {
+
+ DataStream<Long> stream = env.addSource(new
GeneratingSourceFunction(NUM_LONGS));
+
+ stream
+ // -------------- first vertex, chained to the
src ----------------
+ .filter(new LongRichFilterFunction())
+
+ // -------------- second vertex, applying the
co-map ----------------
+ .connect(stream).flatMap(new
LeftIdentityCoRichFlatMapFunction())
+
+ // -------------- third vertex - the stateful
one that also fails ----------------
+ .map(new IdentityMapFunction())
+ .startNewChain()
+
+ // -------------- fourth vertex - reducer and
the sink ----------------
+ .groupBy(0)
+ .reduce(new OnceFailingReducer(NUM_LONGS))
+ .addSink(new SinkFunction<Tuple1<Long>>() {
+ @Override
+ public void invoke(Tuple1<Long> value) {
+ // do nothing
+ }
+ });
+ }
+
+ @Override
+ public void postSubmit() {
+ List[][] checkList = new List[][]{
GeneratingSourceFunction.completedCheckpoints,
+ IdentityMapFunction.completedCheckpoints,
+ LongRichFilterFunction.completedCheckpoints,
+
LeftIdentityCoRichFlatMapFunction.completedCheckpoints};
+
+ for(List[] parallelNotifications : checkList) {
+ for (int i = 0; i < PARALLELISM; i++){
+ List<Long> notifications =
parallelNotifications[i];
+ assertTrue("No checkpoint notification was
received.",
+ notifications.size() > 0);
+ assertFalse("Failure checkpoint was marked as
completed.",
+
notifications.contains(OnceFailingReducer.failureCheckpointID));
+ assertTrue("Checkpoint notification was
received multiple times",
+ notifications.size() == new
HashSet<Long>(notifications).size());
--- End diff --
Maybe it would be also good to check whether we received anything after the
failed checkpoint.
> Properly test checkpoint notifications
> --------------------------------------
>
> Key: FLINK-2423
> URL: https://issues.apache.org/jira/browse/FLINK-2423
> Project: Flink
> Issue Type: Improvement
> Components: Streaming
> Reporter: Gyula Fora
> Assignee: Márton Balassi
>
> Checkpoint notifications (via the CheckpointNotifier interface) are currently
> not properly tested.
> A test should be included to verify that checkpoint notifications are
> eventually called on successful checkpoints, and that they are only called
> once per checkpointID.
--
This message was sent by Atlassian JIRA
(v6.3.4#6332)