资讯动态

Java面试——Flink原理及应用(二)

发布时间:2026/8/27 8:20:39 来源:尧图企业网站定制
Flink原理及应用2、Flink的应用2.1、Flink的安装2.1.1、Flink的独立集群模式安装2.1.2、Flink的高可用模式安装2.2、Flink实战案例2.2.1、构建一个基于Java的Flink应用2.2.2、Flink DataStream API2.2.3、Flink DataSet API2.2.4、Flink Table API和SQL2、Flink的应用Flink被广泛应用于流式计算和批量计算中其高吞吐量、低延迟和简单易用的API得到业界的一致好评。下面对Flink的安装和Flink API的使用进行简单介绍。2.1、Flink的安装Flink可以在独立的集群模式下运行也可以在YARN、Mesos、Docker、Kubernetes等虚拟化环境中运行还可以在AWS、Google Computer Engine等云环境中运行。下面介绍Flink的独立集群模式安装和高可用模式安装对于其他运行模式请读者参照官网。2.1.1、Flink的独立集群模式安装1安装Java环境Flink要求JDK版本大于1.8。2各个服务器之间互信配置。3从Apache官网下载Flink最新安装包目前最新版本为1.9。安装包中1.9.1为Flink的版本2.11为Flink对应的Scala版本。下载页面如图所示。4将Flink安装包复制到待安装目录执行以下命令解压Flink安装包并进入安装目录。5执行以下命令编辑Flink配置文件flink-conf.yaml。6在flink-conf.yaml中修改如下核心配置。对上述配置参数说明如下。①jobmanager.heap.sizeJob Manager JVM的堆内存大小。由于Job Manager主要用于Job分发和运行状态监控等具体的计算任务均在Task Manager上被执行因此不需要太大的内存资源。②taskmanager.heap.sizeTask Manager JVM的堆内存大小。Flink计算任务均在Task Manager上被执行且Flink任务主要在内存中完成快速计算因此一般会把操作系统主要的内存都分配给该任务。如果服务器只有一个Flink Task Manager则建议内存分配策略为操作系统预留10%的内存资源taskmanager.heap.size预留60%的内存资源其余内存资源留给堆外内存使用。③taskmanager.numberOfTaskSlots每个TaskManager分配的Slot数量。每个Slot都对应一个串行的 Pipeline一般配置的建议值为CPU核数这样可以避免CPU上下文切换使得一个CPU持续为某个任务服务。④parallelism.defaultFlink集群的默认并行度默认值为1。⑤io.tmp.dirsFlink临时数据的存放地址。⑥state.backendBackend计算状态的存储方式包括jobmanagerMemery StateBackend​、filesystemFsStateBackend和rocksdbRockDBStateBockend​。⑦state.checkpoints.dirCheckPoint的存储地址用于计算状态的存储和故障恢复。存储地址需要与state.backend的配置对应。当Flink任务启动时可以指定CheckPoint地址让应用程序从指定状态开始运行。从指定状态启动的命令如下。⑧state.savepoints.dirSavePoint的存储地址需要与state.backend的配置对应。⑨taskmanager.network.memory.fraction用于网络缓冲区的JVM内存份数该参数决定了Task Manager可以同时具有多少个流数据交换通道及通道的缓冲程度。⑩taskmanager.network.memory.min网络缓冲区的内存最小值。⑪taskmanager.network.memory.max网络缓冲区的内存最大值由于Flink集群中的数据大多来自外部系统比如Kafka消息系统、Socket端口、S3等这些数据均要通过网络的形式在集群中传输其值一般为该集群能处理的最大数据量。7执行以下命令配置Master地址。8在masters配置文件中加入Master服务地址和端口号具体如下。9执行以下命令配置Slave地址。10在slaves配置文件中加入Slave服务地址具体如下。11执行以下命令启动集群。12启动日志如下。13执行Jps查看进程。14在浏览器地址栏中输入127.0.0.18081访问Flink集群监控界面如图所示。2.1.2、Flink的高可用模式安装在Flink集群中Job Manager是任务的管理节点TaskManager是任务的计算节点因此Task Manager可以按照计算需求进行水平扩展。Flink的高可用HighAvailabilityHA主要指Job Manager的高可用。Job Manager的高可用依赖ZooKeeper进行节点状态的管理和选举工作。Job Manager高可用模式为主Leader从Flower模式在正常情况下由JobManager主节点为整个集群提供服务Job Manager从节点处于Standby状态。在主节点宕机后集群会在现有从节点中选举一个节点作为主节点对外提供服务在之前主节点启动后将以从节点的角色加入集群。一个包含2个JobManager和3个Task Manager的Flink HA集群如图所示。Flink高可用模式的配置如下。1ZooKeeper的安装、配置和启动具体步骤请读者参照ZooKeeper章节。2执行以下命令编辑conf/flink-conf.yaml配置文件。3在conf/flink-conf.yaml配置文件中加入以下HA配置。上述配置中的high-availabilityzookeeper用于声明Flink集群以高可用模式运行且高可用模式依赖ZooKeeper服务。high-availability.zookeeper.quorum 129.168.1.1002181配置了高可用模式下的ZooKeeper地址high-availability.storageDir配置了高可用模式下的状态存储地址该地址必须是一个Job Manager和Task Manager均可访问的地址比如HDFS、S3、Cpeh、NFS等。Flink会将集群运行过程的主要数据存储在该地址上当JobManager发生意外宕机时集群会从high- availability.storageDir中获取最近一次的状态数据进行状态还原和故障恢复。4按照以下配置修改conf/master配置文件。上述配置定义了2个Job Manager服务地址分别为129.168.1.1008081和129.168.1.1018081。5按照以下配置修改conf/salver配置文件。上述配置定义了3个Task Manager服务地址分别为192.168.1.102、192.168.1.103、192.168.1.104。6执行以下命令启动集群。2.2、Flink实战案例Flink支持基于Java、Scala和Python进行开发。下面以Java为基础介绍Flink的使用。2.2.1、构建一个基于Java的Flink应用构建一个简单的Flink应用的步骤如下。1新建一个名为Flink的Maven项目。2在pom.xml中添加以下Flink依赖包。在上述依赖中flink-java是Flink开发必需的核心包flink-streaming-java_2.11是Flink流数据开发的依赖包。其中scope为provided的意思指将依赖包进行编译但是不将其打包到JAR文件中这样做不但能够减小JAR文件的大小还能避免客户端依赖包版本和服务端依赖包版本不一致引起的JAR包冲突问题。3Maven构建配置在pom.xml中添加Maven构建插件。4新建一个SocketFlinkDemo类在类中定义一个简单的Flink统计程序。上述代码定义了一个简单的监听和获取Socket上输入的单词并对单词进行统计的Flink程序。其核心步骤如下。1定义获取StreamExecutionEnvironment实例envStreamExecutionEnvironment是Flink程序运行的上下文环境每个流式计算都在一个StreamExecutionEnvironment上运行。2定义数据源上述代码通过调用env的socketTextStream方法定义一个数据流来实时监听和获取对应的Socket上的数据。3定义转换流上述代码定义了3个转换流分别是flatMap、keyBy、reduce操作。flatMap用于遍历每个数据并获取数据内容构造WordWithCount数据结构并输出keyBy以word字段为Key对数据进行分组reduce对单词执行统计计算并返回计算结果。上述代码还定义了一个长度为5s的时间窗口来定时触发计算。4定义数据输出为数据流添加一个Sink算子将计算结果输出在上述代码中Sink通过覆写invoke方法获取计算结果并打印计算结果。5触发任务执行调用env的execute触发并执行任务。在上述计算中使用到的数据模型WordWithCount如下。6打开终端输入nc-l 8000在8000端口上构建一个Socket端口监听。7程序运行在SocketFlinkDemo上单击右键执行该Flink程序Flink会在程序内部启动一个程序级别的集群Flink Application Cluster并运行该程序。FlinkApplication Cluster是一个专用的Flink Cluster仅用于执行单个Flink Job。Flink Cluster的生命周期与FlinkJob的生命周期对应。在运行后可以看到如下Flink启动核心日志。注意以下日志仅为与集群启动状态有关的核心日志其他日志省略以#开头的注释用于方便读者阅读在真实日志中不存在该注释。- 8在Socket端口上输入数据在终端输入以下字符串。- 9查看Flink程序统计结果。2.2.2、Flink DataStream APIDataStream是Flink基于数据流封装的标准流式处理API支持数据过滤Filter​、状态更新UpdateState​、窗口定义Windows Defining和数据聚合Aggregating等操作。DataStream的数据来源多样包括Socket服务、消息系统Kafka、Kinesis、RabbitMQ​、文件系统。其计算结果可以通过Sink操作写入文件、标准数据库存储、消息系统、日志等任何应用程序可以输出的目标地址。DataStream API编程的核心流程是获取数据源DataSource​、执行数据转换Data Transformation操作、输出计算结果Data Sink​。** 1数据源。**数据源是流式计算的起点Flink通过数据源将外界数据接入系统进行计算。应用程序客户端通过StreamExecutionEnvironment.addSourcesourceFunction向数据流加入一个数据源。Flink默认实现了多种数据源具体如表所示。2数据转换操作。数据转换操作将一个或多个DataStream转换为一个新的DataStream。应用程序通过连接和组合多个数据转换操作形成一个数据流拓扑图。常用的数据转换操作如下。1Map。转换操作名称Map。转换操作说明输入DataStream的元素通过Map函数转换返回一个新的元素。转换操作类型将DataStream转换为DataStream。使用示例将DataStream中的每个Integer类型的元素都计算平方后输出。2FlatMap。转换操作名称FlatMap。转换操作说明输入DataStream的元素通过FlatMap函数转换返回零个、一个或多个新的元素。转换操作类型将DataStream转换为DataStream。使用示例将字符串按照空格分隔输出。3Filter。转换操作名称Filter。转换操作说明对每个元素都进行布尔函数运算返回布尔函数运算结果为true的元素。转换操作类型将DataStream转换为DataStream。使用示例返回DataStream中的偶数。4KeyBy。转换操作名称KeyBy。转换操作说明在逻辑上将数据流划分为多个分区所有具有相同Key的元素都被分配到同一个分区。keyBy​内部基于Hash函数实现分区。有多种指定Key的方式。转换操作类型将DataStream转换为KeyedStream。使用示例以someKey为Key对数据分区以Tuple的第一个元素为Key对数据分区。注意POJO类型的数据如果没有实现hashCode方法则不能进行分区。5Reduce。转换操作名称Reduce。转换操作说明对KeyedStream上的数据进行Reduce操作返回一个新的元素。转换操作类型将KeyedStream转换为DataStream。使用示例对KeyedStream上的元素求和。6Fold。转换操作名称Fold。转换操作说明对KeyedStream上的元素进行叠加操作将当前元素和上次叠加操作的元素进行组合输出一个新的元素。转换操作类型将KeyedStream转换为DataStream。使用示例对序列12345执行Fold后输出输出结果为​“start-1”​​“start-1-2”​​“start-1-2-3”……7Aggregation。转换操作名称Aggregation。转换操作说明对KeyedStream上的元素进行聚合操作min返回统计结果中的最小值minBy返回统计结果中最小值对应的元素。转换操作类型将KeyedStream转换为DataStream。使用示例对KeyedStream中的元素进行求和、求最小值、求最大值。8Window。转换操作名称Window。转换操作说明Window操作根据设置例如在过去5s内到达的数据将每个KeyedStream中的数据都分配到Window中Flink将流数据分配到不同的Window中执行计算。转换操作类型将KeyedStream转换为WindowedStream。使用示例定义一个每5s执行一次的Window翻滚窗口。9WindowAll。转换操作名称WindowAll。转换操作说明WindowAll操作根据设置例如在过去5s内到达的数据将所有DataStream中的数据都分配到一个Window中WindowAll为非并行转换操作。执行WindowAll操作后的所有元素都会被分配到一个Task上进行计算。转换操作类型将DataStream转换为AllWindowedStream。使用示例定义一个每5s执行一次的WindowAll翻滚窗口。10Window Apply。转换操作名称Window Apply。转换操作说明自定义窗口计算函数将WindowedStream或AllWindowedStream中的元素经过Apply函数计算转换为DataStream。转换操作类型将WindowedStream或AllWindowedStream转换为DataStream。使用示例定义一个每10s触发一次的TimeWindowAll翻滚窗口统计每个Window中的数据个数和最大值。上述代码通过timeWindowAllTime.seconds10​定义了一个每10s触发一次的TimeWindowAll翻滚窗口并在apply方法中实现了对窗口中的数据values进行遍历计算窗口中数据的最大值最后将结果输出。同时为了便于读者理解apply方法统计并打印了窗口中数据的个数和窗口的开始时间和结束时间。向socketTextStream中输入数据1、3、5打印出以下结果。通过输出能够明显地看出窗口的开始时间1575123820000ms和结束时间1575123830000ms相差10s也就是TimeWindowAll翻滚窗口的时间长度。11Window Reduce。转换操作名称Window Reduce。转换操作说明对WindowedStream上的数据执行Reduce操作并将Reduce的结果返回。转换操作类型将WindowedStream转换为DataStream。使用示例对WindowedStream中的数据类型为Tuple2StringInteger的第二个值执行求和操作。12Window Fold。转换操作名称Window Fold。转换操作说明对WindowedStream上的数据执行叠加操作。转换操作类型将WindowedStream转换为DataStream。使用示例将WindowedStream中的数据使用“_”连接起来。13Aggregation on Window。转换操作名称Aggregation on Window。转换操作说明按照指定的统计函数对WindowedStream上的数据执行统计。转换操作类型将WindowedStream转换为DataStream。使用示例对WindowedStream中的元素进行求和、最小值、最大值统计。14Union。转换操作名称Union。转换操作说明将多个DataStream联合Union起来构成一个新的DataStream。转换操作类型将多个DataStream转换为一个DataStream。使用示例将dataStream1和otherStream2、otherStream3等多个dataStream联合起来构成一个新的dataStream。15Window Join。转换操作名称Window Join。转换操作说明根据指定的Key将两个DataStream上的数据流进行Join。转换操作类型将两个DataStream转换为一个DataStream。使用示例将otherStream和dataStream进行JoinJoin的结果为一个每5s执行一次的翻滚窗口。16Connect。转换操作名称Connect。转换操作说明将两个DataStream连接生成一个新的ConnectedStreams数据流。转换操作类型将两个DataStream转换为一个ConnectedStreams。使用示例定义两个DataStream并将其连接起来。17CoMap和CoFlatMap。转换操作名称CoMap和CoFlatMap。转换操作说明CoMap、CoFlatMap的功能与Map、FlatMap的功能类似不同之处在于CoMap、CoFlatMap作用于ConnectedStreams上。转换操作类型将ConnectedStreams转换为DataStream。使用示例在connectedStreams上分别执行map和flatMap操作。18Split。转换操作名称Split。转换操作说明将DataStream按照指定的规则拆分Split为多个SplitStream。转换操作类型将DataStream转换为SplitStream。使用示例将someDataStream拆分为两个dataStream一个outputName为even另一个outputName为odd。19Select。转换操作名称Select。转换操作说明从SplitStream中通过选择器选择部分数据组成一个新的DataStream。转换操作类型将SplitStream转换为DataStream。使用示例将someDataStream数据流根据选择器分为三个dataStream一个outputName为even一个outputName为odd另外一个outputName为even和odd。20Custom Partitioning。转换操作名称Custom Partitioning。转换操作说明对DataStream上的数据按照指定Key进行分区操作。转换操作类型将DataStream转换为DataStream。使用示例将dataStream分别按照“someKey”分区和按照第一个值自动分区。21Random Partitioning。转换操作名称Random Partitioning。转换操作说明对DataStream上的数据进行随机分区操作随机分区的结果一般比较均匀。转换操作类型将DataStream转换为DataStream。使用示例将dataStream随机分区。22Broadcasting。转换操作名称Broadcasting。转换操作说明将DataStream中的数据一般为配置信息等公共数据流广播Broadcasting到每个分区上。转换操作类型将DataStream转换为DataStream。使用示例对dataStream中的数据进行广播。3数据输出。Flink可以将计算的结果数据输出到File、Socket、外部存储或者将结果打印到日志中。常用的Flink数据输出操作如表所示。除了上述输出Flink还提供了丰富的连接器用于将结果输出到Kafka、Cassandra、Redis等消息系统和数据库。1将结果输出到Kafka。如下代码将dataStream中的数据输出到Kafka。2将结果输出到RabbitMQ。如下代码将dataStream中的数据输出到RabbitMQ。3将结果输出到Cassandra。如下代码将dataStream中的数据输出到Cassandra。除了基于SQL的方式将结果输出到CassandraFlink还执行以POJO形式将结果输出到Cassandra。4将结果输出到Kinesis。如下代码将dataStream中的数据输出到Kinesis。5将结果输出到ElasticSearch。如下代码将dataStream中的数据输出到ElasticSearch。6将结果输出到Hadoop FileSystem。如下代码将dataStream中的数据输出到HadoopFileSystem。7将结果输出到Flume。如下代码将dataStream中的数据输出到Flume。8将结果输出到Redis。如下代码将dataStream中的数据输出到Redis。9自定义Sink。如下代码将自定义一个SinkLog将dataStream中的数据通过SinkLog输出。2.2.3、Flink DataSet APIFlink封装了DataSet API用于处理批量数据。Flink处理批量数据的流程为Flink系统内部将接收的数据转换成DataSet数据集然后将DataSet数据集并行分布在集群的每个节点上最后基于DataSet数据集进行各种转换操作例如Map、Filter等​最终通过DataSink操作将结果数据集输出到外部系统。1数据源。Flink可以基于通用文件格式创建数据集。其创建数据集的方法位于ExecutionEnvironment。常用的基于文件系统创建DataSet的操作如下。1readTextFilepath​按行读取文件并将其作为字符串返回。2readTextFileWithValuepath​按行读取文件并将其作为可变字符串返回。3readCsvFilepath​按行读取文件使用逗号或其他字符分隔并返回元组或对象的数据集。4readFileOfPrimitivespathdelimiter​解析基本数据类型字符串或整数的文件以换行符或其他字符分隔。5readSequenceFileKeyValuepath​从指定路径读取并解析SequenceFile返回KeyValue元组。常用的基于集合创建DataSet的操作如下。1fromCollectionSeq​用对象集合创建数据集。集合中的所有元素必须属于同一类型。2fromCollectionIterator​用迭代器创建数据集指定迭代器返回元素的数据类型。3fromElementselements_*​根据给定的对象序列创建数据集。所有对象必须属于同一类型。4fromParallelCollectionSplittableIterator​并行地从迭代器创建数据集。指定迭代器返回元素的数据类型。5generateSequencefromto​并行生成给定间隔的数字序列。常用的基于文件系统创建DataSet的操作如下。1readFileinputFormatpath​接收文件输入格式为inputFormat的数据。2createInputinputFormat​接收通用输入格式Flink会自动识别文件格式。2数据转换操作。DataSet的Transformation操作和DataStream基本类似常用方法如下。1Map根据一个元素生成一个新的元素。2FlatMap根据一个元素生成多个元素。3MapPartition将DataStream进行分区操作。4Filter对每个数据都执行布尔函数只保存函数返回true的数据。5Distinct对数据集中的元素去重并返回新的数据集。6Reduce作用于整个DataSet合并该数据集的元素。7Aggregate对一组数据求聚合操作。8ReduceGroup通过将此数据集中的所有元素传递给函数创建一个新的数据集。该函数可以使用收集器输出零个或多个元素也可以作用于完整的数据集迭代器会返回完整的数据集的元素。9Join将两个DataSet连接生成一个新的DataSet。两个数据集对指定要连接的Key进行Join默认是一个Inner Join。可以使用JoinFunction将该组连接的元素转化为单个元素也可以使用FlatJoinFunction将该组元素转化为任意多个元素包括none​。10OuterJoin两个数据集执行左连接leftOuterJoin​、右连接rightOuterJoin或全外连接fullOuterJoin​。与JoinInner Join的区别在于如果在另一侧没有找到匹配的数据则保存这一侧的记录。11CoGroup对一个或多个字段中的每个输入都进行分组然后加入组。每组都调用转换函数。12Union构建两个数据集的并集。13Cross构建两个输入数据集的笛卡尔积。可选择使用CrossFunction将元素对转换为单个元素。14Rebalance均匀地重新对数据集进行并行分区以消除数据偏差。15Hash Partition根据给定的Key对数据集做Hash分区。可以是Position Keys、Expression Keys或者KeySelector Functions。16Range Partition根据给定的Key对一个数据集进行Range分区。可以是Position Keys、Expression Keys或者Key Selector Functions。17Sort Partition本地按照指定顺序在指定字段上对数据集的所有分区进行排序可以指定Field Position或Filed Expression。** 3数据输出。**DataSet通过Data Sink将DataStream中的数据输出到外部系统。DataSet常用的输出操作如下。1writeAsText​​将元素以字符串形式写入文件。2writeAsCsv…​将元组字段以逗号分隔写入文件可配置行和字段分隔符。3print​​将每个元素都调用toString​打印输出。4write​​自定义文件输出方法的基类支持自定义对象到字节的转换。5output​​通用输出方法。4DataSet API使用实例。以下代码实现一个最基本的DataSet处理流程步骤如下。1定义获取Flink的执行环境ExecutionEnvironment。2从本地文件读取文件到DataSet。3在DataSet上执行转换操作调用filter操作只返回偶数数据。4输出结果通过print方法打印执行结果。2.2.4、Flink Table API和SQLFlink有两种关系型APITable API和SQL。Flink通过Table API和SQL来统一流处理和批处理。Table API是Scala和Java的集成查询API允许以非常直观的方式组合执行关系运算例如Selection、Filter和Join的查询。Flink的SQL基于遵循SQL标准的ApacheCalcite实现。无论批处理DataSet还是流处理DataStream​在两个接口中指定的查询都具有相同语义且具有相同结果。Table API、SQL接口、DataSet、DataStream相互之间紧密集成在一起Flink各个API的关系如图11-27所示。应用可以轻松地在各种类型API库之间切换。例如可以使用CEP库从DataStream中提取数据然后使用TableAPI或SQL执行查询扫描、筛选和聚合批处理最终将结 果输出。1Flink关系型API的原理。Flink基于Apache Calcite实现SQL语义解析。Flink通过利用Calcite的查询优化框架与SQL解释器来完成SQL的解析、查询优化、逻辑树生成得到Calcite的RelRoot类的一棵逻辑执行计划树并最终生成Flink的Table。Table中的执行计划会转化成DataSet或DataStream执行具体的数据计算。Flink关系型API的执行流程如图所示。1Table Source、DataSet、DataStream注册Catalog。2调用Table API或SQL执行SQL查询。3SQL解析、验证并生成统一的Calcite的逻辑执行计划Logical Plan​。4根据数据源的性质DataSet、DataStream使用不同规则进行优化。5最终优化后的执行计划将被转换成常规的FlinkDataSet或DataStream程序。2Flink Table API和SQL的使用流程。Table API和DataSet、DataStream关联密切可以首先通过一个DataSet或DataStream创建出一个Table然后调用类似Filter、Join、Select的关系型转换操作来转换为一个新的Table对象最后将一个Table对象转换为一个DataSet或DataStream并将结果输出。从内部实现上来说所有应用于Table的转换操作都变成一棵逻辑表操作树在Table对象被转换回DataSet或者DataStream之后转换器会将逻辑表操作树转化为对应的DataSet或者DataStream操作符。其使用流程如下。1创建一个TableEnvironmentTableEnvironment对象是Table API和SQL集成的核心主要用于注册一个Table、注册一个外部的Catalog、执行SQL查询、注册一个用户自定义的Function、将DataStream或DataSet转换成Table等操作。2注册Table将一个Table注册到TableEnvironment具体步骤为将一个TableSourceMySQL、HBase、CSV、Kakfa、RabbitMQ等注册到TableEnvironment将一个外部的Catalog注册到TableEnvironment访问外部系统的数据或文件将DataStream或DataSet注册为Table。3查询Table上的数据执行Table API或SQL语句进行查询查询的输出结果是一个新的Table。4输出Table为了将Table输出可以使用TableSink。TableSink是一个通用接口支持各种各样的文件格式例如CSV、Parquet、Avro​也支持各种各样的外部系统例如JDBC、HBase、Cassandra​、ElasticSearch​同样支持各种各样的消息服务例如Kafka、RabbitMQ​。批量数据的导出Table使用BatchTableSink。流数据的导出Table使用AppendStreamTableSink、RetractStreamTableSink和UpsertStreamTableSink。5解析Query并执行在步骤4中输出的Table和SQL查询被解析成DataStream或DataSet。一次查询为一个Logical Query Plan解析Logical Query Plan分为优化Logical Plan、将Logical Plan转化为DataStream或DataSet两步。在Table API和SQL解析完毕后其查询会被当作普通DataStream或DataSet进行执行。如下代码实现了一个基于Table API的批量数据查询操作。上述代码将List中Persion的数据转换为Table并注册Table最后基于SQL执行查询并将计算结果转换为DataSet输出。Table不但能用于批量数据查询还能基于实时流数据执行查询。如下代码从Kafka中获取实时流数据并将数据转换为实时Table执行查询。上述代码将Kafka中的数据转换为StreamTable然后基于该StreamTable执行Table API和SQL的查询操作。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价