flink-streaming

    7热度

    1回答

    在我们的项目中,我们有一个Flink(1.1.3)流作业,它从一个kafka队列读取数据,执行映射函数转换并写入另一个队列。 直到我们将流出的REST请求作为流的一部分引入之后,这一切运行良好。 要做到这一点,我们使用了PlayFramework WSClient(因为它是在我们的堆栈的其他地方使用),并以这种方式创造了它的代码: val config = new AhcWSClientConfi

    0热度

    1回答

    Spark流提供终止awaitTermination()的API。有没有类似的API可以在几秒钟后正常关闭flink流?

    0热度

    1回答

    我试图运行一个简单的Apache Flink脚本与卡夫卡指令,但我一直有执行问题。 该脚本应该读取来自kafka制作者的消息,详细阐述它们,然后再发送回另一个主题,即处理结果。 我从这里得到这个例子: http://apache-flink-user-mailing-list-archive.2336050.n4.nabble.com/Simple-Flink-Kafka-Test-td4828.

    0热度

    1回答

    我想用Apache Flink一个接一个地批处理两个文件。 对于一个具体的例子:假设我想给每一行分配一个索引,这样来自第二个文件的行跟随第一行。相反,这样做的,下面的代码行交错的两个文件: val env = ExecutionEnvironment.getExecutionEnvironment val text1 = env.readTextFile("/path/to/file1")

    0热度

    1回答

    这是我第一次使用Apache Flink(1.3.1)并且有问题。更详细地说,我正在使用flink-core,flink-cep和flink-streaming库。我的应用程序是一个Akka ActorSystem,它消耗来自RabbitMQ的消息,并且各个角色处理这些消息。在一些演员中,我想实例化来自Flink的StreamExecutionEnvironment并处理传入的消息。因此我写了一个

    1热度

    1回答

    Apache Flink缓存任务的传出,然后发送下一个任务进行处理。缓冲会影响延迟,因为我知道即使没有填充缓冲区,缓冲也会发送数据到下一个任务。 我该如何改变缓冲超时?我找不到任何文件。 是每个Flink群集或每个TaskManager的配置?它可以根据任务/操作员进行配置吗? 因为我知道Flink缓冲区,即使任务在同一个TaskManager上。在这种情况下,它将影响同一个TaskManager

    0热度

    1回答

    我在Apache Flink的scala中运行一个简单的脚本。 我从Apache Kafka制作人处读取数据。这是我的代码。 import java.util.Properties import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment import org.apache.flink.stream

    0热度

    1回答

    问题陈述: 试图评估的Apache弗林克建模先进的实时低延迟分布分析 使用案例摘要: 的仪器I1提供复杂的分析,I2 ,I3 ...等各有产品定义P1,P2,P3;配置用户参数(动态)U1,U2,U3 &需要流媒体市场数据M1,M2,M3 ... 仪器分析功能(A1,A2)在计算复杂性方面复杂,其中一些可能需要300-400ms但可以并行计算。 从上面清楚地看到,市场数据流将比分析功能&需要消耗最

    1热度

    1回答

    我在使用Flink的Table API和/或Flink的SQL支持(Flink 1.3.1,Scala 2.11)在streaming环境中。我开始用DataStream[Person]和Person是一个案例类,看起来像: Person(name: String, age: Int, attributes: Map[String, String]) 所有工作正常,直到我开始把attribut

    1热度

    1回答

    我想加入来自Kafka生产者的两个流(json)。 如果我筛选数据,代码将起作用。但是当我加入他们时似乎不起作用。我想打印到控制台的联合流,但没有出现。 这是我的代码 import java.util.Properties import org.apache.flink.streaming.api.scala._ import org.apache.flink.streaming.connec