【问题标题】:publish to kafka topic within a map method在 map 方法中发布到 kafka 主题
【发布时间】:2018-05-21 18:07:39
【问题描述】:

从映射函数 (SCALA) 中写入 kafka 主题?

  1. FLINK 应用程序中读取 kafka 主题
  2. 在地图函数中处理数据
  3. 问题陈述 - 在 map 函数中,我正在遍历一个列表。对于列表中的每个元素,我想发布到一个 kafka 主题。
  4. 当我从地图获取输出并接收它时,它可以工作,但如果我尝试从地图方法中推送到主题,它不会
  5. 是否可以从 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


【解决方案1】:

不要在 map 函数中创建 kafka 生产者,也不要尝试在 map 中写入 kafka 主题。老实说,我不能引用任何话说这是一个坏主意……但这是一个坏主意。

相反。将您的地图函数更改为 flatMap(请参见此处的第一个示例:https://ci.apache.org/projects/flink/flink-docs-release-1.3/dev/datastream_api.html)。

因此,在您的循环中,您只需执行 collector.collect(recordToPublishToKafka),而不是在每个循环中都创建一个 kafka 生产者。

您的接收器将在收集到它们时将它们发布出来。

【讨论】:

  • 谢谢,有帮助。我同意你的看法(我仍在学习 Flink/Kafka 工具集,因此是个坏主意……)会采纳你的建议。非常感谢和问候
猜你喜欢
  • 1970-01-01
  • 2018-10-28
  • 2020-09-08
  • 2020-11-14
  • 2020-01-18
  • 2018-10-27
  • 2020-06-24
  • 2019-11-08
  • 2016-01-17
相关资源
最近更新 更多