【发布时间】:2019-01-11 19:27:28
【问题描述】:
我们有一个场景,我们想在集群#1 上使用来自 kafka 主题的数据,但在集群#2 上创建 KTable 主题(重新分区和更改日志)。
频道绑定 -
spring.cloud.stream.bindings.member.destination: member
spring.cloud.stream.bindings.member.consumer.useNativeDecoding: true
spring.cloud.stream.bindings.member.consumer.headerMode: raw
spring.cloud.stream.kafka.streams.bindings.member.consumer.keySerde: org.apache.kafka.common.serialization.Serdes$StringSerde
spring.cloud.stream.kafka.streams.bindings.member.consumer.valueSerde: io.confluent.kafka.streams.serdes.avro.GenericAvroSerde
创建 Ktable -
protected KTable<String, GenericRecord> createKTable(String field, KStream<String, GenericRecord> stream, String stateStore) {
return stream
.map((s, genericRecord) -> KeyValue.pair(field, genericRecord))
.groupByKey()
.reduce((oldVal, newVal) -> newVal, Materialized.as(stateStore));
}
所以成员主题在 cluster#1 上,但是我们想在不同的 cluster 上创建下面的 ktable 主题,不知道在这种情况下如何使用两个不同的 kafka binder -
application-member-store-repartition
application-member-store-changelog
【问题讨论】:
-
您是否尝试过不使用 Spring Cloud Stream 而直接使用 Kafka Streams?我认为 Spring Cloud Stream 中的机制也将相同。如果 Kafka Streams 允许,绑定器可以以某种方式委托它。
标签: apache-kafka apache-kafka-streams spring-cloud-stream spring-kafka