flink-commits mailing list archives

Site index · List index
Message view « Date » · « Thread »
Top « Date » · « Thread »
From aljos...@apache.org
Subject [2/2] flink git commit: [FLINK-7429] [kinesis] Unify empty state restore behaviour across 1.1 / 1.2 / 1.3
Date Fri, 25 Aug 2017 13:18:09 GMT
[FLINK-7429] [kinesis] Unify empty state restore behaviour across 1.1 / 1.2 / 1.3

Prior to this commit, when restoring empty state from previous Flink
versions, the behaviour was different for each version. For older
versions, restoring empty state results in `null`. For newer versions,
restoring empty state results in an empty map.

We want that an empty map represents "this is a restored run, but there
was no state for us", and a null to represent" this is not a restored
run".


Project: http://git-wip-us.apache.org/repos/asf/flink/repo
Commit: http://git-wip-us.apache.org/repos/asf/flink/commit/6876abd3
Tree: http://git-wip-us.apache.org/repos/asf/flink/tree/6876abd3
Diff: http://git-wip-us.apache.org/repos/asf/flink/diff/6876abd3

Branch: refs/heads/release-1.3
Commit: 6876abd3730f205986c7df0a9084cc0f60284270
Parents: 8dc5fd0
Author: Tzu-Li (Gordon) Tai <tzulitai@apache.org>
Authored: Fri Aug 18 21:17:27 2017 +0800
Committer: Aljoscha Krettek <aljoscha.krettek@gmail.com>
Committed: Thu Aug 24 13:25:08 2017 +0200

----------------------------------------------------------------------
 .../connectors/kinesis/FlinkKinesisConsumer.java     | 15 ++++++---------
 .../kinesis/FlinkKinesisConsumerMigrationTest.java   |  9 ++++++---
 2 files changed, 12 insertions(+), 12 deletions(-)
----------------------------------------------------------------------


http://git-wip-us.apache.org/repos/asf/flink/blob/6876abd3/flink-connectors/flink-connector-kinesis/src/main/java/org/apache/flink/streaming/connectors/kinesis/FlinkKinesisConsumer.java
----------------------------------------------------------------------
diff --git a/flink-connectors/flink-connector-kinesis/src/main/java/org/apache/flink/streaming/connectors/kinesis/FlinkKinesisConsumer.java
b/flink-connectors/flink-connector-kinesis/src/main/java/org/apache/flink/streaming/connectors/kinesis/FlinkKinesisConsumer.java
index ea76ccc..aad622e 100644
--- a/flink-connectors/flink-connector-kinesis/src/main/java/org/apache/flink/streaming/connectors/kinesis/FlinkKinesisConsumer.java
+++ b/flink-connectors/flink-connector-kinesis/src/main/java/org/apache/flink/streaming/connectors/kinesis/FlinkKinesisConsumer.java
@@ -366,15 +366,12 @@ public class FlinkKinesisConsumer<T> extends RichParallelSourceFunction<T>
imple
 		LOG.info("Subtask {} restoring offsets from an older Flink version: {}",
 			getRuntimeContext().getIndexOfThisSubtask(), sequenceNumsToRestore);
 
-		if (restoredState.isEmpty()) {
-			sequenceNumsToRestore = null;
-		} else {
-			sequenceNumsToRestore = new HashMap<>();
-			for (Map.Entry<KinesisStreamShard, SequenceNumber> stateEntry : restoredState.entrySet())
{
-				sequenceNumsToRestore.put(
-						KinesisStreamShard.convertToStreamShardMetadata(stateEntry.getKey()),
-						stateEntry.getValue());
-			}
+		sequenceNumsToRestore = new HashMap<>();
+
+		for (Map.Entry<KinesisStreamShard, SequenceNumber> stateEntry : restoredState.entrySet())
{
+			sequenceNumsToRestore.put(
+				KinesisStreamShard.convertToStreamShardMetadata(stateEntry.getKey()),
+				stateEntry.getValue());
 		}
 	}
 

http://git-wip-us.apache.org/repos/asf/flink/blob/6876abd3/flink-connectors/flink-connector-kinesis/src/test/java/org/apache/flink/streaming/connectors/kinesis/FlinkKinesisConsumerMigrationTest.java
----------------------------------------------------------------------
diff --git a/flink-connectors/flink-connector-kinesis/src/test/java/org/apache/flink/streaming/connectors/kinesis/FlinkKinesisConsumerMigrationTest.java
b/flink-connectors/flink-connector-kinesis/src/test/java/org/apache/flink/streaming/connectors/kinesis/FlinkKinesisConsumerMigrationTest.java
index 3e43f13..a55a5c0 100644
--- a/flink-connectors/flink-connector-kinesis/src/test/java/org/apache/flink/streaming/connectors/kinesis/FlinkKinesisConsumerMigrationTest.java
+++ b/flink-connectors/flink-connector-kinesis/src/test/java/org/apache/flink/streaming/connectors/kinesis/FlinkKinesisConsumerMigrationTest.java
@@ -50,6 +50,7 @@ import java.util.concurrent.atomic.AtomicReference;
 
 import static org.junit.Assert.assertEquals;
 import static org.junit.Assert.assertNotEquals;
+import static org.junit.Assert.assertTrue;
 import static org.mockito.Mockito.doAnswer;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.when;
@@ -70,7 +71,7 @@ public class FlinkKinesisConsumerMigrationTest {
 	 */
 	private final MigrationVersion flinkGenerateSavepointVersion = null;
 
-	private final static HashMap<StreamShardMetadata, SequenceNumber> TEST_STATE = new
HashMap<>();
+	private static final HashMap<StreamShardMetadata, SequenceNumber> TEST_STATE = new
HashMap<>();
 	static {
 		StreamShardMetadata shardMetadata = new StreamShardMetadata();
 		shardMetadata.setStreamName("fakeStream1");
@@ -99,7 +100,9 @@ public class FlinkKinesisConsumerMigrationTest {
 		writeSnapshot("src/test/resources/kinesis-consumer-migration-test-flink" + flinkGenerateSavepointVersion
+ "-snapshot", TEST_STATE);
 
 		// write empty state snapshot
-		writeSnapshot("src/test/resources/kinesis-consumer-migration-test-flink" + flinkGenerateSavepointVersion
+ "-empty-snapshot", new HashMap<>());
+		writeSnapshot(
+			"src/test/resources/kinesis-consumer-migration-test-flink" + flinkGenerateSavepointVersion
+ "-empty-snapshot",
+			new HashMap<StreamShardMetadata, SequenceNumber>());
 	}
 
 	@Test
@@ -118,7 +121,7 @@ public class FlinkKinesisConsumerMigrationTest {
 		testHarness.open();
 
 		// assert that no state was restored
-		assertEquals(null, consumerFunction.getRestoredState());
+		assertTrue(consumerFunction.getRestoredState().isEmpty());
 
 		consumerOperator.close();
 		consumerOperator.cancel();


Mime
View raw message