Я новичок в apache моргании. У меня есть один проект flink scala, который использует данные из кластера kafka, и мне нужно передать результат потока в качестве параметра, чтобы использовать 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()
}
}
Любая помощь?