[
https://issues.apache.org/jira/browse/FLINK-4844?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=15582783#comment-15582783
]
ASF GitHub Bot commented on FLINK-4844:
---------------------------------------
Github user aljoscha commented on a diff in the pull request:
https://github.com/apache/flink/pull/2648#discussion_r83662569
--- Diff:
flink-streaming-connectors/flink-connector-kafka-base/src/main/java/org/apache/flink/streaming/connectors/kafka/FlinkKafkaConsumerBase.java
---
@@ -80,38 +81,38 @@
//
------------------------------------------------------------------------
private final List<String> topics;
-
+
/** The schema to convert between Kafka's byte messages, and Flink's
objects */
protected final KeyedDeserializationSchema<T> deserializer;
/** The set of topic partitions that the source will read */
protected List<KafkaTopicPartition> subscribedPartitions;
-
+
/** Optional timestamp extractor / watermark generator that will be run
per Kafka partition,
* to exploit per-partition timestamp characteristics.
* The assigner is kept in serialized form, to deserialize it into
multiple copies */
private SerializedValue<AssignerWithPeriodicWatermarks<T>>
periodicWatermarkAssigner;
-
+
/** Optional timestamp extractor / watermark generator that will be run
per Kafka partition,
- * to exploit per-partition timestamp characteristics.
+ * to exploit per-partition timestamp characteristics.
* The assigner is kept in serialized form, to deserialize it into
multiple copies */
private SerializedValue<AssignerWithPunctuatedWatermarks<T>>
punctuatedWatermarkAssigner;
- private transient OperatorStateStore stateStore;
+ private transient ListState<Serializable> offsetsStateForCheckpoint;
--- End diff --
This can can have a more concrete type. You changed
`OperatorStateStore.getSerializableListState` to this:
```
<T extends Serializable> ListState<T> getSerializableListState(String
stateName) throws Exception;
```
> Partitionable Raw Keyed/Operator State
> --------------------------------------
>
> Key: FLINK-4844
> URL: https://issues.apache.org/jira/browse/FLINK-4844
> Project: Flink
> Issue Type: New Feature
> Reporter: Stefan Richter
> Assignee: Stefan Richter
>
> Partitionable operator and keyed state are currently only available by using
> backends. However, the serialization code for many operators is build around
> reading/writing their state to a stream for checkpointing. We want to provide
> partitionable states also through streams, so that migrating existing
> operators becomes more easy.
--
This message was sent by Atlassian JIRA
(v6.3.4#6332)