【问题标题】:can i traverse the items in a KTable from an external method我可以从外部方法遍历 KTable 中的项目吗
【发布时间】: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


    【解决方案1】:

    经过长时间的调查,解决方案是查询它后面的存储(rocksDB)而不是表。

    如此处所述:confluent

    正确的解决方案是使用 GlobalKTable 来避免如here 所讨论的“状态存储可能已迁移到另一个实例”错误。

    这段代码在 kafka 0.10.2.1 中为我工作:

        private val accessTokenTable: GlobalKTable[String, String] = builder.globalTable(token_topic_name, token_store_string)
    
        private val stream: KafkaStreams = new KafkaStreams(builder, streamingConfig)
        stream.cleanUp()
        stream.start()
        val store: ReadOnlyKeyValueStore[String,String] = stream.store(token_store_string,QueryableStoreTypes.keyValueStore[String,String]())
    

    ....

        override def refreshTokensFlow = {
    
           store.all.asScala.map( tuple => {
           // logic goes here
               System.out.println(tuple.key + ": " + tuple.value)
           }
        }
    

    【讨论】:

      【解决方案2】:

      正确的解决方案是使用 GlobalKTable 来避免 here 讨论的“状态存储可能已迁移到另一个实例”错误。

      由于您回答了自己的问题,并且显然在您的后续行动中遇到了另一个问题,所以让我扩展您在回答中所说的内容,以帮助该问题线程的其他读者。

      • 如果您使用的是 KTable(已分区 = KTable 的每个“实例”只能看到总表数据的一部分)通常,您需要做的是防范此异常并重试。思考:try-catch-retry。
      • 如果您使用的是 GlobalKTable,那么您可以回避这个问题,因为 GlobalKTable 的每个实例都有整个表数据的完整副本。

      注意:通常,您不会在 KTable 与 GlobalKTable 之间做出决定,因为您想防止“状态存储可能已迁移”的情况,而是因为这两种抽象为您的应用程序提供了不同的语义。例如,使用 KTable 而不是 GlobalKTable 有很多很好的理由——如果你这样做了,你只需要了解我们刚刚在这里讨论的内容(这也包含在文档中,但显然不明显/考虑到您确实遇到了这个问题,这已经足够清楚了)。

      希望这会有所帮助!

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2012-10-05
        • 1970-01-01
        • 1970-01-01
        • 2022-11-11
        • 2013-01-17
        • 2023-03-25
        • 2023-04-11
        • 2022-07-27
        相关资源
        最近更新 更多