【发布时间】:2018-02-27 23:47:17
【问题描述】:
我有一个 kafka 主题和一个听它的 KTable。
我想写一个 http POST 请求,它会遍历 ktable 中的当前项目,对它们执行一些操作并写回主题
所以基本上我有:
private val accessTokenTable: KTable[String, String] = builder.table(token_topic_name, tokenStoreString)
val stream: KafkaStreams = new KafkaStreams(builder, streamingConfig)
stream.cleanUp()
stream.start()
....
override def refreshTokens = {
accessTokenTable.mapValues {
new ValueMapper[String, String] {
override def apply(value: String) = {
value
}
}
}.print(token_topic_name)
}
当我尝试调用此方法时,不会打印/写入主题
我错过了什么?我唯一的选择是将ktable中的消息写入hashmap并从那里读取吗?它错过了 ktables 的全部意义?
【问题讨论】:
标签: scala apache-kafka apache-kafka-streams