【发布时间】:2018-04-05 20:38:43
【问题描述】:
我希望通过 GlobalKTable 加入我的一个流,但在此过程中遇到了问题。我也正在听 3 个主题。一个更新原始主题、一个更新主题和一个会话主题。
我的 update-raw 主题流将未序列化的更新请求转换为序列化请求([String, String] -> [String,Update])。这比推送到我的更新主题。
val updateTransformStore = Stores.inMemoryKeyValueStore("updateTransformState")
val updateTransformStoreBuilder = Stores.keyValueStoreBuilder(updateTransformStore,stringSerde,updateSerde)
builder.addStateStore(updateTransformStoreBuilder)
val rawUpdateStream = builder.stream("update-raw", Consumed.`with`(Serdes.String,Serdes.String))
.filter((_,value) => filterUpdateRequest(value))
rawUpdateStream.transform(updateTransformer,"updateTransformState")
.to("update-stream",Produced.`with`(stringSerde,updateSerde))
我的更新主题是 GlobalKTable。
val globalMaterialized: Materialized[String,UpdateInfo,KeyValueStore[Bytes, Array[Byte]]] = Materialized.as("global-update-store").withKeySerde(stringSerde).withValueSerde(updateSerde)
val updateTable: GlobalKTable[String,UpdateInfo] = builder.globalTable("update-stream", globalMaterialized)
我的最后一个主题是我的会话主题
val kvMapper = new KeyValueMapper[String, String, String] {
override def apply(key: String, value: String): String = {
val jsonVal = new ObjectMapper().readTree(value).asInstanceOf[ObjectNode]
val session = jsonVal.findValue("session").findValue("id").toString replaceAll( "\"", "")
session
}
}
val vJoiner = new ValueJoiner[String,UpdateInfo,JsonWithUpdateInfo] {
override def apply(value1: String, value2: UpdateInfo): JsonWithUpdateInfo = {
if(value2 == null) {
JsonWithUpdateInfo(value1,0,"default")
}
else {
JsonWithUpdateInfo(value1, value2.info1, value2.info2)
}
}
}
val filteredStream = builder.stream("session", Consumed.`with`(Serdes.String, Serdes.String))
.filter((_, value) => filterRequest(value))
val joinedStream:KStream[String,JsonWithUpdateInfo] = filteredStream.join(updateTable,kvMapper,vJoiner)
joinedStream.print(stringSerde,jsonSerde)
在将一个更新和一个会话推送到各自的主题后,我的程序似乎在filteredStream.join 尝试执行后停止执行。我的打印永远不会工作,也不会运行任何运行进一步命令的尝试。在我的 kvMapper 和 vJoiner 中抛出调试打印也不会产生输出。尝试加入此表和流时,我有什么遗漏吗?
谢谢
【问题讨论】:
标签: scala apache-kafka-streams