-
对地图函数进行Flink两阶段提交,以实现完全一次的语义
<strong>背景</strong>: 我们有一个Flink管道,该管道由多个源,多个接收器和沿该管道的多个运算符 -
在执行纱线应用程序终止并再次运行之后,flink将从上一个偏移量恢复吗?
我使用FlinkKafkaConsumer来使用kafka并启用检查点。现在,我对偏移量管理和检查点机制有些困惑。 我已经 -
在flink流传输过程中一次读取文件的两行
我想使用flink流处理文件,其中两行属于同一流。在第一行中有一个标题,在第二行中有一个对应的文本 -
重构复杂的业务逻辑以实现链接?
我们有一个非常复杂的现有程序,它的核心业务逻辑类似于: 使用kafka消息(id,...)->检查并验 -
Flink停止作业API返回“ 308永久重定向”
我正在跟踪<a href="https://ci.apache.org/projects/flink/flink-docs-stable/monitoring/rest_api.html#jobs-jobid-1" rel="nofollow noref -
Flink是否具有“ distribute by”命令(类似于spark)
我想通过键对数据进行随机排序并将数据写入接收器。 我不使用分组依据,因为它不需要我进行 -
为什么状态在Apache Flink的coprocessfunction中总是返回null?
我连接了两个流,然后调用<code>process</code>来实现我的逻辑以获得结果。以下是我的Flink代码的流程。 -
如何在Flink SQL中计数
我想在flink SQL中做count(0),但是它给出了类似的异常 org.apache.flink.client.program.ProgramInvocationExcept -
如何在flink程序中动态添加jar
我有一个flink程序,其中一些转换逻辑来自用户定义的lambda函数。我必须在运行时通过下载jar和invoke方法 -
什么环境_config用于Beam启动flink
我希望在运行Beam wordcount.py演示时获得有关如何设置<code>--environment_config</code>的指导。 在DirectRunner -
ClassNotFoundException:org.apache.flink.streaming.connectors.rabbitmq.common.RMQConnectionConfig $ Builder
我正在使用flink 1.9.0和Rabbitmq连接器读取数据,我可以成功编译代码,但是当我运行代码时,出现以下错 -
多节点集群上的Flink与Spark部署模式
在Spark中,我熟悉的三个<strong>集群</strong>(非本地)部署选项: <ul> <li>独立</li> <li> Mesos </li> <li>纱 -
Flink:MaxOutOfOrderness和AllowedLateness之间的区别
在Flink中,有两件事提供相似的行为。两者有什么区别。 <ol> <li> <strong> MaxOutOfOrderness </strong>:与Bound -
使用Apache Flink流处理缓冲转换后的消息(例如1000条计数)
我正在使用Apache Flink进行流处理。 从源(例如:Kafka,AWS Kinesis Data Streams)订阅消息,然后使用Fli -
为什么在本地计算机上可以获得正确的结果,但是当我在Apche Flink中的多个节点上运行时情况却有所不同?
下面是我的Flink逻辑。我所做的是,首先我训练了值,然后创建了自动编码的模型。然后,我测试了以下 -
numLateRecordsDropped:对运营商意味着什么
<a href="https://i.stack.imgur.com/5GklU.png" rel="nofollow noreferrer"><img src="https://i.stack.imgur.com/5GklU.png" alt="enter image -
在Apache Flink Broadcast流中应用基于窗口的规则
我在Apache Flink的BroadcastStream中有一组规则。 我可以在事件流中应用新规则。 但是我无法弄清楚如果我的 -
Flink从卡夫卡加入流
我正在处理来自Kafka的2个流,但是其中一个并不那么频繁。 我尝试执行加入并成功,但是仅几个 -
FlinkML的状态如何?
没有看到有关FlinkML的最新讨论-它快死了还是死了? 最近有趣的一些实时用法有哪些例子? -
对于akka.pattern.AskTimeoutException,Flink作业提交失败
<code>akka.pattern.AskTimeoutException</code>的工作提交失败,我尝试设置<code>akka.ask.timeout: "60s"</code>,但 -
如何将Flink中的时间窗口保存到文本文件?
我开始在Java中使用ApacheFlink。 我的目标是在一分钟的时间窗口中使用一个ApacheKafka主题,这将应用 -
收集后Flink TestHarness输出未清除
我进行了以下测试: <pre><code>testHarness.processElement2(new StreamRecord<>(element1)); testHarness.processElement1( -
Flink TaskManager livenessProbe无法正常工作
我正在遵循<a href="https://ci.apache.org/projects/flink/flink-docs-stable/ops/deployment/kubernetes.html#session-cluster-resource-de -
如何使用Apache Flink在Cassandra中删除行?
在Apache Flink中,可以很容易地通过<code>CassandraSink</code>在Cassandra中插入一行。但是我找不到删除行的方法 -
在相同字段上未更改的Flink键控会引起随机播放吗?
<pre><code>dataStream.map(func1).keyBy("key") //(1) .process(func2).keyBy("key") //(2) .timeWindow().aggregate(func3 -
在Apache Flink的protobuf事件中发布反序列化事件
我正在flink应用中读取运动学的事件。事件采用protobuf格式。如果我在flink应用程序中使用<code>'com.googl -
如何扩大/缩小正在运行的Flink集群?
我现在正在Kubernetes上运行Flink。我假设如果更新TaskManager部署的副本,Kubernetes会为我增加/减少TM吊舱的 -
由于背压导致检查点超时
我们有一项长期的工作,可以将记录存储到Elasticsearch中。由于ES群集在午夜进行路由,因此接收器有时 -
将Apache Flink连接到Elasticsearch的问题
我在Flink站点内使用了一段代码,将Apache Flink连接到Elastic Search。我想通过maven项目从NetBeans软件运行这段 -
StreamingFileSink无法重命名正在进行的文件
当StreamingFileSink尝试重命名正在进行的文件时,我偶尔会遇到问题。参见以下异常: FS已安装在Goo