【问题标题】:Alpakka Akka Stream unable to read from kafkaAlpakka Akka Stream 无法从 kafka 读取
【发布时间】:2019-06-06 22:17:04
【问题描述】:

我已经基于 alpakka 项目构建了一个非常简单的 akka 流,但它不会从 kafka 读取任何内容,即使它连接并创建了一个消费者组。我为流创建了一个隐式 Actor System 和 Materializer。

val done = Consumer.committableSource(consumerSettings,
Subscriptions.topics(kafkaTopic))
.map(msg => msg.committableOffset)
.mapAsync(1) { offset =>
offset.commitScaladsl()
}
.runWith(Sink.ignore)
  • [stream.actor.dispatcher] 将此消息发送给 KafkaConsumerActor “请求消息,requestId:1,分区:Set(kafka-topic-0)”
  • KafkaConsumerActor 似乎没有收到消息,但是当主管要求 Actor 关闭时,它确实收到消息并关闭。

关于为什么它无法读取 Kafka 而没有错误或异常的任何线索?

【问题讨论】:

  • 在启用更多日志时,我发现 - 调试] [06/06/2019 18:25:29.655] [akka-kafka-akka.kafka.default-dispatcher-15] [akka:/ /akka-kafka-poc/system/kafka-consumer-1] 收到来自 Actor[akka://akka-kafka-poc/deadLetters] 的已处理消息 Poll(akka.kafka.internal.KafkaConsumerActor@20eb829b,true)消费者正在接收来自演员 deadLetters 的消息。不确定是 kafka-consumer 演员被杀死还是向 kafka-consumer 发送投票的调度员已经死亡。

标签: apache-kafka akka akka-stream


【解决方案1】:

我不明白为什么我的 akka 流没有使用来自 kafka 代理的消息,但是当我实现与 Runnable Graph 相同的流时,它就可以工作了。

我使用的示例 - https://www.programcreek.com/scala/akka.stream.scaladsl.RunnableGraph

【讨论】:

    猜你喜欢
    • 2019-03-16
    • 2020-05-19
    • 2018-01-19
    • 1970-01-01
    • 2021-07-09
    • 1970-01-01
    • 1970-01-01
    • 2019-09-02
    • 2019-05-12
    相关资源
    最近更新 更多