You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在DirectRunner上运行Golang编写的连接PubSub的Apache Beam Pipeline

在Golang中用DirectRunner运行连接PubSub的Apache Beam流式Pipeline

Golang版的Apache Beam DirectRunner不需要显式设置流式处理选项——它会自动识别数据源的类型(比如PubSub这种流式源),并切换到流式模式运行。以下是具体的实现步骤和示例:

  • 配置PubSub输入源
    使用Beam Go SDK的pubsub包创建流式输入,比如从订阅或主题读取:

    import (
      "context"
      "github.com/apache/beam/sdks/v2/go/pkg/beam"
      "github.com/apache/beam/sdks/v2/go/pkg/beam/io/pubsub"
      "github.com/apache/beam/sdks/v2/go/pkg/beam/runners/direct"
    )
    
    func main() {
      // 初始化Pipeline
      beam.Init()
      p := beam.NewPipeline()
      s := p.Root()
    
      // 从PubSub订阅读取流式数据
      sub := "projects/your-project/subscriptions/your-sub"
      input := pubsub.NewSubscription(s, sub)
    
      // 添加你的数据处理逻辑,比如打印消息内容
      beam.ParDo0(s, func(msg []byte) {
        println("Received message:", string(msg))
      }, input)
    
      // 用DirectRunner运行Pipeline
      if err := direct.Execute(context.Background(), p); err != nil {
        panic(err)
      }
    }
    
  • 关键注意事项

    • 确保本地已配置GCP凭据:可以通过设置环境变量GOOGLE_APPLICATION_CREDENTIALS指向你的服务账号密钥文件,或者执行gcloud auth application-default login完成本地认证。
    • DirectRunner运行流式Pipeline时会持续运行,直到你手动终止进程(比如Ctrl+C)。
    • 如果需要调整流式运行的参数(比如窗口触发规则、水印设置),直接通过Beam的窗口API配置即可,无需额外的Runner专属选项。

内容的提问来源于stack exchange,提问作者Tomoya Ishizuka

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.17 13:15:33