【发布时间】:2020-11-02 15:42:30
【问题描述】:
我正在尝试使用 Consumer.committableSource 通过 Akka 从 Kafka 读取数据。然后我想将数据写入共享文件夹的文件中。
提交时,我们通常使用 via(Committer.flow(committerSettings) 之类的东西。
但是,这个方法不返回Kafka流的值,所以之后我不能调用.runWith(LogRotatorSink.withSinkFactory(rotator, sink))之类的东西来写入数据。
这是没有提交的代码:
Consumer.committableSource(settings, Subscriptions.topics(kafkaTopics.toSet))
.via(processor)
.prepend(headerCSVSource)
.via(CsvFormatting.format(delimiter =
CsvFormatting.SemiColon))
.runWith(LogRotatorSink.withSinkFactory(rotator, sink))
这是我认为我需要的:
Consumer
.committableSource(settings, Subscriptions.topics(kafkaTopics.toSet))
.via(processor)
.prepend(headerCSVSource)
.via(CsvFormatting.format(delimiter =
CsvFormatting.SemiColon))
.via(Committer.flow(committerSettings))
.runWith(LogRotatorSink.withSinkFactory(rotator, sink))
但这不起作用,因为via(Committer.flow) 不返回流值(但 Flow[Committable, Done, NotUsed])。
我需要的是仅在将数据写入文件后才提交偏移量。 如果您觉得其他选项(例如使用 plainSource / auto-commit)更合适,我愿意考虑。
【问题讨论】: