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

使用Apache Beam在Flink上运行Go流处理作业遇上下文错误求助

我正在学习Apache Beam,尝试构建一个流处理示例项目,目标是从Kafka主题"word"读取数据并打印到控制台。已通过Docker在本地部署了独立集群模式的Flink(端口8081)和Kafka(端口9091)。由于缺乏清晰的示例文档,自行编写了代码,但运行时出现错误。

原代码

package main

import (
    "context"
    "flag"
    "time"

    "github.com/apache/beam/sdks/v2/go/pkg/beam"
    "github.com/apache/beam/sdks/v2/go/pkg/beam/io/xlang/kafkaio"
    "github.com/apache/beam/sdks/v2/go/pkg/beam/register"
    "github.com/apache/beam/sdks/v2/go/pkg/beam/runners/flink"
    "google.golang.org/appengine/log"
)

var (
//  expansionAddr = flag.String("expansion_addr", "",
//      "Address of Expansion Service. If not specified, attempts to automatically start an appropriate expansion service.")
//  bootstrapServers = flag.String("bootstrap_servers", "",
//      "(Required) URL of the bootstrap servers for the Kafka cluster. Should be accessible by the runner.")
//  topic = flag.String("topic", "kafka_taxirides_realtime", "Kafka topic to write to and read from.")

)

func init() {
    register.DoFn2x0[context.Context, []byte](&LogFn{})
}

// LogFn is a DoFn to log rides.
type LogFn struct{}

// ProcessElement logs each element it receives.
func (fn *LogFn) ProcessElement(ctx context.Context, elm []byte) {
    log.Infof(ctx, "Word info: %v", string(elm))
}

// FinishBundle waits a bit so the job server finishes receiving logs.
func (fn *LogFn) FinishBundle() {
    time.Sleep(2 * time.Second)
}

func main() {
    flag.Parse()
    //beam initialization
    beam.Init()
    ctx := context.Background()
    //creating pipeline object and scope
    pipeline := beam.NewPipeline()
    scope := pipeline.Root()

    //reading from kafka IO --> This is not a native support as of now for beam and golang.
    //it uses a cross-compiled library from java to acheive the kafka connector

    //defining kafka details
    brokerAddr := ""
    bootstrapServer := "bootstrap-server:kafka-1:9091"
    topic := "word"
    // input reader is our consumer which reads from the input topic, this is defined as per kafkaio docs.
    inputReader := kafkaio.Read(scope, brokerAddr, bootstrapServer, []string{topic})
    vals := beam.DropKey(scope, inputReader)
    beam.ParDo0(scope, &LogFn{}, vals)

    if _, err := flink.Execute(ctx, pipeline); err != nil {
        log.Errorf(ctx, "Failed to execute job: %v", err)
    }
}

运行报错信息

2023/09/10 01:06:49 Downloaded: C:\tmp\artifacts\beam-sdks-java-io-expansion-service-2.49.0-m4yWpU_pIZFwgP3wHCijYZg7hfO6Eg5-Dx3eXCSRTb0.jar (sha256: 9b8c96a54fe921917080fdf01c28a361983b85f3ba120e7e0f1dde5c24914dbd, size: 59370861)
panic: not an App Engine context

问题原因与修复方案

核心原因

报错的本质是代码使用了google.golang.org/appengine/log日志库,这个库要求传入App Engine专属的上下文对象,但Beam运行时提供的是标准Go context,不满足该库的校验逻辑,因此触发panic。

修复步骤

  1. 替换日志依赖:移除google.golang.org/appengine/log,改用Go标准库的log包,该包无需特殊上下文支持。
  2. 调整日志调用:标准库log的方法不需要上下文参数,直接使用log.Printf、log.Fatalf等方法即可。
  3. 验证Kafka地址:确保bootstrapServer地址能被Flink集群访问,本地Docker场景下,如果是主机访问容器,建议用localhost:9091(需提前映射容器端口到主机);如果是容器内部通信,使用容器服务名+端口。

修改后的完整代码

package main

import (
    "context"
    "flag"
    "log"
    "time"

    "github.com/apache/beam/sdks/v2/go/pkg/beam"
    "github.com/apache/beam/sdks/v2/go/pkg/beam/io/xlang/kafkaio"
    "github.com/apache/beam/sdks/v2/go/pkg/beam/register"
    "github.com/apache/beam/sdks/v2/go/pkg/beam/runners/flink"
)

var (
    expansionAddr     = flag.String("expansion_addr", "", "Address of Expansion Service. If not specified, attempts to automatically start an appropriate expansion service.")
    bootstrapServers  = flag.String("bootstrap_servers", "localhost:9091", "(Required) URL of the bootstrap servers for the Kafka cluster. Should be accessible by the runner.")
    topic             = flag.String("topic", "word", "Kafka topic to read from.")
)

func init() {
    register.DoFn2x0[context.Context, []byte](&LogFn{})
}

// LogFn is a DoFn to log received elements.
type LogFn struct{}

// ProcessElement logs each element it receives.
func (fn *LogFn) ProcessElement(ctx context.Context, elm []byte) {
    log.Printf("Word info: %v", string(elm))
}

// FinishBundle waits a bit to ensure logs are fully received.
func (fn *LogFn) FinishBundle() {
    time.Sleep(2 * time.Second)
}

func main() {
    flag.Parse()
    beam.Init()
    ctx := context.Background()

    pipeline := beam.NewPipeline()
    scope := pipeline.Root()

    // Read from Kafka topic
    inputReader := kafkaio.Read(scope, *expansionAddr, *bootstrapServers, []string{*topic})
    vals := beam.DropKey(scope, inputReader)
    beam.ParDo0(scope, &LogFn{}, vals)

    if _, err := flink.Execute(ctx, pipeline); err != nil {
        log.Fatalf("Failed to execute job: %v", err)
    }
}

额外说明

  • 恢复了原代码中注释的flag参数,支持通过命令行灵活配置Kafka地址、主题等信息
  • 如果自动启动Expansion Service失败,可以手动启动Java版Expansion Service,再通过--expansion_addr参数指定服务地址

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 04:43:15