【发布时间】:2019-02-21 09:39:27
【问题描述】:
我对从 Stream 内部调用交互式查询有特殊要求。这是因为我需要创建一个新的 Stream,它应该在 State Store 中包含数据。以下截断代码:
tempModifiedDataStream.to(topic.getTransformedTopic(), Produced.with(Serdes.String(), Serdes.String()));
GlobalKTable<String, String> myMetricsTable = builder.globalTable(
topic.getTransformedTopic(),
Materialized.<String, String, KeyValueStore<Bytes, byte[]>>as(
topic.getTransformedStoreName() /* table/store name */)
.withKeySerde(Serdes.String()) /* key serde */
.withValueSerde(Serdes.String()) /* value serde */
);
KafkaStreams streams = new KafkaStreams(builder.build(), kStreamsConfigs());
KStream<String, String> tempAggrDataStream = tempModifiedDataStream
.flatMap((key, value) -> {
try {
List<KeyValue<String, String>> result = new ArrayList<>();
ReadOnlyKeyValueStore<String, String> keyValueStore =
streams .store(
topic.getTransformedStoreName(),
QueryableStoreTypes.keyValueStore());
在最后一行中,要访问 State Store,我需要拥有 KafkaStreams 对象,并且当我创建 KafkaStreams 对象时拓扑已完成。这种方法的问题在于“tempAggrDataStream”因此不是拓扑的一部分,并且该部分代码不会被执行。而且我不能移动下面的 KafkaStreams 定义,否则我不能调用交互式查询。
我对 Kafka Streams 有点陌生;这对我来说是不是很傻?
【问题讨论】:
-
您能否更准确地描述一下,您想要实现什么?为什么需要这家国营商店?也许
join是你需要的(tempModifiedDataStream.join(myMetricsTable)? -
每当更新 Ktable 时,我需要将 Ktable 的全部内容输出到新的 Topic。即使我将原始 Stream 与它的 Ktable 进行左连接,我认为我无法实现这一点。
-
我仍然无法满足您的需求。当至少一个值发生变化时,您想复制整个主题的内容吗?每次更新后目标主题都会改变吗?
-
我的要求很简单。每当 KTable 的内容发生变化时,我想将 KTable 的所有键值对写入主题。例如,假设我的 KTable 当前包含
, , 一个新的更新是 。那时我想将诸如 {"asia":"55","europe":"66","australia":"77"} 之类的 json 写入主题。即遍历 KTable 的所有条目并将数据输出到某处 -
这不是 Kafka Streams 的好用例。另外,如果您说“新主题”,您的意思是每次更新都有一个新主题吗?还是每次更新都会重复使用一个主题?
标签: apache-kafka apache-kafka-streams kafka-interactive-queries