【问题标题】:Testing Spring cloud stream with kafka stream binder: using TopologyTestDriver I get the error of "The class is not in the trusted packages"使用 kafka 流绑定器测试 Spring 云流:使用 TopologyTestDriver 我收到“该类不在受信任的包中”的错误
【发布时间】:2020-10-24 17:31:44
【问题描述】:

我有一个使用 kafka 流绑定器的简单流处理器(不是消费者/生产者)。

@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)
    }}
}

我正在尝试使用 TopologyTestDriver 对其进行测试:

@SpringBootTest(
        webEnvironment = SpringBootTest.WebEnvironment.NONE,
        classes = [Application::class, FooProcessor::class]
)
class FooProcessorTests {
    var testDriver: TopologyTestDriver? = null
    val INPUT_TOPIC = "input"
    val OUTPUT_TOPIC = "output"

    val inputKeySerde: Serde<FooName> = JsonSerde<FooName>()
    val inputValueSerde: Serde<FooAddress> = JsonSerde<FooAddress>()
    val outputKeySerde: Serde<FooName> = JsonSerde<FooName>()
    val outputValueSerde: Serde<FooAddressPlus> = JsonSerde<FooAddressPlus>()

    fun getStreamsConfiguration(): Properties? {
        val streamsConfiguration = Properties()
        streamsConfiguration[StreamsConfig.APPLICATION_ID_CONFIG] = "TopologyTestDriver"
        streamsConfiguration[StreamsConfig.BOOTSTRAP_SERVERS_CONFIG] = "dummy:1234"
        streamsConfiguration[JsonDeserializer.TRUSTED_PACKAGES] = "*"
        streamsConfiguration["spring.kafka.consumer.properties.spring.json.trusted.packages"] = "*"
        return streamsConfiguration
    }

    @Before
    fun setup() {
        val builder = StreamsBuilder()
        val input: KStream<FooName, FooAddress> = builder.stream(INPUT_TOPIC, Consumed.with(inputKeySerde, inputValueSerde))
        val processor = FooProcessor()
        val output: KStream<FooName, FooAddressPlus> = processor.processFoo().apply(input)
        output.to(OUTPUT_TOPIC, Produced.with(outputKeySerde, outputValueSerde))
        testDriver = TopologyTestDriver(builder.build(), getStreamsConfiguration())
    }

    @After
    fun tearDown() {
        try {
            testDriver!!.close()
        } catch (e: RuntimeException) {
            // https://issues.apache.org/jira/browse/KAFKA-6647 causes exception when executed in Windows, ignoring it
            // Logged stacktrace cannot be avoided
            println("Ignoring exception, test failing in Windows due this exception:" + e.localizedMessage)
        }
    }

    @org.junit.Test
    fun testOne() {
        val inputTopic: TestInputTopic<FooName, FooAddress> =
                testDriver!!.createInputTopic(INPUT_TOPIC, inputKeySerde.serializer(), inputValueSerde.serializer())
        val key = FooName()
        key.name = "sherlock"
        val value = FooAddress()
        value.name = "sherlock"
        value.address = "Baker street"
        inputTopic.pipeInput(key, value)
        val outputTopic: TestOutputTopic<FooName, FooAddressPlus> =
                testDriver!!.createOutputTopic(OUTPUT_TOPIC, outputKeySerde.deserializer(), outputValueSerde.deserializer())
        val message = outputTopic.readValue()

        assertThat(message.name).isEqualTo(key.name)
        assertThat(message.address).isEqualTo(value.address)
    }
}

运行它时,我在inputTopic.pipeInput(key, value) 行中收到此错误

类 'package.FooAddress' 不在受信任的包中:[java.util, java.lang]。如果您认为此类可以安全反序列化,请提供其名称。如果序列化仅由受信任的来源完成,您还可以启用全部信任 ()。*

关于如何解决这个问题的任何想法?在getStreamsConfiguration() 中设置这些属性没有帮助。请注意,这是一个流处理器,而不是消费者/生产者。

非常感谢!

【问题讨论】:

  • 非常感谢@GaryRussell。我会检查我的另一个问题。顺便说一句,感谢出色的 Spring Cloud Stream 框架
  • 不幸的是,我无法让它工作,可能是我做错了什么,所以我非常感谢您的指导@GaryRussell 在我的getStreamsConfiguration 方法中我添加了:streamsConfiguration["spring.cloud.streams.kafka.streams.binder.configuration.spring.json.trusted.packages"] = "*" 但仍然是由于is not in the trusted packages: [java.util, java.lang],测试失败。谢谢!
  • 我将研究为什么它不能与测试驱动程序一起工作 - 它对 Spring 一无所知,因此添加 spring 属性将无济于事。
  • 非常感谢@GaryRussell。如果有帮助,我可以提供带有代码的 git repo。

标签: unit-testing kotlin spring-cloud-stream


【解决方案1】:

当 Kafka 自己创建 Serde 时,它​​会通过调用 configure() 来应用属性。

由于您要自己实例化 Serde,因此您需要在其上调用 configure(),并传入属性映射。

这就是受信任的包属性被传播到反序列化器的方式。

或者,您可以在解串器上调用setTrustedPackages()

【讨论】:

  • 天哪,它有效!真的不能说我有多感激。非常感谢。
  • 只是出于好奇:我如何在解串器上调用setTrustedPackages()?我无权访问它,是吗?我只有我创建的 serdes。非常感谢!
  • 有一个构造函数,您可以在其中提供预配置的(反)序列化程序public JsonSerde(JsonSerializer&lt;T&gt; jsonSerializer, JsonDeserializer&lt;T&gt; jsonDeserializer)。或者,您可以将deserializer() 返回的值转换为JsonDeserializer&lt;?&gt;
  • 我看到你是新来的 - 见stackoverflow.com/help/someone-answers
  • 非常感谢@GaryRussell 的提示,按照程序给出答案。
【解决方案2】:

因此,为了完整起见,下面是按照@GaryRussell 建议配置 serde 时代码的外观:

private fun getStreamsConfiguration(): Properties? {
    // Don't set the trusted packages here since topology test driver does not know about Spring
    val streamsConfiguration = Properties()
    streamsConfiguration[StreamsConfig.APPLICATION_ID_CONFIG] = "TopologyTestDriver"
    streamsConfiguration[StreamsConfig.BOOTSTRAP_SERVERS_CONFIG] = "dummy:1234"
}
@Before
fun setup() {
    val builder = StreamsBuilder()
    // Set the trusted packages for all serdes
    val config = mapOf<String, String>(JsonDeserializer.TRUSTED_PACKAGES to "*")
    inputKeySerde.configure(config, true)
    inputValueSerde.configure(config, false)
    outputKeySerde.configure(config, true)
    outputValueSerde.configure(config, false)
}

其余代码仍如问题中所述。所有功劳归功于@GaryRusell。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-05-11
    • 2021-08-07
    • 2021-12-25
    • 2021-05-08
    • 2019-07-05
    • 2020-02-24
    相关资源
    最近更新 更多