代码之家  ›  专栏  ›  技术社区  ›  Dth

如何在内存Kafka Streams状态存储上启用缓存

  •  3
  • Dth  · 技术社区  · 8 年前

    我想减少向下游发送的数据数量,因为我只关心给定键的最后一个值,所以我通过以下方式读取主题中的数据:

    KTable table = build.table("inputTopic", Materialized.as("myStore"));
    

    为什么?因为在幕后,数据正在被缓存,如前所述 here ,并且仅当 犯罪间隔太太 隐藏物最大字节数。缓冲 踢进来。

    到目前为止还不错,但在这种情况下,我根本没有利用RocksDB,所以我想用内存存储的默认实现来代替它。我隐式启用缓存,以防万一。

    Materialized.as(Stores.inMemoryKeyValueStore("myStore")).withCachingEnabled();
    

    但是,这不起作用-数据没有被缓存,每个记录都被发送到下游。

    是否有其他方法可以启用缓存?或者也许有更好的方法来实现我的目标?

    1 回复  |  直到 8 年前
        1
  •  3
  •   Dth    8 年前

    看来我错了,内存状态存储缓存工作正常。我将简要说明我是如何测试它的,也许有人会发现它很有用。我制作了一个非常基本的Kafka Streams应用程序,它只读取抽象为KTable的主题。

    public class Main {
        public static void main(String[] args) {
            StreamsBuilder builder = new StreamsBuilder();
            Logger logger = LoggerFactory.getLogger(Main.class);
    
            builder.table("inputTopic", Materialized.as(Stores.inMemoryKeyValueStore("myStore")).withCachingEnabled())
                    .toStream()
                    .foreach((k, v) -> logger.info("Result: {} - {}", k, v));
    
            new KafkaStreams(builder.build(), getProperties()).start();
        }
    
        private static Properties getProperties() {
            Properties properties = new Properties();
            properties.put(APPLICATION_ID_CONFIG, "testApp");
            properties.put(BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
            properties.put(COMMIT_INTERVAL_MS_CONFIG, 10000);
            properties.put(CACHE_MAX_BYTES_BUFFERING_CONFIG, 10 * 1024 * 1024L);
            properties.put(DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName());
            properties.put(DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName());
            return properties;
        }
    }
    

    然后我运行了卡夫卡的控制台制作人:

    /kafka-console-producer.sh --broker-list localhost:9092 --topic inputTopic --property "parse.key=true" --property "key.separator=:"
    

    并发送了几条消息:a:a、a:b、a:c。应用程序中只能看到最后一条消息,因此缓存工作正常。

    2018-03-06 21:21:57信息主:26-结果:a-c

    我还稍微更改了流以检查 aggregate 方法

    builder.stream("inputTopic")
            .groupByKey()
            .aggregate(() -> "", (k, v, a) -> a + v, Materialized.as(Stores.inMemoryKeyValueStore("aggregate")))
            .toStream()
            .foreach((k, v) -> logger.info("Result: {} - {}", k, v));
    

    我用同一个键连续快速发送了几条消息,只收到了一个结果,因此数据没有立即发送到下游-完全按照预期。

    推荐文章