【问题标题】:KTable and KStream Space considerations understandingKTable和KStream Space注意事项理解
【发布时间】:2021-05-28 17:42:05
【问题描述】:

我们正在使用 KSQLDB 执行 POC,但有一些疑问:-

我有一个名为 USERPROFILE 的 Kafka 主题,它有大约 1 亿条唯一记录和 10 天的保留政策。此 Kafka 主题继续从其底层 RDBMS 表中实时接收 INSERT/UPDATE 类型的事件。

以下是这个kafka主题中接收到的记录的简单结构:-

{"userid":1001,"firstname":"Hemant","lastname":"Garg","countrycode":"IND","rating":3.7} 

1.) 我们已经在上述主题上打开了一个 Kafka Stream :-

create STREAM userprofile_stream (userid INT, firstname VARCHAR, lastname VARCHAR, countrycode VARCHAR, rating DOUBLE) WITH (VALUE_FORMAT = 'JSON', KAFKA_TOPIC = 'USERPROFILE')>;

2.) 因为,给定的 userId 可以更新,我们只想要唯一的记录(对于每个 userId),我们还在上述主题上打开了另一个 Kafka 表:-

ksql> create TABLE userprofile_table(userid VARCHAR PRIMARY KEY, firstname VARCHAR, lastname VARCHAR, countrycode VARCHAR, rating DOUBLE) WITH (KAFKA_TOPIC = 'USERPROFILE', VALUE_FORMAT = 'DELIMITED');

问题是:-

  • 打开 KTable 是否需要额外的磁盘空间?例如,Kafka 主题有 1 亿条记录,相同的记录是否也会出现在 KTable 中,或者只是底层 kafka 主题的一些虚拟视图?

  • 我们打开的信息流也有同样的问题。打开 KStream 是否需要磁盘(经纪人服务器的)额外空间?例如,Kafka 主题有 1 亿条记录,相同的记录是否也会出现在 KStream 中,或者它只是底层 kafka 主题的一些虚拟视图?

  • 说,我们在 5 月 1 日收到 id 为 1001 的记录,然后在 5 月 11 日,该记录将不再在 Kafka 主题上可用,但是该记录是否仍会出现在 kstream / Ktable 上? KStream / KTable 是否有一些保留政策,就像我们对 Topic 的保留政策一样?

答案将不胜感激。

-- 最好的 阿地亚

【问题讨论】:

    标签: apache-kafka apache-kafka-streams confluent-platform ksqldb ktable


    【解决方案1】:

    ksqlDB 服务器由Kafka Streams 提供支持。因此,当你创建流或表时,服务器会分别创建一个 KStream 或 KTable。

    此外,KStream 和 KTables 由 Kafka 中的主题支持。因此,在 ksqlDB 服务器上创建流和表将在您的 Kafka 集群上创建实际的主题。

    话虽如此,来自 ksqlDB 的流和表是根据需要具体化并进行了相当优化的,Confluent 的这两篇文章通过良好的视觉帮助提供了更多关于内部行为的见解:

    您甚至可以自己查看创建的数据。为了举例,我创建了:

    • 来自原始主题的MESSAGES_STREAM
    • 来自上方流的MATERIALIZED_MESSAGES_STTREAM
    • 第一个流中的 MESSAGES

    以下是创建命令供参考:

    ksql> CREATE STREAM messages_stream (user_id BIGINT KEY, message VARCHAR) 
      WITH (KAFKA_TOPIC = 'hello_topic_json', VALUE_FORMAT='JSON');
    
    ksql> CREATE STREAM materialized_messages_stream AS 
      SELECT user_id, UCASE(message) 
      FROM messages_stream 
    EMIT CHANGES;
    
    ksql> CREATE TABLE messages AS
      SELECT user_id, count(*) as msg_count
      FROM messages_stream
      GROUP BY user_id
    EMIT CHANGES;
    

    通过查看 ksqlDB 中的详细信息,我们可以看到第一个流使用原始主题作为源:

    ksql> describe extended MESSAGES_STREAM;
    
    Name                 : MESSAGES_STREAM
    Type                 : STREAM
    Timestamp field      : Not set - using <ROWTIME>
    Key format           : KAFKA
    Value format         : JSON
    Kafka topic          : hello_topic_json (partitions: 1, replication: 1)
    -- […]
    
    ksql> describe extended MATERIALIZED_MESSAGES_STREAM;
    
    Name                 : MATERIALIZED_MESSAGES_STREAM
    Type                 : STREAM
    Timestamp field      : Not set - using <ROWTIME>
    Key format           : KAFKA
    Value format         : JSON
    Kafka topic          : MATERIALIZED_MESSAGES_STREAM (partitions: 1, replication: 1)
    -- […]
    
    ksql> describe extended MESSAGES;
    
    Name                 : MESSAGES
    Type                 : TABLE
    Timestamp field      : Not set - using <ROWTIME>
    Key format           : KAFKA
    Value format         : JSON
    Kafka topic          : MESSAGES (partitions: 1, replication: 1)
    -- […]
    

    查看集群上声明的主题,我们可以看到第二个流、表及其在后台创建的 changelog 主题:

    $ ./kafka-topics.sh --bootstrap-server localhost:29092 --list
    MATERIALIZED_MESSAGES_STREAM
    MESSAGES
    __consumer_offsets
    __transaction_state
    _confluent-ksql-ksql_docker_command_topic
    _confluent-ksql-ksql_dockerquery_CTAS_MESSAGES_1-Aggregate-Aggregate-Materialize-changelog
    _schemas
    hello_topic_json
    

    您还可以看到流和表之间的保留策略不同。前者会删除旧记录,后者会压缩数据:

    $ ./kafka-topics.sh --bootstrap-server localhost:29092 --topic MATERIALIZED_MESSAGES_STREAM --describe
    Topic: MATERIALIZED_MESSAGES_STREAM PartitionCount: 1   ReplicationFactor: 1    Configs: cleanup.policy=delete
    
    $ ./kafka-topics.sh --bootstrap-server localhost:29092 --topic MESSAGES --describe
    Topic: MESSAGES PartitionCount: 1   ReplicationFactor: 1    Configs: cleanup.policy=compact
    

    TL;DR,回到你的问题:

    1. 是的,打开 KTable 需要空间,但很可能不会是已用字节的 1 对 1 映射。
    2. 在您的情况下,流很可能会使用主题作为参考,并且不会占用更多空间,因为您没有发生任何数据转换。
    3. 表的保留策略是压缩,因此您的条目在表上仍然可用。但是,在您的信息流中,数据将与参考主题中的数据一样可用。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-02-23
      • 1970-01-01
      • 1970-01-01
      • 2020-10-17
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多