Github user kl0u closed the pull request at:
https://github.com/apache/flink/pull/5230
---
Github user kl0u commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r166344325
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/KeyedStateBackend.java
---
@@ -38,6 +38,37 @@
*/
void setCurrentKey(K new
Github user kl0u commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r166344317
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/KeyedStateBackend.java
---
@@ -38,6 +38,37 @@
*/
void setCurrentKey(K new
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r166337384
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/KeyedStateBackend.java
---
@@ -38,6 +38,37 @@
*/
void setCurrentKey(K
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r166337949
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/KeyedStateBackend.java
---
@@ -38,6 +38,37 @@
*/
void setCurrentKey(K
Github user kl0u commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r165399750
--- Diff:
flink-streaming-java/src/main/java/org/apache/flink/streaming/api/graph/StreamGraphGenerator.java
---
@@ -586,7 +586,7 @@ private StreamGraph
genera
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r165398176
--- Diff:
flink-streaming-java/src/test/java/org/apache/flink/streaming/api/DataStreamTest.java
---
@@ -753,6 +763,182 @@ public void onTimer(
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r165395176
--- Diff:
flink-streaming-java/src/main/java/org/apache/flink/streaming/api/graph/StreamGraphGenerator.java
---
@@ -586,7 +586,7 @@ private StreamGraph
ge
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r165397632
--- Diff:
flink-streaming-java/src/test/java/org/apache/flink/streaming/api/DataStreamTest.java
---
@@ -753,6 +763,182 @@ public void onTimer(
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r165394684
--- Diff:
flink-streaming-java/src/main/java/org/apache/flink/streaming/api/functions/co/BaseBroadcastProcessFunction.java
---
@@ -0,0 +1,105 @@
+/*
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r165395860
--- Diff:
flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/co/CoBroadcastWithKeyedOperator.java
---
@@ -0,0 +1,323 @@
+/*
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r165387915
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/OperatorBackendSerializationProxy.java
---
@@ -57,18 +63,25 @@ public int getVersion()
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r165385868
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/DefaultOperatorStateBackend.java
---
@@ -148,21 +170,27 @@ public void close() throws
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r165388777
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/OperatorBackendSerializationProxy.java
---
@@ -77,16 +90,33 @@ public void read(DataIn
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r165378004
--- Diff:
flink-core/src/main/java/org/apache/flink/api/common/state/ReadWriteBroadcastState.java
---
@@ -0,0 +1,90 @@
+/*
+ * Licensed to the Apac
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r165377821
--- Diff:
flink-core/src/main/java/org/apache/flink/api/common/state/ReadOnlyBroadcastState.java
---
@@ -0,0 +1,59 @@
+/*
+ * Licensed to the Apach
Github user tzulitai commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r163938277
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/RegisteredBroadcastBackendStateMetaInfo.java
---
@@ -0,0 +1,232 @@
+/*
+ * Lic
Github user tzulitai commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r163933180
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/OperatorBackendSerializationProxy.java
---
@@ -77,16 +90,29 @@ public void read(DataIn
Github user tzulitai commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r163936390
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/OperatorStateHandle.java
---
@@ -36,8 +36,9 @@
* The modes that determine how
Github user tzulitai commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r163938102
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/RegisteredBroadcastBackendStateMetaInfo.java
---
@@ -0,0 +1,232 @@
+/*
+ * Lic
Github user tzulitai commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r163935065
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/OperatorBackendSerializationProxy.java
---
@@ -77,16 +90,29 @@ public void read(DataIn
Github user tzulitai commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r163935632
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/OperatorBackendStateMetaInfoSnapshotReaderWriters.java
---
@@ -211,4 +297,34 @@ public
Github user tzulitai commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r163932401
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/OperatorBackendSerializationProxy.java
---
@@ -64,11 +70,18 @@ public int getVersion()
Github user tzulitai commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r163942745
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/OperatorBackendSerializationProxy.java
---
@@ -77,16 +90,29 @@ public void read(DataIn
Github user tzulitai commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r163932956
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/OperatorBackendSerializationProxy.java
---
@@ -35,18 +35,24 @@
public st
Github user tzulitai commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r163934364
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/OperatorBackendStateMetaInfoSnapshotReaderWriters.java
---
@@ -112,13 +142,40 @@ publi
Github user tzulitai commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r163930707
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/DefaultOperatorStateBackend.java
---
@@ -137,7 +155,12 @@ public ExecutionConfig getEx
Github user tzulitai commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r163932340
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/OperatorBackendSerializationProxy.java
---
@@ -64,11 +70,18 @@ public int getVersion()
Github user tzulitai commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r163936903
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/OperatorStateHandle.java
---
@@ -36,8 +36,9 @@
* The modes that determine how
Github user kl0u commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r163872816
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/BackendWritableBroadcastState.java
---
@@ -0,0 +1,42 @@
+/*
+ * Licensed to the Ap
Github user kl0u commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r163870691
--- Diff:
flink-core/src/main/java/org/apache/flink/api/common/state/BroadcastState.java
---
@@ -0,0 +1,48 @@
+/*
+ * Licensed to the Apache Software F
Github user kl0u commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r163870666
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/DefaultOperatorStateBackend.java
---
@@ -513,17 +630,100 @@ public void addAll(List values
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r163841958
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/BackendWritableBroadcastState.java
---
@@ -0,0 +1,42 @@
+/*
+ * Licensed to th
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r163850320
--- Diff:
flink-runtime/src/test/java/org/apache/flink/runtime/state/OperatorStateBackendTest.java
---
@@ -753,4 +935,27 @@ private static Environment crea
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r163849551
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/DefaultOperatorStateBackend.java
---
@@ -601,21 +805,43 @@ public void addAll(List val
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r163840020
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/DefaultOperatorStateBackend.java
---
@@ -513,17 +630,100 @@ public void addAll(List va
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r163839057
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/RoundRobinOperatorStateRepartitioner.java
---
@@ -228,36 +221,58 @@ private Group
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r163838390
--- Diff:
flink-core/src/main/java/org/apache/flink/api/common/state/BroadcastState.java
---
@@ -0,0 +1,48 @@
+/*
+ * Licensed to the Apache Softwa
Github user kl0u commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r160904948
--- Diff:
flink-streaming-java/src/main/java/org/apache/flink/streaming/api/datastream/BroadcastConnectedStream.java
---
@@ -0,0 +1,216 @@
+/*
+ * Lice
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r159222752
--- Diff:
flink-streaming-java/src/main/java/org/apache/flink/streaming/api/graph/StreamGraphGenerator.java
---
@@ -586,7 +586,7 @@ private StreamGraph
ge
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r159224065
--- Diff:
flink-streaming-java/src/main/java/org/apache/flink/streaming/api/functions/co/KeyedBroadcastProcessFunction.java
---
@@ -0,0 +1,178 @@
+/*
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r159215700
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/RoundRobinOperatorStateRepartitioner.java
---
@@ -42,10 +42,10 @@
@Overri
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r159216046
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/state/AbstractKeyedStateBackend.java
---
@@ -99,7 +99,7 @@
private final ExecutionCo
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r159222186
--- Diff:
flink-runtime/src/test/java/org/apache/flink/runtime/state/SerializationProxiesTest.java
---
@@ -287,16 +339,60 @@ public void
testOperatorState
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r159223446
--- Diff:
flink-streaming-java/src/test/java/org/apache/flink/streaming/api/operators/co/CoBroadcastWithKeyedOperatorTest.java
---
@@ -0,0 +1,734 @@
+/
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r159215876
--- Diff:
flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/StateAssignmentOperation.java
---
@@ -613,7 +613,6 @@ private static void checkSt
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r159222356
--- Diff:
flink-streaming-java/src/main/java/org/apache/flink/streaming/api/datastream/BroadcastConnectedStream.java
---
@@ -0,0 +1,216 @@
+/*
+ *
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r159223798
--- Diff:
flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/co/CoBroadcastWithKeyedOperator.java
---
@@ -0,0 +1,369 @@
+/*
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/5230#discussion_r159222137
--- Diff:
flink-runtime/src/test/java/org/apache/flink/runtime/state/SerializationProxiesTest.java
---
@@ -287,16 +339,60 @@ public void
testOperatorState
GitHub user kl0u opened a pull request:
https://github.com/apache/flink/pull/5230
[FLINK-8345] Add iterator of keyed state on broadcast side of connected
streams.
*Thank you very much for contributing to Apache Flink - we are happy that
you want to help us improve Flink. To help th
50 matches
Mail list logo