【问题标题】:Kafka streams co-partitioning vs interactive query卡夫卡流共同分区与交互式查询
【发布时间】:2020-10-04 18:40:58
【问题描述】:

我有以下拓扑:

topology.addSource(WS_CONNECTION_SOURCE, new StringDeserializer(), new WebSocketConnectionEventDeserializer()
            , utilService.getTopicByType(TopicType.CONNECTION_EVENTS_TOPIC))
            .addProcessor(SESSION_PROCESSOR, WSUserSessionProcessor::new, WS_CONNECTION_SOURCE)
            .addStateStore(sessionStoreBuilder, SESSION_PROCESSOR)
            .addSink(WS_STATUS_SINK, utilService.getTopicByType(TopicType.ONLINE_STATUS_TOPIC),
                    stringSerializer, stringSerializer
                    , SESSION_PROCESSOR)

            //WS session routing
            .addSource(WS_NOTIFICATIONS_SOURCE, new StringDeserializer(), new StringDeserializer(),
                    utilService.getTopicByType(TopicType.NOTIFICATION_TOPIC))
            .addProcessor(WS_NOTIFICATIONS_ROUTE_PROCESSOR, SessionRoutingEventGenerator::new,
                    WS_NOTIFICATIONS_SOURCE)
            .addSink(WS_NOTIFICATIONS_DELIVERY_SINK, new NodeTopicNameExtractor(), WS_NOTIFICATIONS_ROUTE_PROCESSOR)
            .addStateStore(userConnectedNodesStoreBuilder, WS_NOTIFICATIONS_ROUTE_PROCESSOR, SESSION_PROCESSOR);  

如您所见,有 2 个源主题。状态存储是从第一个主题构建的,第二个流程读取状态存储。当我开始拓扑时,我看到这些流线程被分配了两个源主题的相同分区(共同分区)。我认为这是因为状态存储是由第二个主题流访问的。

这在功能上运行良好。但是有一个性能问题。当第一个源主题的输入数据量激增(更新状态存储)时,第二个主题的处理就会延迟。

对我来说,第二个主题应该尽快处理。处理第一个主题的延迟很好。

我正在考虑以下策略:

Current configuration:
     WS_CONNECTION_SOURCE - 30 partitions
     WS_NOTIFICATIONS_SOURCE - 30 partitions
     streamThreads: 10
     appInstances: 3 

New configuration:
    WS_CONNECTION_SOURCE - 15 partitions
    WS_NOTIFICATIONS_SOURCE - 30 partitions
    streamThreads: 10
    appInstances: 3
    Since there is no co-partitioning, tasks has to use interactive query to access store

思路是在10个线程中,5个线程只处理第二个主题,这样可以在第一个主题激增时缓解当前的问题。

这是我的问题:

1. Is this strategy correct? To avoid co-partitioning and use interactive query
2. Is there a chance that Kafka will assign 10 partitions of WS_CONNECTION_SOURCE 
   to one instance since there are 10 threads and one instance won't get any?
3. Is there any better approach to solve the performance problem?

【问题讨论】:

    标签: apache-kafka apache-kafka-streams partitioning


    【解决方案1】:

    状态存储和交互式查询是 Kafka Streams 的抽象。 要使用交互式查询,您必须定义状态存储(使用 Kafka Streams API)并强制您对输入主题具有相同数量的分区。 我认为您的解决方案行不通。交互式查询用于公开在 Kafka 流之外查询状态存储的能力(不适用于处理器 API 中的访问)

    也许您可以查看您的 SESSION_PROCESSOR 源代码并从其他拓扑中提取更多工作到 Process 并将结果发布到中间主题,然后基于该构建该状态存储。

    另外:

    目前 Kafka Streams 不支持输入主题的优先级。有关于源主题优先级的 KIP:KIP-349。不幸的是,链接的 Jira 票证已关闭,因为不会修复 (https://issues.apache.org/jira/browse/KAFKA-6690)

    【讨论】:

    • 我尝试了我提到的方法(不同的分区和交互式查询),它非常慢。看起来交互式查询不适用于生产用途。现在我回到原始解决方案,源主题的分区数相等。两个处理器都有不同的任务,不能合并。 SESSION_PROCESSOR 写入状态存储,NOTIFICATION_ROUTE_PROCESSOR 读取状态存储并执行不同类型的任务
    • 将更多工作从一个处理器转移到另一个处理器如何改善这种情况?
    • @cppcoder,我不是说合并那些处理器。而是将SESSION_PROCESSOR 拆分为两个。一个带有heavy 计算并将结果发布到中间主题,另一个lightweight 读取它并仅更新state.store。繁重的可以作为不同的应用程序运行,并且不会影响lightweight one 的性能
    • 明白你的意思。但截至目前,SESSION_PROCESSOR 所做的只是检查现有状态存储并根据当前记录更新存储中的值。但是,它确实将新事件转发到另一个主题。会不会很重?
    • 我在处理器中看到这个System.currentTimeMillis() - context.timestamp() 的值异常高。为什么它在处理器本身内部接收缓慢?处理延迟小于 100 毫秒。但上述延迟以分钟为单位显示
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-01-15
    • 2021-03-01
    • 1970-01-01
    • 2023-01-03
    相关资源
    最近更新 更多