我是 apache flink 的新手。我有一个使用来自 kafka 集群的数据的 flink scala 项目,我需要将流结果作为参数传递以使用返回转换后的流的 api。这是我的代码
class Testing {
def main(args: Array[String]): Unit = {}
def streamTest(): Unit = {
val env: StreamExecutionEnvironment = StreamExecutionEnvironment.getExecutionEnvironment
val properties = new Properties()
properties.setProperty("bootstrap.servers", "test1.server.local:9092,test2.server.local:9092,test3.server.local:9092")
val consumer_test = new FlinkKafkaConsumer[String]("topic_test", new SimpleStringSchema(), properties)
consumer_test.setStartFromEarliest()
val stream = env.addSource(consumer_test).setParallelism(5)
val api_test = "http://api-test.server.local/test/?msg=%s"
// Here I need pass stream as parameter to api and return transformed stream
env.execute()
}
}
任何帮助?