首页 > TAG信息列表 > apache-flink

Java-Flink-多个源的集成测试

我有一个Flink作业,正在使用此处描述的方法进行集成测试:https://ci.apache.org/projects/flink/flink-docs-stable/dev/stream/testing.html#integration-testing 作业从两个来源获取输入,这两个来源合并在CoFlatMapFuntion中.在测试环境中,我当前正在使用两个简单的SourceFunction

java-flink-测量背压

我使用Flink进行了一些测试,以与其他流媒体平台进行比较.测试的数据源是一个kafka主题,其通信量各不相同,Im试图找出flink是否跟上潮流. 有没有办法知道kafka消费者承受了多少“背压” flink? IE是否跟上步伐?解决方法:Apache Kafka项目提供了一些工具,可以从Zookeeper中获取主题和消

java-如何增加Flink taskmanager.numberOfTaskSlots以在没有Flink服务器的情况下运行它(在IDE或胖子中)

我有一个关于在IDE中或作为胖子运行Flink流作业而不将其部署到Flink服务器的问题. 问题是,当我的工作中有多个任务槽时,无法在IDE中运行它. public class StreamingJob { public static void main(String[] args) throws Exception { // set up the streaming execution envi

java-Apache Flink:如何计算DataStream中事件的总数

我有两个原始流,我正在加入这些流,然后我要计算已加入的事件总数是多少,尚未加入的事件有多少.我通过在joinedEventDataStream上使用地图来执行此操作,如下所示 joinedEventDataStream.map(new RichMapFunction<JoinedEvent, Object>() { @Override publ

如何在flink环境下初始化flink作业的spring资源

我最近遇到了一些关于开发flink作业的问题,这些问题引入了spring和hibernate,而且这个工作将在flink集群上运行.所以我需要在任务管理器而不是作业管理器上运行flink操作符之前初始化spring资源.但是我找不到任何合适的StreamExecutionEnvironment方法来做到这一点. 我尝试了下面这

java – Cassandra Pojo Sink Flink中的动态表名

我是Apache Flink的新手.我正在使用Pojo Sink将数据加载到Cassandra中.现在,我在@Table注释的帮助下指定表和键空间名称. 现在,我想在运行时动态传递表名和键空间名,以便我可以将数据加载到用户指定的表中.有没有办法实现这个目标?解决方法:@Table是一个CQL注释,用于定义此类实体映

java – Apache Flink:由TupleSerializer引起的NullPointerException

当我执行我的Flink应用程序时,它给我这个NullPointerException: 2017-08-08 13:21:57,690 INFO com.datastax.driver.core.Cluster - New Cassandra host /127.0.0.1:9042 added 2017-08-08 13:22:02,427 INFO org.apache.flink.runtime.taskmanager.Task

java – JpaRepository不在自定义RichSinkFunction中自动装配

我创建了一个自定义的Flink RichSinkFunction并试图在这个自定义类中自动装配JpaRepository,但我不断得到一个NullPointerException. 如果我在构造函数中自动装配它,我可以看到找到了JpaRepo – 但是当调用invoke方法时,我收到一个NullPointerException. public interface Messag

java – flink:在窗口流上应用多个聚合

我有一些数据以id,float,float,float形式出现.我想按顺序min(),max()和sum()字段,并按id值分组. 使用flatMap我有一个Tuple4的位,但我不知道如何将它发送到下一步. 是)我有的: dataStream.flatMap(new mapper()).keyBy(0) .timeWindowAll(Time.of(5, TimeUnit.SECONDS)).min(1

java – flink:scala版本冲突?

我试图在IntelliJ中编译here的kafka样本.经过很多与依赖关系的讨论已经遇到了这个问题,我似乎无法过去: 15/10/25 12:36:34 ERROR actor.ActorSystemImpl: Uncaught fatal error from thread [flink-akka.actor.default-dispatcher-4] shutting down ActorSystem [flink] java.lang