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