You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何向Kafka Processor传递参数?附Kotlin代码示例

如何向Kafka Transformer传递单例参数

当然可以实现向Transformer传递单例参数,你只需要调整Transformer的构造逻辑,利用TransformerSupplier的lambda捕获外部变量即可,具体修改步骤如下:

  1. 修正BadPingsMarker类的构造与初始化逻辑
    Kafka Streams的Transformer接口中,init方法仅接收ProcessorContext一个参数,无法自定义额外入参。因此需要通过构造函数将BadPingIdentifier注入到Transformer中,调整后的类代码如下:
class BadPingsMarker(private val pingIdentifier: BadPingIdentifier) : Transformer<ID, Ping, KeyValue<ID, Ping>> {
    var state: KeyValueStore<String, Tuple<String, String>>? = null
    val logger: Logger = LogManager.getLogger(BadPingsMarker::class.java)!!

    override fun init(context: ProcessorContext) {
        state = context.getStateStore(MY_STATE_STORE) as KeyValueStore<String, Tuple<String, String>>
    }

    override fun transform(key: ID, value: Ping): KeyValue<ID, Ping> {
        val somevalue = value.somevalue
        val state_checker = state!![MY_STATE_STORE_A]

        // 这里可以正常使用pingIdentifier实现标记逻辑
        return KeyValue(key, value)
    }

    override fun close() {}
}
  1. 在identifyBadPings函数中传递单例参数
    在TransformerSupplier的lambda中直接创建BadPingsMarker实例并传入单例对象,Kotlin中外部变量默认是effectively final,可以被lambda捕获:
private fun identifyBadPings(
    pingStream: KStream<ID, Ping>,
    mySignletonBadPingIdentifier: BadPingIdentifier
): KStream<ID, Ping> {

    return pingStream.transform(
        TransformerSupplier { BadPingsMarker(mySignletonBadPingIdentifier) },
        MY_STATE_STORE
    )
}

这样就能成功将单例的BadPingIdentifier传递到BadPingsMarker中使用了。

内容的提问来源于stack exchange,提问作者Christos Hadjinikolis

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.28 03:23:16