Hacktoberfest 2026:维护者为十月标记出来的 issue,仍然开放、适合新手。 浏览 Hacktoberfest issue

Proposal: Watermarking

未关闭
#354 4 条评论 1 个 reaction 已指派 0 人 在 GitHub 查看

还没有人认领这个 Issue。

评估

难度
5/5
预计耗时
一周以上
新手友好度
25/100
Issue 类型
功能
描述清晰度
需要澄清
活跃度
停滞
技术栈
python

调研方向

Start by inspecting the Streamz Dataframe entry points for windowing and aggregation operations; the issue names no files or tests. Clarify how the event-time column and watermark threshold should be configured, how late events are handled, and which aggregation behaviors define done. Use the timestamp table as the basis for acceptance tests.

由索引模型根据 Issue 内容生成。

描述

This is a feature that exists in other streaming data system with windowing and aggregation functions. Its purpose is to support late arriving data in a window or aggregation. It will require that Streamz becomes semi-aware of the data structure because we will need to specify a column that represents the event time. So, for clarity, there are two timestamps here:

  1. The time at which the event enters the data pipeline
  2. The time at which the event is created in the source. We'll call this the "event time".

Due to various latencies in a distributed system, an event that should be included into an aggregation arrives too late into the pipeline to be counted. As an example, if you have a window of 5 minutes, but an event that has an event time within those 5 minutes arrives 3 minutes after the window closes, it will not be included in any aggregations.

What watermarking will do is keep the window open for a specified amount of time to include all of the data. So, if we have a 5 minute window and a watermarking threshold of 5 minutes, the window will include all events in the first 5 minutes and all events in the second 5 minutes if the event time belongs to the previous 5 minutes. If the event time is outside of the window, it will be dropped. This may be the key to implementing this, because we may just be able to include all data and then just drop data that is outside of the watermark threshold.

Here is an example of windows of 5 seconds and a watermark threshold of 5 seconds.

Arrives Event Time Included
00:01:01 00:01:01 Yes - Is inside of the window time
00:01:02 00:01:01 Yes - Is inside of the window time
00:01:02 00:01:02 Yes - Is inside of the window time
00:01:03 00:01:02 Yes - Is inside of the window time
00:01:04 00:01:04 Yes - Is inside of the window time
00:01:06 00:01:04 Yes - Is with-in watermark threshold
00:01:09 00:01:04 Yes - Is with-in watermark threshold
00:01:11 00:01:04 No - Arrived too late
00:01:12 00:01:12 No - Is outside of window

I've been spending the last few days trying to figure out where this would fit into Streamz because it seems like Streamz doesn't determine what gets included in a batch. So, I'm thinking this could be implemented a few places.

My current thinking is that the Streamz Dataframe would need new parameters for the watermark threshold time and the event time column. And, when operations like windowing or aggregations are performed, it would take into account the watermarking threshold.

As always, feedback is greatly appreciated here. I'd like to include something like this so that it works for everyone.

Also, let me know if this explanation isn't clear.

主要语言
Python
星标
1.3k
派生
151
平均合并
17 小时 39 分钟
30 天内合并 PR
1

环境准备

  • 提供 Dockerfile 或 Docker Compose 文件
  • 没有 Pull Request 模板
  • 阅读贡献指南

从这里开始

  1. 先读完整个 Issue,再读项目的贡献指南。
  2. 在 Issue 下留言说明你要接手 —— 这能避免两个人做同样的事。
  3. Fork 仓库,在一个分支上完成修改。
  4. 提交 Pull Request,并在描述里引用这个 Issue 编号。

python-streamz/streamz 的其他 Issue

查看 python-streamz/streamz 的全部 Issue

相似的 Issue

更多 Python Issue

把新 issue 发到你的邮箱

精选适合新手参与的 GitHub issue 摘要。