【问题标题】:Kafka Streams: action on n-th eventKafka Streams:对第 n 个事件的操作
【发布时间】:2020-04-01 12:33:05
【问题描述】:

我正在尝试找到对 Kafka Streams 中的第 n 个事件执行操作的最佳方法。

我的情况:我有一个带有一些事件的输入流。我必须通过 eventType == login 过滤它们,并在每个 n 次登录(比如说,第五次)为同一个 accountId 发送这个事件到输出流。

经过一些调查和不同的尝试,我得到了下面的代码版本(我正在使用 Kotlin)。

data class Event(
    val payload: Any = {},
    val accountId: String,
    val eventType: String = ""
)
// intermediate class to keep the key and value of the original event
data class LoginEvent(
    val eventKey: String,
    val eventValue: Event
)
fun process() {
        val userLoginsStoreBuilder = Stores.keyValueStoreBuilder(
            Stores.persistentKeyValueStore("logins"),
            Serdes.String(),
            Serdes.Integer()
        )
        val streamsBuilder = StreamsBuilder().addStateStore(userCheckInsStoreBuilder)
        val inputStream = streamsBuilder.stream<String, String>(inputTopic)

        inputStream.map { key, event ->
            KeyValue(key, json.readValue<Event>(event))
        }.filter { _, event -> event.eventType == "login" }
             .map { key, event -> KeyValue(event.accountId, LoginEvent(key, event)) }
             .transform(
                    UserLoginsTransformer("logins", 5),
                    "logins"
                )
             .filter { _, value -> value }
             .map { key, _ -> KeyValue(key.eventKey, json.writeValueAsString(key.eventValue)) }
             .to("fifth_login", Produced.with(Serdes.String(), Serdes.String()))

        ...
    }
class UserLoginsTransformer(private val storeName: String, private val loginsThreshold: Int = 5) :
    TransformerSupplier<String, CheckInEvent, KeyValue< LoginEvent, Boolean>> {

    override fun get(): Transformer<String, LoginEvent, KeyValue< LoginEvent, Boolean>> {
        return object : Transformer<String, LoginEvent, KeyValue< LoginEvent, Boolean>> {
            private lateinit var store: KeyValueStore<String, Int>

            @Suppress("UNCHECKED_CAST")
            override fun init(context: ProcessorContext) {
                store = context.getStateStore(storeName) as KeyValueStore<String, Int>
            }

            override fun transform(key: String, value: LoginEvent): KeyValue< LoginEvent, Boolean> {
                val counter = (store.get(key) ?: 0) + 1
                return if (counter == loginsThreshold) {
                    store.delete(key)
                    KeyValue(value, true)
                } else {
                    store.put(key, counter)
                    KeyValue(value, false)
                }
            }

            override fun close() {
            }
        }
    }
}

我最担心的是 transform 函数在我的情况下不是线程安全的。我已经检查了在我的案例中使用的 KV 存储的实现,这是 RocksDB 存储(非事务性),因此值可能会在读取和比较之间更新,并且错误的事件将被发送到输出。

我的其他想法:

  1. 将物化视图用作没有转换器的商店,但我坚持实施。
  2. 创建将使用 TransactionalRocksDB 的自定义持久 KV 存储(不确定是否值得)。
  3. 创建将在内部使用 ConcurrentHashMap 的自定义持久 KV 存储(在我们预期的许多用户的情况下,它可能会导致高内存消耗)。

另一个注意事项:我使用的是 Spring Cloud Stream,所以也许这个框架有一个适合我的案例的内置解决方案,但我没有找到它。

如果有任何建议,我将不胜感激。提前致谢。

【问题讨论】:

    标签: apache-kafka apache-kafka-streams spring-cloud-stream spring-cloud-stream-binder-kafka


    【解决方案1】:

    我最担心的是转换函数在我的情况下不是线程安全的。我已经检查了在我的案例中使用的 KV 存储的实现,这是 RocksDB 存储(非事务性),因此值可能会在读取和比较之间更新,并且错误的事件将被发送到输出。

    没有理由担心。如果您使用多个线程运行,每个线程将拥有自己的 RocksDB 存储整体数据的一个分片(请注意,整体状态是基于输入主题分区进行分片的,并且单个分片永远不会被不同的线程处理)。因此,您的代码将正常工作。您唯一需要确保的是,该数据是由accountId 划分的,这样单个帐户的登录事件就会进入同一个分片。

    如果您输入的数据在写入您的输入主题时已经被accountId 分区,则您无需执行任何操作。如果没有,并且您可以控制上游应用程序,那么在上游的应用程序生产者中使用自定义分区器来获得您需要的分区可能是最简单的。如果您无法更改上游应用程序,则需要在将accountId 设置为新键后重新分区数据,即在调用transform() 之前执行through()

    【讨论】:

    • 我无法控制上游,所以我需要重新分区。还有一点需要注意 - 我们不是在谈论多个线程(JVM 线程),而是关于正在阅读同一主题的多个消费者。
    • 不同的消费者会总是使用 Kafka Streams 在不同的线程上运行。
    猜你喜欢
    • 1970-01-01
    • 2020-01-07
    • 1970-01-01
    • 1970-01-01
    • 2022-08-19
    • 2015-04-23
    • 2020-02-23
    • 2017-12-22
    • 1970-01-01
    相关资源
    最近更新 更多