【问题标题】:Spring Cloud Stream generates unnecessary complex Kafka topologies, why?Spring Cloud Stream 会生成不必要的复杂 Kafka 拓扑,为什么?
【发布时间】:2019-09-19 14:25:34
【问题描述】:

我有一个 KStream 应用程序,其中包含一堆 KStream、连接和其他操作。我启用了logging.level.org.springframework.kafka.config=debug 来验证正在生成的拓扑,并发现了很多根本没有意义的节点。

然后我将应用程序简化为:

interface ShippingKStreamProcessor {

    @Input("input")
    fun input(): KStream<Int, Customer>

}

@Suppress("UNCHECKED_CAST")
@Configuration
class ShippingKStreamConfiguration {

    @StreamListener
    fun process(@Input("input") input: KStream<Int, Customer> {}

}

奇怪的是,这样一个简单的 KStream 声明会生成这个复杂的拓扑:

2019-04-30 23:47:03.881 DEBUG 2944 --- [           main] o.s.k.config.StreamsBuilderFactoryBean   : Topologies:
   Sub-topology: 0
    Source: KSTREAM-SOURCE-0000000000 (topics: [customer])
      --> KSTREAM-MAPVALUES-0000000001
    Processor: KSTREAM-MAPVALUES-0000000001 (stores: [])
      --> KSTREAM-BRANCH-0000000003, KSTREAM-PROCESSOR-0000000002
      <-- KSTREAM-SOURCE-0000000000
    Processor: KSTREAM-BRANCH-0000000003 (stores: [])
      --> KSTREAM-BRANCHCHILD-0000000004, KSTREAM-BRANCHCHILD-0000000005
      <-- KSTREAM-MAPVALUES-0000000001
    Processor: KSTREAM-BRANCHCHILD-0000000004 (stores: [])
      --> KSTREAM-MAPVALUES-0000000007
      <-- KSTREAM-BRANCH-0000000003
    Processor: KSTREAM-BRANCHCHILD-0000000005 (stores: [])
      --> KSTREAM-PROCESSOR-0000000006
      <-- KSTREAM-BRANCH-0000000003
    Processor: KSTREAM-MAPVALUES-0000000007 (stores: [])
      --> none
      <-- KSTREAM-BRANCHCHILD-0000000004
    Processor: KSTREAM-PROCESSOR-0000000002 (stores: [])
      --> none
      <-- KSTREAM-MAPVALUES-0000000001
    Processor: KSTREAM-PROCESSOR-0000000006 (stores: [])
      --> none
      <-- KSTREAM-BRANCHCHILD-0000000005

原生 Kafka 应用程序中的相同简单流会产生更合乎逻辑的拓扑:

fun main(args: Array<String>) {

    val builder = StreamsBuilder()

    val streamsConfiguration = Properties()
    streamsConfiguration[StreamsConfig.APPLICATION_ID_CONFIG] = "kafka-shipping-service"
    streamsConfiguration[StreamsConfig.BOOTSTRAP_SERVERS_CONFIG] = "http://localhost:9092"
    streamsConfiguration[AbstractKafkaAvroSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG] = "http://localhost:8081"

    val serdeConfig = mapOf(
        AbstractKafkaAvroSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG to "http://localhost:8081",
        AbstractKafkaAvroSerDeConfig.VALUE_SUBJECT_NAME_STRATEGY to TopicRecordNameStrategy::class.java.name
    )

    //val byteArraySerde = Serdes.ByteArray()
    val intSerde = Serdes.IntegerSerde()
    val customerSerde = SpecificAvroSerde<Customer>()
    customerSerde.configure(serdeConfig, false)

    val customerStream = builder.stream<Int, Customer>("customer",
        Consumed.with(intSerde, customerSerde)) as KStream<Int, Customer>

    val topology = builder.build()
    println(topology.describe())

    val streams = KafkaStreams(topology, streamsConfiguration)
    streams.start()
}

拓扑:

Topologies:
   Sub-topology: 0
    Source: KSTREAM-SOURCE-0000000000 (topics: [customer])
      --> none

Spring Cloud Stream 生成如此复杂拓扑的原因是什么?

【问题讨论】:

