Apache Flink是一个用于可扩展批处理和流数据处理的开源平台。 Flink在一个系统中支持批量和流分析。分析程序可以用Java和Scala中简洁优雅的API编写。
Apache Flink:使用TableFunction的LEFT JOIN不会返回预期的结果
Flink版本:1.3.1我创建了两个表,一个是来自内存,另一个是来自UDTF。当我测试join并离开join时,他们返回了相同的结果。我所期待的是离开加入的行多于......
Apache Flink表1.4:表上的外部SQL执行可能吗?
是否可以在外部查询现有的StreamTable,而无需上传.jar获取执行环境并检索表环境?我等了Apache Flink Table 1.4发布,......
我无法在网上找到有关此内容的更多信息。我想知道是否有可能构建一个Flink应用程序,可以动态消耗匹配正则表达式模式的所有主题并同步这些...
Flink以什么格式保持运营商的管理状态(用于检查点或逻辑运算符之间的通信(即沿着作业图的边缘)?文档读取......
Flink抛出java.io.NotSerializableException
我做了自定义KeyedDeserializationSchema来反序列化kafka消息并使用它如下所示:object Job {case class KafkaMsg [K,V](key:K,value:V,topic:String,partiton:Int,offset:...
从Flink文档中,我知道可以使用迭代运算符实现循环。由于Flink代码被延迟评估,因此无法使用while循环评估终止条件。 ...
我正在运行一个基于EventTime测试窗口的简单示例。我可以使用处理时间生成输出,但是当我使用EventTime时,没有输出。请帮我理解我...
我正在做一个Flink项目。该项目的主要思想是读取JSON的数据流(网络日志),关联它们,并生成一个新的JSON,它是不同JSON的组合......
我有一个类在我的Flink流作业中扩展了RichFlatmapFunction。我在open()方法中创建一个Jedis实例并在close()方法中关闭它(jedis.close()),以便所有记录......
Apache Flink:如何将源中的模式应用到另一个数据流?
我有一个事件的数据流,以及另一个模式的数据流。用户在运行时提供模式,他们需要通过Kafka主题。我需要在...上应用每个模式
使用WindowStream.apply()函数无法应用WindowFunction
我使用Apache Flink和Scala相对较新,我只是掌握了一些基本功能。我试图实现自定义WindowFunction。问题是 ...
数据流被分区并分发到每个插槽以进行处理。现在我可以得到每个分区任务的结果。将某些函数应用于...的结果的最佳方法是什么?
Flink:Clustermanager失去的Cluster Execution错误
我在Flink上运行实时流媒体程序,有1名主人和2名工人。一个工作程序在单独的计算机上运行, 而另一个工作程序在主计算机上运行。我是 ...
寻找有关存储/访问Flink参考数据的位置的一些建议。这里的用例非常简单 - 我有一个包含国家列表的列文本文件。我正在播放Twitter数据,然后......
Flink CEP如何管理间歇性状态?它存放在哪里?它只是在内存中还是有支持状态的快速持久存储?文档没有提到这个......
我正在尝试使用带有Flink的EMR将一些输出写入S3。我使用的是Scala 2.11.7,Flink 1.3.2和EMR 5.11。但是,我收到以下错误:java.lang.NoSuchMethodError:org.apache.hadoop ....
我为Apache Flink流程实现了MapFunction。它正在解析传入的元素并将它们转换为其他格式,但有时会出现错误(即传入的数据无效)。我看到两个......
我已经设置了一个2节点的独立Apache Flink集群。对于少量数据(70 MB),2的并行性需要更多的时间(2分30秒)来处理,因为1的并行性只需要18 ...
Flink的CoProcessFunction不会触发onTimer
我尝试聚合两个流,如val joinedStream = finishResultStream.keyBy(_。searchId).connect(startResultStream.keyBy(_。searchId))。process(new SomeCoProcessFunction)然后......
我想从flink包Toletum.pruebas中读一个kafka主题; import org.apache.flink.api.common.functions.MapFunction; import org.apache.flink.api.java.utils.ParameterTool; ...