ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Apache Beam Java 实战:使用 WithTimestamps 为 PCollection 元素添加时间戳

Apache Beam Java 实战:使用 WithTimestamps 为 PCollection 元素添加时间戳 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载导读Apache Beam 的窗口Windowing机制依赖元素的时间戳来划分窗口、计算水位线Watermark与处理迟到数据但并非所有数据源都会自带时间戳像TextIO这类有界源Bounded Source读取到的文件内容在进入PCollection时并不携带任何时间信息。本文以 Beam Katas 训练营中WithTimestamps练习task.md为主线讲解如何用WithTimestampstransform 为元素批量指派时间戳并深入 SDK 源码剖析其底层实现同时给出可复现的完整代码与测试验证方案帮助你为后续的固定窗口、迟到数据处理打下基础。为什么需要手动添加时间戳Beam 官方编程指南明确指出一个关键事实有界源Bounded Sources不会为元素提供时间戳。典型例子就是TextIO.read()读取本地文件或 GCS 文件时每一行文本只作为普通字符串进入管道Beam 无法得知该行业务上对应什么时刻。这在做事件类数据处理时是致命的——例如日志分析、订单流统计每条记录的真实发生时间往往记录在数据内容本身如 JSON 字段、CSV 列而非文件名中。具体场景如下读取文件后需要按业务时间事件发生时间而不是到达时间做窗口聚合使用Window.into(FixedWindows.of(...))等窗口函数时Beam 需要根据元素时间戳计算该元素落入哪个窗口WithTimestamps.java 的类注释明确说明时间戳用于将元素分配到BoundedWindow流式计算中水位线推进、迟到数据判定都以元素时间戳为基础。因此当数据源不提供时间戳而业务又需要时间维度时必须在进入PCollection后显式补上时间戳。这也是 Beam Katas 中 Adding Timestamp 这一课要解决的训练目标。Kata 任务基于 Event.getDate() 为元素指派时间戳任务目标原文档给出了明确的练习要求Kata:Please assign each element a timestamp based on theEvent.getDate().即读入一批Event对象将其内部字段dateJoda-Time 的DateTime类型转换为元素时间戳使得每个元素携带其业务发生时间。输入数据结构练习提供了Event数据模型Event.javaimport java.io.Serializable; import java.util.Objects; import org.joda.time.DateTime; public class Event implements Serializable { private String id; private String event; private DateTime date; public Event(String id, String event, DateTime date) { ... } public String getId() { return id; } public String getEvent() { return event; } public DateTime getDate() { return date; } // 含 equals / hashCode / toString }三个字段分别表示事件 ID、事件名称和业务发生时间date使用 Joda-Time 的DateTimedate正是我们要提取为时间戳的字段。标准解法WithTimestamps练习给出的标准答案Task.java非常简洁核心只有一行static PCollectionEvent applyTransform(PCollectionEvent events) { return events.apply(WithTimestamps.of(event - event.getDate().toInstant())); }关键点解析WithTimestamps.of(SerializableFunctionT, Instant fn)接收一个从T到Instant的映射函数对PCollection中每个元素求值将其结果作为该元素的新时间戳注意返回类型是InstantJoda-Time因此需要调用event.getDate().toInstant()把DateTime转换为InstantWithTimestamps是一个PTransformPCollectionT, PCollectionT即输入输出是同一个元素类型只改变元素附带的时间戳元数据不改变元素内容本身。完整的可运行主程序如下main方法部分public class Task { public static void main(String[] args) { PipelineOptions options PipelineOptionsFactory.fromArgs(args).create(); Pipeline pipeline Pipeline.create(options); PCollectionEvent events pipeline.apply( Create.of( new Event(1, book-order, DateTime.parse(2019-06-01T00:00:0000:00)), new Event(2, pencil-order, DateTime.parse(2019-06-02T00:00:0000:00)), new Event(3, paper-order, DateTime.parse(2019-06-03T00:00:0000:00)), new Event(4, pencil-order, DateTime.parse(2019-06-04T00:00:0000:00)), new Event(5, book-order, DateTime.parse(2019-06-05T00:00:0000:00)) ) ); PCollectionEvent output applyTransform(events); output.apply(Log.ofElements()); pipeline.run(); } static PCollectionEvent applyTransform(PCollectionEvent events) { return events.apply(WithTimestamps.of(event - event.getDate().toInstant())); } }管道流程为Create.of(events)构造 5 条测试事件 →applyTransform用WithTimestamps补时间戳 →Log.ofElements()来自 Log.java 的日志工具打印带时间戳的元素 →pipeline.run()执行。输入数据与预期结果输入 5 条事件的date分别为 2019-06-01 至 2019-06-05每天 00:00:00UTC0。执行WithTimestamps后每个元素的时间戳即为对应日期例如元素(1, book-order)的时间戳为2019-06-01T00:00:00.000Z。这些时间戳将直接影响后续窗口归属若后续应用 1 天的固定窗口5 个事件会恰好落入 5 个不同的窗口。另一种等价实现ParDo outputWithTimestamp原文档同时提示也可以应用一个 ParDo transform输出带设定时间戳的新元素。同目录下的姊妹练习ParDo/task.md给出了这种实现static PCollectionEvent applyTransform(PCollectionEvent events) { return events.apply(ParDo.of(new DoFnEvent, Event() { ProcessElement public void processElement(Element Event event, OutputReceiverEvent out) { out.outputWithTimestamp(event, event.getDate().toInstant()); } })); }两种方式的对比维度WithTimestampsParDo outputWithTimestamp代码量一行 Lambda 即可需要自定义 DoFn 与 ProcessElement适用场景仅需整体重打时间戳需要在处理逻辑中按元素逐个、条件化设置时间戳元素内容不修改仅改元数据可输出新元素或改变元素推荐度首选语义更清晰需要精细控制时使用从源码看WithTimestamps本质上也是基于ParDo封装的其expand方法内部调用input.apply(AddTimestamps, ParDo.of(new AddTimestampsDoFn(fn, allowedTimestampSkew)))WithTimestamps.java可见两者是高层便捷 API 与底层原语的关系。源码深挖WithTimestamps 的底层实现为了理解其行为边界我们直接阅读核心 SDK 源码WithTimestamps.java工厂方法of(SerializableFunctionT, Instant fn)要求函数必须可序列化SerializableFunction因为 Beam 管道会序列化到各 Worker 执行fn为 null 时会通过checkNotNull直接抛出异常第 80 行。默认允许的时间戳偏移构造时allowedTimestampSkew默认为Duration.ZERO第 71 行即新时间戳只能向后未来移动不能早于输入元素的原始时间戳。若回拨超过允许偏移执行时会抛出IllegalArgumentException。回拨开关已废弃withAllowedTimestampSkew(Duration)可放宽回拨上限new Duration(Long.MAX_VALUE)表示无限偏移第 84-99 行。但源码以Deprecated标记并明确警告允许元素落后于水位线会导致其被判定为迟到数据若超过下游Window.withAllowedLateness的容忍范围可能被静默丢弃。日常开发应避免使用该 API。null 时间戳若映射函数对某元素返回 null执行时会抛出NullPointerException第 60-61 行注释。窗口保持不变输出元素仍保留输入元素所在的窗口若希望按新时间戳重新划分窗口需要再应用一次Window.into(WindowFn)第 57-58 行注释。这些细节解释了为什么用WithTimestamps要保证date字段非空、且转换出的时间一般晚于元素原始时间等实践约束。用 PAssert 验证时间戳是否生效练习自带单元测试TaskTest.java它展示了验证时间戳的标准手法值得在实际项目中复用PCollectionKVEvent, Instant timestampedResults results.apply(KVEvent, Instant, ParDo.of(new DoFnEvent, KVEvent, Instant() { ProcessElement public void processElement(Element Event event, ProcessContext context, OutputReceiverKVEvent, Instant out) { out.output(KV.of(event, context.timestamp())); } }) ); PAssert.that(results).containsInAnyOrder(events); // 元素内容保持不变 PAssert.that(timestampedResults) .containsInAnyOrder( KV.of(events.get(0), events.get(0).getDate().toInstant()), KV.of(events.get(1), events.get(1).getDate().toInstant()), KV.of(events.get(2), events.get(2).getDate().toInstant()), KV.of(events.get(3), events.get(3).getDate().toInstant()), KV.of(events.get(4), events.get(4).getDate().toInstant()) ); testPipeline.run().waitUntilFinish();验证思路分两步通过context.timestamp()在DoFn中取回每个元素的当前时间戳与Event.getDate().toInstant()逐一比对用PAssert.containsInAnyOrder断言两者完全一致用PAssert.that(results).containsInAnyOrder(events)确认WithTimestamps只改变元数据、不改变元素内容。测试基于TestPipelineTestPipeline.create()运行无需真实 Runner可在gradlew下直接执行验证。实践要点与常见误区结合原文档、练习源码与 SDK 实现总结以下实战要点有界源无时间戳凡是TextIO、AvroIO等批量读取的场景元素默认时间戳为纪元开始或数据源定义值做事件时间窗口前必须先补时间戳时间戳字段应尽早提取建议在管道入口读入后第一个 transform就完成WithTimestamps避免中间 transform 丢失上下文保持窗口语义一致WithTimestamps不改写窗口归属需要按新时间戳分窗时紧随其后应用Window.into(...)固定窗口练习见 Fixed Time Window 目录避免时间戳回拨默认Duration.ZERO禁止回拨回拨元素会成为迟到数据甚至被丢弃不要轻易动用已废弃的withAllowedTimestampSkew注意空值与类型映射函数返回 null 会直接导致运行失败DateTime.toInstant()与Instant的转换是 Joda-Time 中的常见操作务必确认数据字段非空。小结通过本文你已完整掌握在 Apache Beam Java SDK 中为PCollection元素添加时间戳的标准方案理解有界源不带时间戳这一前提用一行WithTimestamps.of(event - event.getDate().toInstant())完成任务并知道其底层由ParDoAddTimestampsDoFn实现、默认不允许时间戳回拨等边界。在此基础上可以进一步探索同一课程下的ParDo手动实现以及后续的固定窗口Fixed Time Window 练习逐步建立起对 Beam 事件时间模型的完整认知。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam Kotlin Katas 实战用 WithTimestamps 为 PCollection 元素添加时间戳Apache Beam Kotlin Katas 实战用 WithTimestamps 为 PCollection 元素添加时间戳 导读 在 Apache B大数据批处理流处理数据工程Apache Beam Katas 实战用 ParDo 为 PCollection 元素添加时间戳Apache Beam Katas 实战用 ParDo 为 PCollection 元素添加时间戳 导读 本文以 Apache Beam 官方 Katas大数据批处理流处理数据工程Apache Beam Kotlin 实战用 ParDo 为 PCollection 元素添加时间戳Adding TimestampApache Beam Kotlin 实战用 ParDo 为 PCollection 元素添加时间戳Adding Timestamp 导读 在 Apach大数据批处理流处理数据工程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表