【问题标题】:How to send multiple files to kafka producer using akka stream in Scala如何在Scala中使用akka流向kafka生产者发送多个文件
【发布时间】:2019-02-02 23:36:45
【问题描述】:

我正在尝试使用 akka 流向 kafka 生产者发送多个数据,同时我自己编写了生产者,但正在努力使用 akka-streamIO 来获取多个文件,这些文件将是我想要发送给我的数据kafka Producer 这是我的代码:

 object App {

  def main(args: Array[String]): Unit = {
    val file = Paths.get("233339.8.1231731728115136.1722327129833578.log")
//    val file = Paths.get("example.csv")
//
//    val foreach: Future[IOResult] = FileIO.fromPath(file)
//      .to(Sink.ignore)
//      .run()

    println("Hello from producer")

    implicit val system:ActorSystem = ActorSystem("producer-example")
    implicit val materializer:Materializer = ActorMaterializer()

    val producerSettings = ProducerSettings(system,new StringSerializer,new StringSerializer)

    val done: Future[Done] =
      Source(1 to 955)
        .map(value => new ProducerRecord[String, String]("test-topic", s"$file : $value"))
        .runWith(Producer.plainSink(producerSettings))

    implicit val ec: ExecutionContextExecutor = system.dispatcher
    done onComplete  {
      case Success(_) => println("Done"); system.terminate()
      case Failure(err) => println(err.toString); system.terminate()
    }

  }

}

【问题讨论】:

  • 当您取消注释该代码时,当前发生了什么?这适用于一个文件吗?您尝试过获得多个什么?
  • 它对一个文件有效,但对多个文件无效...这是我发现的: val file = Paths.get("example.csv") val foreach: Future[IOResult] = FileIO.fromPath(file) .to(Sink.ignore) .run()
  • 如果你尝试多个线程会怎样。每个文件一个?
  • 你是什么意思?我只想能够发送与我现在使用的不同的源,并且能够通过 Kafka 生产者发送我的 955 个日志文件,这些文件是文本文件......
  • 对... new FileSenderThread(filename).start()... 将开始从文件创建自己的生产者

标签: scala apache-kafka akka-stream


【解决方案1】:

给定多个文件名:

val fileNames : Iterable[String] = ???

可以创建一个Source,它发出使用flatMapConcat 连接在一起的文件的内容:

val chunkSize = 8192

val chunkSource : Source[ByteString, _] = 
  Source.apply(fileNames)
        .map(fileName => Paths get fileName)
        .flatMapConcat(path => FileIO.fromPath(path, chunkSize))

这将发出固定大小的ByteString 值,这些值都是chunkSize 长度,除了可能更小的最后一个值。

如果你想用一些分隔符分隔行,那么你可以使用Framing:

val delimiter : ByteString = ???

val maxFrameLength : Int = ???

val framingSource : Source[ByteString, _] =
  chunkSource.via(Framing.delimiter(delimiter, maxFrameLength))

【讨论】:

  • 谢谢拉蒙,但我有 955 个日志文件,它们是文本文件。我正在尝试更改我的生产者(kafka)发送的数据的来源,并且我希望能够发送多个文件,而不是像我在代码中所做的那样发送文件的编号和名称。谢谢!
  • @tupacshakur 我的解决方案不发送名称,而是发送实际内容。
  • 是的,我知道 Ramon,我只是想知道如何将您的解决方案集成到我的代码中,以便将数据作为 kafka 生产者的价值发送。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-01-02
  • 2018-04-19
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多