【发布时间】:2021-02-07 07:01:24
【问题描述】:
我有一个简单的流处理器(不是消费者/生产者),看起来像这样 (Kotlin)
@Bean
fun processFoo():Function<KStream<FooName, FooAddress>, KStream<FooName, FooAddressPlus>> {
return Function { input-> input.map { key, value ->
println("\nPAYLOAD KEY: ${key.name}\n");
println("\nPAYLOAD value: ${value.address}\n");
val output = FooAddressPlus()
output.address = value.address
output.name = value.name
output.plus = "$value.name-$value.address"
KeyValue(key, output)
}}
}
FooName、FooAddress 和 FooAddressPlus 这些类与处理器在同一个包中。
这是我的配置文件:
spring.cloud.stream.kafka.binder:
brokers: localhost:9093
spring.cloud.stream.function.definition: processFoo
spring.cloud.stream.kafka.streams.binder.functions.processFoo.applicationId: foo-processor
spring.cloud.stream.bindings.processFoo-in-0:
destination: foo.processor
spring.cloud.stream.bindings.processFoo-out-0:
destination: foo.processor.out
spring.cloud.stream.kafka.streams.binder:
deserializationExceptionHandler: logAndContinue
configuration:
default.key.serde: org.apache.kafka.common.serialization.Serdes$StringSerde
default.value.serde: org.apache.kafka.common.serialization.Serdes$StringSerde
commit.interval.ms: 1000
运行处理器时出现此错误:
The class '<here_comes_package>.FooAddress' is not in the trusted packages: [java.util, java.lang].
If you believe this class is safe to deserialize, please provide its name.
If the serialization is only done by a trusted source, you can also enable trust all (*).
在使用 Kafka Streams Binder 流处理器时,为所有内容设置可信包的最佳方法是什么? (没有消费者/生产者,只有流处理器)
非常感谢!
【问题讨论】:
-
@jokarls 谢谢。这是我之前尝试过的事情之一,但我有一个流处理器,在那个答案中,他们分别使用消费者和生产者。在这种情况下,他们要么在配置中设置消费者/生产者属性(spring.kafka.consumer.properties.spring.json.trusted.packages),要么在消费者/生产者工厂中设置反序列化器。因为我有一个流处理器,所以我没有可以设置的消费者或生产者道具或工厂,但我可能错了。
标签: kotlin apache-kafka spring-cloud-stream