【问题标题】:How to push datastream to kafka topic by retaining the order using string method in Flink Kafka Problem如何通过在 Flink Kafka 问题中使用字符串方法保留顺序来将数据流推送到 kafka 主题
【发布时间】:2022-01-26 00:43:43
【问题描述】:

我正在尝试每个500 ms 创建一个JSON 数据集,并希望将其推送到Kafka 主题,以便我可以在下游设置一些窗口并执行计算。以下是我的代码:

package KafkaAsSource

import org.apache.flink.api.common.serialization.SimpleStringSchema
import org.apache.flink.streaming.api.datastream.DataStream
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer.Semantic
import org.apache.flink.streaming.connectors.kafka.{FlinkKafkaProducer}
import org.apache.flink.streaming.connectors.kafka.internals.KeyedSerializationSchemaWrapper


import java.time.format.DateTimeFormatter
import java.time.LocalDateTime
import java.util.{Optional, Properties}

object PushingDataToKafka {

  def main(args: Array[String]): Unit = {
    val env = StreamExecutionEnvironment.getExecutionEnvironment
    env.setMaxParallelism(256)
    env.enableCheckpointing(5000)
    val stream: DataStream[String] = env.fromElements(createData())

    stream.addSink(sendToTopic(stream))
  }

  def getProperties(): Properties = {
    val properties = new Properties()
    properties.setProperty("bootstrap.servers", "localhost:9092")
    properties.setProperty("zookeeper.connect", "localhost:2181")

    return properties
  }

  def createData(): String = {
    val minRange: Int = 0
    val maxRange: Int = 1000
    var jsonData = ""
    for (a <- minRange to maxRange) {
      jsonData = "{\n  \"id\":\"" + a + "\",\n  \"Category\":\"Flink\",\n  \"eventTime\":\"" + DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss.SSS").format(LocalDateTime.now) + "\"\n  \n}"
      println(jsonData)
      Thread.sleep(500)
    }
    return jsonData
  }

  def sendToTopic(): Properties = {
    val producer = new FlinkKafkaProducer[String](
      "topic"
      ,
      new KeyedSerializationSchemaWrapper[String](new SimpleStringSchema())
      ,
      getProperties(),
      FlinkKafkaProducer.Semantic.EXACTLY_ONCE
    )
    return producer
  }
}

它给了我以下错误:

type mismatch;
 found   : Any
 required: org.apache.flink.streaming.api.functions.sink.SinkFunction[String]
    stream.addSink(sendToTopic())

修改后的代码:

object FlinkTest {

  def main(ars: Array[String]): Unit = {
    val env = StreamExecutionEnvironment.getExecutionEnvironment()
    env.setMaxParallelism(256)
    var stream = env.fromElements("")
    //env.enableCheckpointing(5000)
    //val stream: DataStream[String] = env.fromElements("hey mc", "1")

    val myProducer = new FlinkKafkaProducer[String](
      "maddy", // target topic
      new KeyedSerializationSchemaWrapper[String](new SimpleStringSchema()), // serialization schema
      getProperties(), // producer config
      FlinkKafkaProducer.Semantic.EXACTLY_ONCE)
    val minRange: Int = 0
    val maxRange: Int = 10
    var jsonData = ""
    for (a <- minRange to maxRange) {
      jsonData = "{\n  \"id\":\"" + a + "\",\n  \"Category\":\"Flink\",\n  \"eventTime\":\"" + DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss.SSS").format(LocalDateTime.now) + "\"\n  \n}"
      println(a)
      Thread.sleep(500)
      stream = env.fromElements(jsonData)
      println(jsonData)
      stream.addSink(myProducer)
    }

    env.execute("hey")
  }

  def getProperties(): Properties = {
    val properties = new Properties()
    properties.setProperty("bootstrap.servers", "localhost:9092")
    properties.setProperty("zookeeper.connect", "localhost:2181")
    return properties
  }
  /*
  def createData(): String = {
    val minRange: Int = 0
    val maxRange: Int = 10
    var jsonData = ""
    for (a <- minRange to maxRange) {
      jsonData = "{\n  \"id\":\"" + a + "\",\n  \"Category\":\"Flink\",\n  \"eventTime\":\"" + DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss.SSS").format(LocalDateTime.now) + "\"\n  \n}"
      Thread.sleep(500)
    }
    return jsonData
  }
  */

}

