【发布时间】: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