【问题标题】:How to join KStream and GlobalKTable?如何加入 KStream 和 GlobalKTable?
【发布时间】: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


    【解决方案1】:

    (来自上面评论部分的对话)

    如果您有任何 null 键或值,它们将被删除。

    为了跟踪调用链,您可能需要在KStreamKTableJoinProcessor#process() 上设置断点,以查看join 运算符的实际作用。

    (OP)似乎这是问题所在,我不知道在加入时我需要一个密钥。我的映射器使用部分值来加入表,所以我认为空键很好,在加入之前应用了转换,它按预期工作。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2020-08-13
      • 2019-09-26
      • 1970-01-01
      • 2018-08-30
      • 2018-09-18
      • 2020-04-21
      • 2020-05-02
      相关资源
      最近更新 更多