如何可视化DirectRunner本地运行的Beam流水线执行过程
本地DirectRunner运行Beam流水线的执行图可视化说明
首先明确结论:本地通过DirectRunner运行Beam流水线时,无法获得和GCP Dataflow环境完全一致的带实时运行指标、托管式交互的执行图页面,但可以通过SDK内置能力导出结构完全一致的流水线拓扑图查看。
- GCP环境下的可视化执行图是Dataflow Runner配套的托管能力:运行流水线时Dataflow服务会自动采集流水线拓扑、各步骤吞吐量/耗时、worker运行状态等数据,实时渲染到云控制台的交互页面,不需要用户做额外配置。
- DirectRunner本身没有自带托管式Web可视化界面,默认运行时不会自动生成可交互的执行图页面。
- 你可以通过Beam SDK内置的图渲染能力,导出完整的流水线执行拓扑:
- 流水线构建完成后,调用SDK内置的渲染方法,将流水线的有向无环图(DAG)导出为标准DOT格式文件
- 本地通过Graphviz等工具将DOT文件渲染为图片,即可查看所有转换步骤、PCollection流转关系,图的结构逻辑和GCP控制台展示的静态执行拓扑完全一致,仅缺少运行时动态指标、worker状态等云端专属的监控数据
- 以Python SDK为例,导出执行图的参考代码如下:
import apache_beam as beam from apache_beam.runners.direct import DirectRunner from apache_beam.runners.render import render # 初始化DirectRunner流水线 p = beam.Pipeline(runner=DirectRunner()) # 自定义流水线处理逻辑 _ = ( p | beam.Create(["gcp", "beam", "direct runner"]) | beam.Map(lambda x: x.upper()) | beam.Map(print) ) # 导出执行拓扑为DOT格式文件 render(p, output_path="local_pipeline.dot") # 运行流水线 p.run()
- Java SDK同样提供了类似的拓扑导出接口,导出的DOT文件可以用相同的本地工具渲染查看。如果不需要图形化展示,也可以调整日志级别为
DEBUG,运行时会在日志中按执行顺序打印所有步骤的文本信息。
注意:本地导出的静态执行图仅展示流水线的逻辑拓扑,不会同步展示DirectRunner运行时的动态指标、数据分片分布、错误栈定位等云端可视化页面带有的交互调试能力。
内容的提问来源于stack exchange,提问作者Dimon Buzz
相关产品推荐
相关产品推荐

