如何向Kafka Processor传递参数?附Kotlin代码示例
如何向Kafka Transformer传递单例参数
当然可以实现向Transformer传递单例参数,你只需要调整Transformer的构造逻辑,利用TransformerSupplier的lambda捕获外部变量即可,具体修改步骤如下:
- 修正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() {} }
- 在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
相关产品推荐
相关产品推荐

