flink-streaming 相关问题

Apache Flink是一个用于可扩展批处理和流数据处理的开源平台。 Flink在一个系统中支持批量和流分析。分析程序可以用Java和Scala中简洁优雅的API编写。

动态调用Flink运算符

我最近开始学习流处理,现在正在Apache Flink尝试。我正在尝试编写一项工作,以读取来自Kafka主题的事件,可能会执行一些无状态链接...

回答 2 投票 0

生成“心跳”类型事件以将事件时间向前推送

我正在构建Flink流应用程序,并希望使用事件时间,因为它可以确保所有设置的计时器都将在发生故障或重播时确定性地触发...

回答 1 投票 0

Flink Kafka:在没有收到消息的时间间隔后,正常关闭来自kafka源的消耗flink的消息

我已将flinkkafkaconsumer作为源添加到我的流执行环境。我想在特定时间内未收到新消息时关闭/停止flink消耗数据(类似于kafka ...

回答 1 投票 0

如何保持Flink运动消费者的幂等性?

我有一个用例,其中,我通过EMR上运行的flink作业(使用flink-kinesis连接器)使用了运动学流中的事件。作业接收事件,对其进行处理并将其下沉到某个数据存储中。...

回答 1 投票 1

keyBy求和两次的问题是什么

下面是我写的简单代码:val env = StreamExecutionEnvironment.getExecutionEnvironment val list = new ListBuffer [Tuple3 [String,Int,Int]] val random = new Random()for(x

回答 1 投票 0

Flink fireTime on windowAll运算符

我有一个带有两个键的管道,因此我正在使用WindowAll运算符。问题在于windowAll运算符处于一种状态,需要及时进行修改,并且没有...

回答 1 投票 0

如何通过另一个流的修订记录来更改flink中处理器的状态?

边际建议可以实现我的目标。但是它仍然没有计划。目前还有其他方法可以实现吗?

回答 1 投票 0

flink进程的功能

我正在研究Flink,并且正在执行一些程序。我看到有很多与Flink相关的过程。当我启动集群时,它将启动作业管理器和任务管理器的过程,并且...

回答 1 投票 1

哪个是读取慢速更改查找并丰富流输入集合的最佳方法?

我正在使用Apache Beam,具有1.5GB的流媒体集合。我的查询表是JDBCio mysql响应。当我在没有侧面输入的情况下运行管道时,我的工作将在大约2分钟内完成。当我...

回答 1 投票 0

在Flink中以(X,Y)为键的缓慢流丰富了由(X)为键的缓慢变化的流

我需要用(userId)键控的缓慢变化的streamB来充实以(userId,startTripTimestamp)键控的快速变化的streamA。我将Flink 1.8与DataStream API一起使用。我考虑2种方法:...

回答 1 投票 1

Flink中的垃圾收集器

我正在Flink中查看堆内存图,我发现堆内存的值始终在增长。当GC自身激活时?是否有用于处理GC的Flink类?

回答 2 投票 1

从flink群集外部访问flink状态的方法是什么?

我是Apache flink的新手,正在构建一个简单的应用程序,在该应用程序中,我从运动学流中读取事件,例如TestEvent {String id,DateTime created_at,Long amount} ...

回答 1 投票 0

Flink Yarn在任务失败时无限重启

[我在具有以下配置的AWS yarn群集上运行flink流作业,主节点-1,核心节点-1,任务节点-3并且启用了jobmanager.execution.failover-strategy:region作为...之一]]] >

回答 1 投票 1

Flink Yarn在任务失败时无限重启

[我在具有以下配置的AWS yarn集群上运行flink流作业,主节点-1,核心节点-1,任务节点-3并且启用了jobmanager.execution.failover-strategy:region作为...之一]]] >

回答 1 投票 1

Apache Flink TaskExecutor关闭

我是Apache Flink的新手,所以我目前正在尝试做一些实验。我正在从Kafka阅读主题,然后将其打印在控制台上。打印大约100k + kafka消息后,它抛出...

回答 1 投票 2

在Flink中的两个不同流中使用相同的运算符

我想在两个不同的流中使用相同的运算符。但是,我收到一个错误,指出UID或该运算符不是唯一的。惰性val opt:DataStream [Foo] => DataStream [Buzz] = src => src ....

回答 1 投票 0

减少参数的太多参数[Scala中的Flink 1.9]

我正在尝试将Flink的增量窗口聚合与ReduceFunction一起用于我正在做的项目,以返回单个值,该值是具有窗口边界的时间窗口中的最小值。 def ...

回答 2 投票 0

Apache Flink广播流中基于应用窗口的规则

我在Apache Flink的BroadcastStream中有一组规则。我可以在事件流中应用新规则。但是我无法弄清楚如果我的规则像...

回答 1 投票 0

数据流在Flink转换中重复使用

我有一个DataStream(比如说inStream),我需要在其上应用两个不同的转换列表以生成两个不同的输出流(比如说outStream1和outStream2)。 inStream是...

回答 1 投票 1

Flink:MaxOutOfOrderness与AllowedLateness之间的差异

在Flink中,有两件事提供相似的行为。两者有什么区别。 MaxOutOfOrderness:与BoundedOutOfOrdernessTimestampExtractor一起使用。允许流的元素...

回答 1 投票 0

© www.soinside.com 2019 - 2024. All rights reserved.