And*_*ose 5 apache-kafka-streams
如果我创建并启动一个 KafkaStream 实例,然后在关闭钩子调用 .close() 中,则不会发生任何事情 - 日志记录表明存在的一个 StreamThread 已进入 PENDING_SHUTDOWN 状态,但随后就永远保持这种状态。我已经尝试了所有方法,但没有运气 - 已阅读 Kafka 流源代码以了解它在关闭期间做了什么,但在我看来,代码意味着 StreamThread 如果处于运行状态,将永远不会停止(这是不可能的) - 我没有在 Kafka JIRA 中看到任何此类性质的错误。
这是我的简单 KafkaStream (scala) 应用程序的相关代码:
val props: Properties = {
val p = new Properties()
p.put(StreamsConfig.APPLICATION_ID_CONFIG, "test-application")
p.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092")
p.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, "1")
p
}
implicit val produced = Produced.`with`(new StringSerde(), new StringSerde())
val builder: StreamsBuilder = new StreamsBuilder()
val in: KStream[String, String] = builder.stream[String, String]("input-topic")
in.map((k,v) =>{
println("Consumed and transforming value")
(k,s"$v_transformed")
}).to("output-topic")
val streams: KafkaStreams = new KafkaStreams(builder.build(), props)
streams.start()
sys.addShutdownHook(streams.stop())
Run Code Online (Sandbox Code Playgroud)
正如您所看到的,它只是从一个主题中读取,对所使用的记录中的值进行微小的更改并将其写入另一个主题。
流启动,当向其发送 SIGINT 时,应用程序将调用 stop()。
当我 ^C 终止该进程时,我看到 Kafka 日志记录表明 StreamThread-1 正在转换为 PENDING_SHUTDOWN,这就是它所能达到的程度。它最终应该(在几秒钟内)达到 NOT_RUNNING 状态,但它从来没有这样做,并且它继续为它继续从输入主题读取的每条记录输出我的 println 语句。
我在这里做错了什么?
更新:根据评论者的建议,我尝试以 60 秒的超时调用 close(),最终得到了这个,但仍然没有关闭:-
[shutdownHook1] INFO o.apache.kafka.streams.KafkaStreams - 流客户端 [-] 流客户端无法在超时内完全停止
| 归档时间: |
|
| 查看次数: |
1330 次 |
| 最近记录: |