【问题标题】:Kafka commit with Akka and LogRotatorKafka 使用 Akka 和 LogRotator 提交
【发布时间】: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)更合适,我愿意考虑。

【问题讨论】:

    标签: apache-kafka akka-stream


    【解决方案1】:

    看起来您需要将流元素传递给一个接收器,并在成功后传递给另一个。

    您可以在流中运行子流。类似这样的事情:

    .via(CsvFormatting.format(delimiter = CsvFormatting.SemiColon))
    .mapAsync(1) { c =>
      Source.single(c).runWith(LogRotatorSink.withSinkFactory(rotator, sink)).map(_ => c)
    }
    .runWith(Committer.sink(committerSettings))
    

    它应该可以工作,但是,经过一番思考,我认为最好不要使用 sink 写入日志,而是使用其他不会终止流的方式。

    【讨论】:

    • 感谢您的回答。我想知道,在您的解决方案中,写入文件失败是否会提交偏移量?
    • 不,不会的。那么你在mapAsync 中得到的是failed Future,它使流失败。
    猜你喜欢
    • 2020-02-15
    • 1970-01-01
    • 2018-02-22
    • 1970-01-01
    • 2020-05-08
    • 2023-01-31
    • 2018-01-19
    • 2021-11-15
    • 2017-09-10
    相关资源
    最近更新 更多