当桌子可以使用相同的标识符对桌子进行分区时,如何组织flink中的数据流
我坚信Flink是解决事件处理问题的完美解决方案。我什至设法制作了一个工作原型,但我不相信它几乎是最佳的。 这是我的场景: 我有两个运动流 一个流包…
Flink getResult()的聚合功能可以更改累加器值?
getResult()函数可以更改累加器的值吗?在累积时仍会有效应值。 代码在下面: public class WeightedAverage implements AggregateFunction { publi…
Flink Oracle连接
我正在使用支持 Flink 1.13 的AWS Kinesis Studio。我看到Flink 1.13不支持Oracle连接。 基于版本1.13的文档,它支持 mysql,Postgresql,derby。 我…
在flink SQL或表API中控制每个操作员并行性
我正在使用Flink SQL来定义从无限源(Google pubsub)读取的流媒体作业,并使用HBase接收器写入Google Bigtable。查询是在流中窗口窗口,并在这些窗口…
Flink缺少窗户处理器(事件时间窗口)和Kafka源
我们有一个流媒体作业,具有20个单独的管道,每个管道都有一个/许多Kafka主题来源,并且一些管道带有窗户处理器,而其他管道则是未窗口的处理器。 当…
如何使用Flink中的单个窗口进行多个汇总?
我是Flink的新手,我想做很多次在Spark中做的事情。 例如,在Spark中,我可以在 ds.groupByKey(???).mapGroups(???) // aggregate 1 .groupByKey(???)…
FLINK-如何解决过滤器功能中的背压?
使用Flink v1.13.2使用RockSDB,我正在用Grafana中的以下命令来测量任务按: sum by(task_name)(flink_task_taskmanager_job_job_job_task_back_ba…
有人面对这个``时间戳始终都应该是非负或无效的。
从Flink应用程序向Kafka主题发布消息时获取此错误。 相同的代码在我们的测试环境中使用相似的生产者配置,但在生产环境中失败。我找不到这个问题的原…
Adpative调度程序未通过Flink 1.14.0识别
我正在尝试使用Flink 1.14的自适应时间表来基于可用资源来运行Flink作业,而不是等待所需的并行性(缩放),但我看不到Flink识别自适应时间表。 Ex: f…
Flink流根据条件进入两个接收器
试图根据条件看到流到两个水槽的可能性。 Requirement is stream have events, all events after transformation need to go to one sink ( assume on…
datastream/table API的统一连接器
我正在为 Flink 1.14.4 编写一个简单的连接器(源/接收器),它主要包装官方的 Kafka 连接器并自动设置自定义序列化器/反序列化器。我对 FLIP-27 中引…
Flink Task Manager如何预先设置几个设置
我需要下载密钥库并在创建源连接器时使用并使用它,目前,我正在覆盖方法但是,这被我的每个源连接器都要求。 @Override public void open(Configurat…