|
21 | 21 | import static org.junit.Assert.assertThat; |
22 | 22 |
|
23 | 23 | import java.io.Serializable; |
24 | | -import java.util.Collection; |
25 | | -import java.util.HashMap; |
26 | | -import java.util.Iterator; |
27 | | -import java.util.List; |
28 | | -import java.util.Map; |
| 24 | +import java.util.*; |
29 | 25 | import java.util.Map.Entry; |
30 | | -import java.util.Set; |
31 | 26 | import java.util.stream.Stream; |
32 | 27 | import javax.annotation.Nonnull; |
33 | 28 | import org.apache.flink.api.common.ExecutionConfig; |
34 | 29 | import org.apache.flink.api.common.JobID; |
35 | | -import org.apache.flink.api.common.accumulators.Accumulator; |
36 | | -import org.apache.flink.api.common.accumulators.DoubleCounter; |
| 30 | +import org.apache.flink.api.common.accumulators.*; |
37 | 31 | import org.apache.flink.api.common.accumulators.Histogram; |
38 | | -import org.apache.flink.api.common.accumulators.IntCounter; |
39 | | -import org.apache.flink.api.common.accumulators.LongCounter; |
40 | 32 | import org.apache.flink.api.common.cache.DistributedCache; |
41 | 33 | import org.apache.flink.api.common.externalresource.ExternalResourceInfo; |
42 | 34 | import org.apache.flink.api.common.functions.BroadcastVariableInitializer; |
43 | 35 | import org.apache.flink.api.common.functions.RuntimeContext; |
44 | | -import org.apache.flink.api.common.state.AggregatingState; |
45 | | -import org.apache.flink.api.common.state.AggregatingStateDescriptor; |
46 | | -import org.apache.flink.api.common.state.ListState; |
47 | | -import org.apache.flink.api.common.state.ListStateDescriptor; |
48 | | -import org.apache.flink.api.common.state.MapState; |
49 | | -import org.apache.flink.api.common.state.MapStateDescriptor; |
50 | | -import org.apache.flink.api.common.state.ReducingState; |
51 | | -import org.apache.flink.api.common.state.ReducingStateDescriptor; |
52 | | -import org.apache.flink.api.common.state.State; |
53 | | -import org.apache.flink.api.common.state.StateDescriptor; |
54 | | -import org.apache.flink.api.common.state.ValueState; |
55 | | -import org.apache.flink.api.common.state.ValueStateDescriptor; |
| 36 | +import org.apache.flink.api.common.state.*; |
56 | 37 | import org.apache.flink.api.common.typeutils.TypeSerializer; |
57 | 38 | import org.apache.flink.api.java.tuple.Tuple2; |
58 | | -import org.apache.flink.metrics.CharacterFilter; |
59 | | -import org.apache.flink.metrics.Counter; |
60 | | -import org.apache.flink.metrics.Gauge; |
61 | | -import org.apache.flink.metrics.Meter; |
62 | | -import org.apache.flink.metrics.MetricGroup; |
63 | | -import org.apache.flink.metrics.SimpleCounter; |
| 39 | +import org.apache.flink.metrics.*; |
64 | 40 | import org.apache.flink.metrics.groups.OperatorMetricGroup; |
65 | | -import org.apache.flink.runtime.state.KeyGroupedInternalPriorityQueue; |
66 | | -import org.apache.flink.runtime.state.Keyed; |
67 | | -import org.apache.flink.runtime.state.KeyedStateBackend; |
68 | | -import org.apache.flink.runtime.state.KeyedStateFunction; |
69 | | -import org.apache.flink.runtime.state.PriorityComparable; |
70 | | -import org.apache.flink.runtime.state.StateSnapshotTransformer.StateSnapshotTransformFactory; |
71 | | -import org.apache.flink.runtime.state.VoidNamespace; |
| 41 | +import org.apache.flink.runtime.state.*; |
72 | 42 | import org.apache.flink.runtime.state.heap.HeapPriorityQueueElement; |
73 | 43 | import org.apache.flink.runtime.state.internal.InternalListState; |
74 | 44 | import org.apache.flink.shaded.guava30.com.google.common.util.concurrent.MoreExecutors; |
@@ -352,10 +322,13 @@ public boolean deregisterKeySelectionListener(KeySelectionListener<Object> liste |
352 | 322 |
|
353 | 323 | @Nonnull |
354 | 324 | @Override |
355 | | - public <N, SV, SEV, S extends State, IS extends S> IS createInternalState( |
356 | | - @Nonnull TypeSerializer<N> namespaceSerializer, |
357 | | - @Nonnull StateDescriptor<S, SV> stateDesc, |
358 | | - @Nonnull StateSnapshotTransformFactory<SEV> snapshotTransformFactory) { |
| 325 | + public <N, SV, SEV, S extends State, IS extends S> IS createOrUpdateInternalState( |
| 326 | + @Nonnull TypeSerializer<N> typeSerializer, |
| 327 | + @Nonnull StateDescriptor<S, SV> stateDescriptor, |
| 328 | + @Nonnull |
| 329 | + StateSnapshotTransformer.StateSnapshotTransformFactory<SEV> |
| 330 | + stateSnapshotTransformFactory) |
| 331 | + throws Exception { |
359 | 332 | throw new UnsupportedOperationException(); |
360 | 333 | } |
361 | 334 |
|
|
0 commit comments