【问题标题】:Kafka Streaming tasks and management of Internal state storesKafka Streaming 任务和内部状态存储的管理
【发布时间】:2020-10-09 18:15:00
【问题描述】:

假设我们在 2 台不同的机器(实例)上启动了 2 个 Streaming-Tasks,具有以下属性:-

public final static String applicationID = "StreamsPOC";
public final static String bootstrapServers = "10.21.22.56:9093";    
public final static String topicname = "TestTransaction";
public final static String shipmentTopicName = "TestShipment";
public final static String RECORD_COUNT_STORE_NAME = "ProcessorONEStore";

使用上述这些属性,流任务的定义如下所示:-

        Map<String, String> changelogConfig = new HashMap();
        changelogConfig.put("min.insyc.replicas", "1");
        // Below line not working.
        changelogConfig.put("topic", "myChangedTopicLog");
       
        StoreBuilder kvStoreBuilder = Stores.keyValueStoreBuilder(
                Stores.persistentKeyValueStore(AppConfigs.RECORD_COUNT_STORE_NAME),
                AppSerdes.String(), AppSerdes.Integer()
        ).withLoggingEnabled(changelogConfig);

        kStreamBuilder.addStateStore(kvStoreBuilder);


        KStream<String, String> sourceKafkaStream = kStreamBuilder.stream
                (AppConfigs.topicname, Consumed.with(AppSerdes.String(), AppSerdes.String()));

现在,正如我所观察到的,kafka 在后台创建了主题(用于备份内部状态存储),名称如下:- StreamsPOC-ProcessorONEStore-changelog

第一个问题是:- 两个不同的流式传输任务是否都维护内部状态存储并将其备份到同一主题?

第二个问题是 ;- 假设 Task-1 在分区 1 上接机并将 写入其本地内部状态存储,并且任务 2 开始在分区 2 上工作并说出它还将 写入其本地各自的状态存储,然后它不会引发数据被覆盖的风险,因为这两个任务都将数据备份到相同的变更日志主题?

第三个问题是:- 我如何指定自定义名称来更改日志主题?

非常感谢您的回复!

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    首先,关于术语的一些思考:“任务”一词在 Kafka Stream 中具有明确的含义,并且作为用户,您不会自己创建任务。当您的程序被执行时,Kafka Streams 会创建 tasks,它们是“独立的计算单元”并为您执行这些任务。 -- 我猜,你所说的“任务”实际上是一个KafkaStreams 客户端(即所谓的instance)。

    如果您使用相同的application.id 启动多个实例,它们将形成一个消费者组,它们将以数据并行的方式分担负载。对于状态存储,每个实例将托管存储的 shard(有时也称为分区)。所有实例都使用相同的主题,并且该主题对每个存储分片都有一个分区。从存储分片到更改日志分区存在 1:1 映射。此外,从输入主题分区到 tasks 存在 1:1 映射,并且在任务和存储分片之间存在 1:1 映射。因此,总体而言,它是一个 1:1:1:1 的映射:对于每个输入主题分区,创建一个任务,每个任务持有状态存储的一个分片,每个存储分片由更改日志主题的一个分区支持。 (即,底线是,输入主题分区的数量决定了您获得多少并行任务和存储分片,并且更改日志主题创建的分区数量与您的输入主题相同。)

    1. 是的,所有实例都使用相同的变更日志主题。
    2. 由于任务是通过分片和变更日志主题分区隔离的,因此它们不会相互覆盖。然而,任务的想法是每个任务处理不同的(非重叠)键空间,因此具有相同&lt;k1,...&gt; 的所有记录应该由相同的任务处理。当然,这条规则可能有例外,如果您的应用程序不使用非重叠键空间,则程序只会被执行(当然,根据您的业务逻辑要求,这可能是正确的或不正确的)。李>
    3. 您似乎已经这样做了:注意,您只能自定义部分更改日志主题名称:&lt;application.id&gt;-&lt;storeName&gt;-changelog -- 即,您可以选择application.idstoreName。不过,整体命名模式是硬编码的。

    【讨论】:

    • 非常感谢 Matthias 的详细回复。它真的很有帮助。
    • 嗨@AdityaGoel,如果这个答案对你有帮助,请感谢人们通过投票和/或接受答案所花费的时间。这真的会帮助其他人从你的问题和这个答案中学习。
    猜你喜欢
    • 1970-01-01
    • 2021-11-21
    • 2023-01-20
    • 2018-09-22
    • 1970-01-01
    • 2018-12-25
    • 2020-05-13
    • 1970-01-01
    • 2023-03-23
    相关资源
    最近更新 更多