You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Flink 1.17移除Kinesis连接器包,寻求替代实现方案

核心替代方案:社区独立维护的Kinesis连接器

Flink 1.17起官方不再内置flink-kinesis-connector,替代方案是使用Apache Flink社区独立维护的Kinesis连接器——它作为第三方连接器存在,完全兼容Flink 1.17及后续版本,核心API和功能与原连接器保持一致。

依赖配置

在项目构建文件中添加对应依赖(以Maven为例):

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-connector-kinesis</artifactId>
    <version>1.17.x</version> <!-- 匹配你的Flink 1.17具体版本,如1.17.1 -->
</dependency>

Gradle项目则使用:

implementation 'org.apache.flink:flink-connector-kinesis:1.17.x'

代码迁移注意事项

  • 原org.apache.flink.streaming.connectors.kinesis包下的类,现已迁移至org.apache.flink.connector.kinesis路径,仅需调整import语句即可,核心类名、方法签名基本无变化:
    • 原FlinkKinesisConsumer → 现KinesisSource
    • 原FlinkKinesisProducer → 现KinesisSink
  • Kinesis区域、凭证、消费起始位置等核心配置逻辑不变,只需对应迁移到新的Source/Sink构建器中。

迁移示例

消费Kinesis流(原FlinkKinesisConsumer迁移)

原代码片段:

FlinkKinesisConsumer<String> consumer = new FlinkKinesisConsumer<>(
    "my-stream",
    new SimpleStringSchema(),
    kinesisConfig
);
env.addSource(consumer);

迁移后代码:

KinesisSource<String> source = KinesisSource.<String>builder()
    .setStreamArn("arn:aws:kinesis:region:account-id:stream/my-stream")
    .setDeserializationSchema(new SimpleStringSchema())
    .setKinesisClientConfig(kinesisConfig)
    .setStartingPosition(StartingPosition.fromTimestamp(System.currentTimeMillis() - 3600000))
    .build();
env.fromSource(source, WatermarkStrategy.noWatermarks(), "Kinesis Source");

写入Kinesis流(原FlinkKinesisProducer迁移)

原代码片段:

FlinkKinesisProducer<String> producer = new FlinkKinesisProducer<>(
    new SimpleStringSchema(),
    kinesisConfig
);
producer.setDefaultStream("my-output-stream");
env.addSink(producer);

迁移后代码:

KinesisSink<String> sink = KinesisSink.<String>builder()
    .setStreamArn("arn:aws:kinesis:region:account-id:stream/my-output-stream")
    .setSerializationSchema(new SimpleStringSchema())
    .setKinesisClientConfig(kinesisConfig)
    .build();
env.sinkTo(sink);

额外说明

  • 该独立连接器会持续更新,兼容后续Flink版本,保留原连接器的Exactly-Once语义、批量写入、消费位点管理等核心特性。
  • 若使用AWS Kinesis Data Firehose,可引入flink-connector-kinesis-firehose依赖,迁移规则与上述一致。

内容的提问来源于stack exchange,提问作者priyadhingra19

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.11 17:22:38