  • 这很有趣。会调查的。我能想到的一个原因是由于活页夹需要做的消息转换。因此,如果您可以启用本机解码/编码,我认为这会稍微减少拓扑。除此之外,还有其他一些路径,binder 可能会添加更多的拓扑,例如使用 DLQ 等。

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


【解决方案1】:

@codependent 在拓扑中有这些额外处理器的原因是因为您使用的是框架提供的 de/serailzers(本机解码和编码默认为 false)。基本上,我们以byte[] 的形式接收来自Kafka 主题的数据,然后在内部进行转换。对于这些转换,我们会经过一些额外的处理器,因此您最终会得到更深层次的拓扑。

这是 Java 中的基本 StreamListener(与上面的内容差不多,但使用更简单的值类型):

@StreamListener
public void process(@Input("input") KStream<Integer, String> input ) {

}

使用活页夹中的标准开箱即用设置,我能够获得与您观察到的相同的更深层次的拓扑。但是,当我如下所示修改应用程序的配置时,

spring.cloud.stream.kafka.streams:
  binder.configuration:
    default.key.serde: org.apache.kafka.common.serialization.Serdes$IntegerSerde
    default.value.serde: org.apache.kafka.common.serialization.Serdes$StringSerde
spring.cloud.stream.bindings.input.consumer.useNativeDecoding: true

我的拓扑简化如下:

2019-05-01 18:02:12.705 DEBUG 67539 --- [           main] o.s.k.config.StreamsBuilderFactoryBean   : Topologies:
   Sub-topology: 0
    Source: KSTREAM-SOURCE-0000000000 (topics: [hello-1])
      --> KSTREAM-MAPVALUES-0000000001
    Processor: KSTREAM-MAPVALUES-0000000001 (stores: [])
      --> none
      <-- KSTREAM-SOURCE-0000000000

这仍然与您从普通 Kafka Streams 应用程序中获得的拓扑不同,但事实证明这是我们可以在 binder 中改进以避免的情况。简而言之,通过切换到 Kafka Streams 提供的本机解码和编码,您可以避免由 binder 构建的所有这些额外级别的拓扑。

在某些情况下,您别无选择,只能依赖 Spring Cloud Stream 提供的反序列化,例如,您从基于 Spring Cloud Stream 的生产者接收数据,该生产者使用了一些特殊的序列化器。我认为在您的情况下确实如此,因为据我所知,您的生产者基于 Spring Cloud Stream 并且使用框架提供的 Avro 序列化程序。在这种情况下,在您的处理器中使用 Kafka Stream 的 Avro Serde 将不起作用,因为这些序列化程序不兼容。所以这里有一些你的选择。

方法#1:

  1. 让您的生产者使用 Kafka 提供的本机序列化程序。
  2. 然后使用在您的 Kafka Streams 应用程序中使用相同序列化器/反序列化器的 Serde。

方法 #2:

  1. 使用 SCSt 提供的消息序列化程序。
  2. 然后使用默认的 Kafka Streams binder 提供的默认反序列化。

#2 的缺点显然是您在上面提到的,即更深的拓扑。根据您的用例和吞吐量,这可能没问题。如果这成为一个真正的性能问题,我们可以在框架完成转换时尝试简化这个过程。

话虽如此,我在 Kafka 活页夹中创建了一个 issue,以便在下一个版本的活页夹中进行更改。欢迎您的反馈、建议、赞成/反对票。

【讨论】:

  • 从 Spring Cloud Stream Kafka Streams binder 3.0 版本开始,默认反序列化由 Kafka 原生完成。因此 Spring 应用生成的拓扑与原生应用生成的拓扑是等价的。
猜你喜欢
  • 1970-01-01
  • 2021-09-28
  • 2019-09-19
  • 1970-01-01
  • 1970-01-01
  • 2018-03-28
  • 2018-04-28
  • 1970-01-01
  • 2018-03-09
相关资源
最近更新 更多