如何在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专属选项。
- 确保本地已配置GCP凭据:可以通过设置环境变量
内容的提问来源于stack exchange,提问作者Tomoya Ishizuka
相关产品推荐
相关产品推荐

