【问题标题】:Kafka Streaming Suppress Feature to get hold of transactions which are late beyond grace periodKafka Streaming Suppress Feature 以获取延迟超过宽限期的交易
【发布时间】:2019-05-15 02:51:04
【问题描述】:

我目前正在为日窗口使用 Kafka 流式 DSL 抑制功能。我们可能会遇到这样一种情况,即某些事件可能会出现得很晚,超出了宽限期。

根据 kafka 流媒体文档,此类事件将被丢弃,这些事件不适合窗口。

请帮帮我。

1) 是否有可能在同一流程中获取此类丢弃的事件?

Apache flink 确实提供了对此类非常晚的事件的保留,并且想知道此类功能是否可用于流式传输。

2) 考虑到数以百万计的事件流经系统,使用 DSL-suppress for day window 将间歇性聚合数据保存在内存中的可行性如何?

任何时间线 kafka 流媒体社区都将很快提供 RockDB 支持,以避免由于内存不足而导致应用程序崩溃。

【问题讨论】:

    标签: apache-kafka-streams


    【解决方案1】:

    我目前正在为日窗口使用 kafka 流式 DSL 抑制功能。我们可能会遇到一些事件可能来得很晚,超出宽限期的情况。

    根据 kafka 流媒体文档,此类事件将被丢弃,这些事件不适合窗口。 [...]

    1) 是否有可能在同一流程中获取此类丢弃的事件?

    您需要延长宽限期。宽限期的重点是允许您定义您可以接受(非常)迟到的事件到达的时间。宽限期实际上可能比窗口大小长——我提到这一点是因为您提到“不适合窗口”。

    在我看来,您似乎接受了迟到的事件,但您不想增加宽限期。为什么?

    Apache flink 确实提供了对此类非常晚的事件的保留,并且想知道此类功能是否可用于流式传输。

    如果您的意思是:对于 Kafka Streams 中的此类非常晚的事件,是否有类似回调之类的东西,那么答案是否定的,没有。

    2) 考虑到数百万个事件流经系统,使用 DSL-suppress for day window 将间歇性聚合数据保存在内存中的可行性如何?

    任何时间线 kafka 流媒体社区都将很快提供 RockDB 支持,以避免由于内存不足而导致应用程序崩溃。

    对于其他读者:RocksDB 已经得到支持,并且是 Kafka Streams 中所有有状态操作的默认状态存储引擎。唯一的例外是 Supress() 功能的当前实现,其中抑制缓冲区尚未通过 RocksDB 维护。

    关于您的问题:KAFKA-7224: Add spill-to-disk for Suppression 的工作正在进行中,但确切的预计到达时间尚不清楚。

    【讨论】:

    • 第 1 部分:规划 2 种类型的用例进行统计计算。 1) 实时 - 按预期工作。 2) 按需延迟交易,可能是较早日期的文件馈送,可能是一周前的。交易提要 --> 输入主题 --> 窗口流应用程序 --> 输出主题 --> CASSANDRA
    • 第 2 部分:不希望将宽限期保留更长的时间,因为手动馈送可能会延迟一周。计划是保留 1 天的窗口和 15 分钟的宽限期,以便按时发出并在 Cassandra 中提供统计数据。
    • 第 3 部分:如果我获得了一种机制来验证特定事务(延迟提要)在流处理期间被当前窗口拒绝,那么我可以直接从 cassandra 获取现有统计信息 + 计算内存中的摘要 + 重新添加到 Cassandra )。 OUTPUT 主题不适用于历史交易统计。主要议程是为历史交易提要使用相同的实时设置。希望我能够澄清。请指导。
    猜你喜欢
    • 2017-11-04
    • 1970-01-01
    • 1970-01-01
    • 2018-03-29
    • 1970-01-01
    • 2017-07-05
    • 1970-01-01
    • 2020-10-29
    • 1970-01-01
    相关资源
    最近更新 更多