如何在Beam YAML管道中用CSV的timestamp作为行时间戳做窗口聚合
如何在Beam YAML中使用CSV的timestamp字段进行窗口聚合
要让窗口聚合基于CSV中的timestamp字段而非默认的处理时间,你需要在读取CSV后添加WithTimestamps转换,将字符串格式的时间戳解析为Beam可识别的时间戳类型。以下是修改后的完整管道配置:
pipeline: type: chain transforms: - type: ReadFromCsv config: path: data/people.csv # 新增:将CSV的timestamp字段解析为行时间戳 - type: WithTimestamps config: timestamp: parse_timestamp(timestamp, "yyyy-MM-dd'T'HH:mm:ss") - type: WindowInto windowing: type: fixed size: 10s # 注意添加时间单位,这里是10秒窗口 - type: Combine config: group_by: country combine: total: value: age fn: sum - type: LogForTesting - type: WriteToJson config: path: ~/out.json options: yaml_experimental_features: Combine
关键说明:
WithTimestamps转换:使用Beam内置的parse_timestamp函数,第一个参数是CSV中的timestamp字段名,第二个参数是该字段的时间格式(与你的CSV格式完全匹配),完成字符串到时间戳的转换。- 窗口大小单位:原配置中的
size:10必须指定单位(如10s表示10秒、5m表示5分钟),否则Beam无法识别窗口时长。 - 窗口聚合逻辑:后续的
Combine操作会基于转换后的行时间戳,对每个10秒窗口内的同国家数据进行年龄求和。
你的CSV数据中timestamp字段采用ISO8601格式,上述parse_timestamp的格式字符串可以正确解析该字段。
内容的提问来源于stack exchange,提问作者By1
相关产品推荐
相关产品推荐

