【发布时间】:2018-05-21 18:07:39
【问题描述】:
从映射函数 (SCALA) 中写入 kafka 主题?
- 在 FLINK 应用程序中读取 kafka 主题
- 在地图函数中处理数据
- 问题陈述 - 在 map 函数中,我正在遍历一个列表。对于列表中的每个元素,我想发布到一个 kafka 主题。
- 当我从地图获取输出并接收它时,它可以工作,但如果我尝试从地图方法中推送到主题,它不会
-
是否可以从 map 方法中发布到主题
// Main Function def main(args: Array[String]) { ... // some list val list_ = ("a", "b", "c", "d") // Setup Properties val props = new Properties() props.setProperty("zookeeper.connect", zookeeper_url + ":" + zookeeper_port) props.setProperty("bootstrap.servers", broker_url + ":" + broker_port) props.setProperty("auto.offset.reset", "earliest") props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer") props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer") ... // Connect to Source val input_stream = env.addSource(new FlinkKafkaConsumer09[String](topic_in, new SimpleStringSchema(), properties)) // Process each Record val stream = input_stream.map(x=> { // loop through list "list_" -> variable in in Main // and publish to topic_out // -- THIS IS MY CURRENT ISSUE !!!) // -- Does not work (No compile issue) // var producer2 = new KafkaProducer[String, String](props) var record = new ProducerRecord(topic_out, "KEY", list(i)) producer2.send(record) producer2.flush() // ... Other process and return processed string }) // publish to different topic of proccessed input string (Works) stream.addSink(new FlinkKafkaProducer09[String](broker_url + ":" + broker_port, other_topic, new SimpleStringSchema()))
【问题讨论】:
-
什么“不起作用”?你为什么要在循环中创建生产者???所有
vars是怎么回事?什么是i????为什么在map内外使用不同的生产者类? -
“不起作用” - 当我将作业提交给 flink 时。它运行工作,但它没有进入地图功能(如果包含生产者代码)“
-
我正在尝试根据输入流中的某个条件循环遍历列表并将感兴趣的项目发布到主题
-
我该怎么做呢? (感谢您的帮助)
标签: scala apache-kafka apache-flink kafka-producer-api