[
https://issues.apache.org/jira/browse/FLINK-3211?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Tzu-Li (Gordon) Tai updated FLINK-3211:
---------------------------------------
Affects Version/s: 1.0.0
> Add AWS Kinesis streaming connector
> -----------------------------------
>
> Key: FLINK-3211
> URL: https://issues.apache.org/jira/browse/FLINK-3211
> Project: Flink
> Issue Type: New Feature
> Components: Streaming Connectors
> Affects Versions: 1.0.0
> Reporter: Tzu-Li (Gordon) Tai
> Original Estimate: 336h
> Remaining Estimate: 336h
>
> AWS Kinesis is a widely adopted message queue used by AWS users, much like a
> cloud service version of Apache Kafka. Support for AWS Kinesis will be a
> great addition to the handful of Flink's streaming connectors to external
> systems and a great reach out to the AWS community.
> After a first look at the AWS KCL (Kinesis Client Library), KCL already
> supports stream read beginning from a specific offset (or "record sequence
> number" in Kinesis terminology). For external checkpointing, KCL is designed
> to use AWS DynamoDB to checkpoint application state, where each partition's
> progress (or "shard" in Kinesis terminology) corresponds to a single row in
> the KCL-managed DynamoDB table.
> So, implementing the AWS Kinesis connector will very much resemble the work
> done on the Kafka connector, with a few different tweaks as following (I'm
> mainly just rewording [~StephanEwen]'s original description [1]):
> 1. Determine KCL Shard Worker to Flink source task mapping. KCL already
> offers worker tasks per shard, so we will need to do mapping much like [2].
> 2. Let the Flink connector also maintain a local copy of application state,
> accessed using KCL API, for the distributed snapshot checkpointing.
> 3. Restart the KCL at the last Flink local checkpointed record sequence upon
> failure. However, when KCL restarts after failure, it is originally designed
> to reference the external DynamoDB table. Need a further look on how to work
> with this so that the Flink checkpoint and external checkpoint in DynamoDB is
> properly synced.
> Most of the details regarding KCL's state checkpointing, sharding, shard
> workers, and failure recovery can be found here [3].
> As for the Kinesis sink connector, it should be fairly straightforward and
> almost, if not completely, identical to the Kafka sink. The Kinesis sink can
> be implemented with AWS KPL (Kinesis Producer Library) [4].
> On the other hand, while KCL and KPL are handy high-level APIs for AWS
> Kinesis, it might be preferable for the user to use the low-level API in AWS
> SDK instead, mainly due to the fact that the utilities of KCL and KPL come at
> cost, such as extra cost for the DynamoDB table for KCL checkpoint and the
> extra latency introduced by KPL [4]. For now, we can implement the initial
> version of the Kinesis connector with KCL and KPL, and leave the work for
> user-configurable switch for the connector to use low-level AWS SDK as future
> to-dos.
> References:
> [1]
> http://apache-flink-user-mailing-list-archive.2336050.n4.nabble.com/Kinesis-Connector-td2872.html
> [2] http://data-artisans.com/kafka-flink-a-practical-how-to/
> [3] http://docs.aws.amazon.com/kinesis/latest/dev/advanced-consumers.html
> [4]
> http://docs.aws.amazon.com//kinesis/latest/dev/developing-producers-with-kpl.html#d0e4998
--
This message was sent by Atlassian JIRA
(v6.3.4#6332)