71、Flink 的 Hybrid Source 详解
创始人
2025-01-08 01:35:18
0
Hybrid Source
1.概述

Hybrid Source 解决了从异构数据源顺序读取输入以生成单个输入流的问题。

示例:从 S3 读取前几天的有界输入,然后使用 Kafka 的最新无界输入,当有界文件输入完成而不中断应用程序时 Hybrid Source 会从 FileSource 切换到 KafkaSource。

在 Hybrid Source 出现之前,需要创建一个具有多个源的拓扑结构,并由用户定义切换机制;使用 HybridSource 之后,从 DataStream API 的角度看,多个源在 Flink 作业图中显示为单个源。

需要依赖如下:

     org.apache.flink     flink-connector-base     1.19.0  
2.下一个源的起始位置

要在一个 Hybrid Source 中排列多个源,除最后一个源外的所有源都需要有界;因此通常需要为源分配一个开始和结束位置。

a)固定起始位置

示例:从文件中读取到预先确定的切换时间,然后继续从 Kafka 中读取,每个源都覆盖了预先已知的范围,可以像直接使用一样预先创建包含的源。

long switchTimestamp = ...; // derive from file input paths  FileSource fileSource =   FileSource.forRecordStreamFormat(new TextLineInputFormat(), Path.fromLocalFile(testDir)).build();  KafkaSource kafkaSource =           KafkaSource.builder()                   .setStartingOffsets(OffsetsInitializer.timestamp(switchTimestamp + 1))                   .build();  HybridSource hybridSource =           HybridSource.builder(fileSource)                   .addSource(kafkaSource)                   .build(); 
b)动态其实位置

示例:文件源需要读取的数据量很大,可能比下一个源可用的保留时间更长,切换需要在 “当前时间-X” 发生。

因此要将下一个源的启动时间设置为切换时间,需要从以前的文件枚举器中转移结束位置,以便通过实现 SourceFactory 来延迟构建KafkaSource。

注意:枚举器需要支持获取结束时间戳。

FileSource fileSource = CustomFileSource.readTillOneDayFromLatest();  HybridSource hybridSource =     HybridSource.builder(fileSource)         .addSource(             switchContext -> {               CustomFileSplitEnumerator previousEnumerator =                   switchContext.getPreviousEnumerator();                              // how to get timestamp depends on specific enumerator               long switchTimestamp = previousEnumerator.getEndTimestamp();                              KafkaSource kafkaSource =                   KafkaSource.builder()                       .setStartingOffsets(OffsetsInitializer.timestamp(switchTimestamp + 1))                       .build();                              return kafkaSource;             },             Boundedness.CONTINUOUS_UNBOUNDED)         .build(); 

相关内容

热门资讯

解迷一下!微信边锋辅助挂件,约... 您好,微信边锋辅助挂件这款游戏可以开挂的,确实是有挂的,需要了解加去威信【136704302】很多玩...
曝光一下!江西中至小程序黑科技... 曝光一下!江西中至小程序黑科技,竹间穿有挂没,原来是有挂(哔哩哔哩)1、每一步都需要思考,不同水平的...
推荐一下!丰城呱呱辅助器,友乐... 推荐一下!丰城呱呱辅助器,友乐广西app下载安装,总是真的是有挂(哔哩哔哩)在进入丰城呱呱辅助器软件...
推荐一下!家乡大贰祈福有用吗,... 推荐一下!家乡大贰祈福有用吗,中至赣州黑科技辅助软件,切实真的是有挂(哔哩哔哩)1)家乡大贰祈福有用...
有挂一下!微玩体育辅助器,we... 有挂一下!微玩体育辅助器,wepkerplus辅助,本来真的有挂(哔哩哔哩)微玩体育辅助器辅助器是一...
专业一下!九九联盟辅助神器,随... 专业一下!九九联盟辅助神器,随意玩俱乐部辅助,本来是真的有挂(哔哩哔哩)亲,关键说明,随意玩俱乐部辅...
解密一下!山西奇迹打锅子辅助,... 解密一下!山西奇迹打锅子辅助,新老夫子较二八年,果然真的是有挂(哔哩哔哩)1、打开软件启动之后找到中...
开挂一下!开心泉州小程序辅助哪... 开挂一下!开心泉州小程序辅助哪里查看,微信小游戏破解版,竟然真的有挂(哔哩哔哩)1、开心泉州小程序辅...
解迷一下!欢乐茶馆脚本辅助,四... 解迷一下!欢乐茶馆脚本辅助,四川家园辅助软件,本来存在有挂(哔哩哔哩)1、四川家园辅助软件辅助器安装...
科普一下!佛手在线辅助器苹果版... 科普一下!佛手在线辅助器苹果版,神殿娱乐控制系统,好像是真的有挂(哔哩哔哩)1、佛手在线辅助器苹果版...