【问题标题】:Invoking Kafka Interactive Queries from inside a Stream从流中调用 Kafka 交互式查询
【发布时间】: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


【解决方案1】:

如果你想实现每次数据修改后主题内容的全部发送,我认为你应该使用Processor API。

您可以使用状态存储创建org.apache.kafka.streams.kstream.Transformer。 对于每个处理消息,它将更新状态存储并将所有内容发送到下游。 这不是很有效,因为它将为每个处理消息转发主题/状态存储的全部内容(可能是数千、数百万条记录)。

如果您只需要 latest 值,将主题 cleanup.policy 设置为 compact 就足够了。而从其他站点使用KTable,它给出了表的抽象(流的快照)

用于转发状态存储的全部内容的示例 Transformer 代码如下。整个工作在transform(String key, String value)方法中完成。

public class SampleTransformer
        implements Transformer<String, String, KeyValue<String, String>> {

    private String stateStoreName;
    private KeyValueStore<String, String> stateStore;
    private ProcessorContext context;

    public SampleTransformer(String stateStoreName) {
        this.stateStoreName = stateStoreName;
    }

    @Override
    @SuppressWarnings("unchecked")
    public void init(ProcessorContext context) {
        this.context = context;
        stateStore = (KeyValueStore) context.getStateStore(stateStoreName);
    }

    @Override
    public KeyValue<String, String> transform(String key, String value) {
        stateStore.put(key, value);
        stateStore.all().forEachRemaining(keyValue -> context.forward(keyValue.key, keyValue.value));
        return null;
    }

    @Override
    public void close() {

    }
}

有关处理器 API 的更多信息可以找到:

如何将处理器 API 与 Stream DSL 结合可以找到:

【讨论】:

  • 从上下文中获取的 stateStore 是“分片的”,因为它将只包含那些属于某个分区或流任务的键。如果有 5 个输入分区,您的流应用程序的单个实例正在使用它们,那么将有 5 个单独的 SampleTransformation 对象,每个对象都有 1 个分区的键值。所以,你只会发布世界状态的一部分。要使用所有键访问 stateStore,我认为您需要查看交互式查询。
猜你喜欢
  • 2021-01-24
  • 2019-06-05
  • 1970-01-01
  • 2019-08-06
  • 1970-01-01
  • 2017-12-19
  • 1970-01-01
  • 1970-01-01
  • 2021-11-08
相关资源
最近更新 更多