Modified Code 给了我 Kafka 主题中的数据,但它不保留顺序。我在循环中做错了什么?另外,必须将 Flink 的版本从 1.13.5 更改为 1.12.2

我最初使用Flink 1.13.5ConnectorsScala2.11。我到底错过了什么?

【问题讨论】:

  • 我也尝试了最新的 flink 版本,即1.14.2,但仍然没有运气。我正在尝试使用KafkaSink 方法,但无法在Scala 中使用它。
  • 我认为您在 main 方法中缺少 env.execute("jobName") 语句。
  • 您的错误非常明确:您的 sendToTopic() 方法返回类型是 Any,而它应该类似于 FlinkKafkaProducer[String]
  • @Niko 感谢您的回复。在发布问题之前,我已经尝试了一切,包括您推荐的建议。我想我仍然缺少一些东西。我也按照文档中的建议进行了尝试,但没有帮助。
  • 我认为从返回String的方法创建数据流时存在问题。

标签: apache-kafka apache-flink


【解决方案1】:

关于这个循环的几点说明:

for (a <- minRange to maxRange) {
    jsonData = 
      "{\n  \"id\":\"" + a + "\",\n  \"Category\":\"Flink\",\n  \"eventTime\":\""
      + DateTimeFormatter
        .ofPattern("yyyy-MM-dd HH:mm:ss.SSS")
        .format(LocalDateTime.now) + "\"\n  \n}"
    println(a)
    Thread.sleep(500)
    stream = env.fromElements(jsonData)
    println(jsonData)
    stream.addSink(myProducer)
}
  • 睡眠发生在 Flink 客户端,仅影响客户端在将作业图提交到集群之前组装作业图的时间。它对作业的运行方式没有影响。

  • 这个循环创建了 10 个独立的管道,这些管道将独立、并行地运行,全部生产到同一个 Kafka 主题。这些管道将相互竞争。


要获得您正在寻找的行为(跨单个管道的全局排序),您需要从单个源(当然是按顺序)生成所有事件,并以并行方式运行作业之一。像这样的事情会做到这一点:

import org.apache.flink.streaming.api.scala.{StreamExecutionEnvironment, _}

object FlinkTest {

  def main(ars: Array[String]): Unit = {
    val env = StreamExecutionEnvironment.getExecutionEnvironment()
    env.setParallelism(1)

    val myProducer = ...
    val jsonData = (i: Long) => ...

    env.fromSequence(0, 9)
      .map(i => jsonData(i))
      .addSink(myProducer)

      env.execute()
  }
}

您可以将 maxParallelism 保留为 256(或其默认值 128);它在这里不是特别相关。 maxParallelism 是keyBy 将密钥散列到的散列桶的数量,它定义了作业可扩展性的上限。

【讨论】:

  • 感谢您的解释,大卫。如何使其在单个管道下连续?实现我需要的最佳方法是什么?我究竟应该改变什么?
  • 当我尝试从shell读取数据时,如何在消费数据时保持顺序?
  • 是不是因为我的代码中设置的maxParallelism 方法正在创建不同的管道?
  • 我尝试使用并行设置1,但它给了我不同的错误。所以,是的,我从概念上理解了你的观点。但是,当我尝试编码时,它给了我错误。 ` val jsonData = (i: Long) => "{\n \"id\":\"" + i + "\",\n \"Category\":\"Flink\",\n \"eventTime \":\"" + DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss.SSS").format(LocalDateTime.now) + "\"\n \n}" env.fromSequence(0, 9 ) .map(i => jsonData(i)) .addSink(myProducer)`
  • 这是 Flink 流。在写入 Kafka 时,实现完全一次结果的唯一方法是使用具有这种效果的 Kafka 事务。但是您可以减少检查点间隔,以使事务更小并更频繁地提交。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-12-19
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-11-04
  • 1970-01-01
相关资源
最近更新 更多