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

如何将本地Apache Beam Kafka管道部署到GCP Dataflow?

Apache Beam管道部署GCP Dataflow失败问题

问题描述

我编写了一个完整的Apache Beam管道,用于订阅和发布Kafka主题,本地运行完全正常。但尝试部署为GCP Dataflow作业时,按官方文档操作后并未创建Dataflow作业,代码仅在本地执行。查阅资料后得知可能需要创建灵活模板,但不确定只需创建模板还是需要修改代码。


管道代码

public static void iot_topic_connection(String IP) {
    System.out.println( "Initiating connection with iot topic" );
    Pipeline pipeline = Pipeline.create();
    PCollection<KafkaRecord<String, String>> pCollectionA = pipeline.apply(KafkaIO.<String, String>read()
            .withBootstrapServers(IP)
            .withTopic("iotA")
            .withKeyDeserializer(StringDeserializer.class)
            .withValueDeserializer(StringDeserializer.class)
    );
    PCollection<KafkaRecord<String, String>> pCollectionB = pipeline.apply(KafkaIO.<String, String>read()
            .withBootstrapServers(IP)
            .withTopic("iotB")
            .withKeyDeserializer(StringDeserializer.class)
            .withValueDeserializer(StringDeserializer.class)
    );

    //mapping and windowing
    PCollection<KV<String, String>> wrdA = pCollectionA
            .apply(
            MapElements.into(TypeDescriptors.kvs(TypeDescriptors.strings(), TypeDescriptors.strings()))
                    .via((KafkaRecord<String, String> record) -> KV.of(record.getKV().getKey(), splitValue(record.getKV().getValue(),0)))
            ).apply(
                    Window.<KV<String,String>>into(FixedWindows.of(Duration.standardSeconds(30)))
                            .triggering(Repeatedly.forever(AfterWatermark.pastEndOfWindow()))
                            .withAllowedLateness(Duration.ZERO).accumulatingFiredPanes()
            );

    PCollection<KV<String, String>> wrdB = pCollectionB
            .apply(
            MapElements.into(TypeDescriptors.kvs(TypeDescriptors.strings(), TypeDescriptors.strings()))
                    .via((KafkaRecord<String, String> record) -> KV.of(record.getKV().getKey(), splitValue(record.getKV().getValue(),0)))
            ).apply(
                    Window.<KV<String,String>>into(FixedWindows.of(Duration.standardSeconds(30)))
                            .triggering(Repeatedly.forever(AfterWatermark.pastEndOfWindow()))
                            .withAllowedLateness(Duration.ZERO).accumulatingFiredPanes()
            );

    TupleTag<String> wrdA_tag = new TupleTag<>();
    TupleTag<String> wrdB_tag = new TupleTag<>();
    //CoGroupbyKey operation =~ join
    PCollection<KV<String, CoGbkResult>> results =
            KeyedPCollectionTuple.of(wrdA_tag, wrdA)
                    .and(wrdB_tag, wrdB)
                    .apply(CoGroupByKey.create());

    //here aggregation needs to be implemented
    PCollection<String> final_data =
            results.apply(
                    ParDo.of(
                            new DoFn<KV<String, CoGbkResult>, String>() {
                                @ProcessElement
                                public void processElement(ProcessContext c) {
                                    Float avg_temp = 0.0f;
                                    Integer array_lenght = 1;
                                    KV<String, CoGbkResult> e = c.element();
                                    //System.out.println(e);
                                    String key = e.getKey();
                                    Iterable<String> wrdA_obj = e.getValue().getAll(wrdA_tag);
                                    Iterable<String> wrdB_obj = e.getValue().getAll(wrdB_tag);
                                    Iterator<String> wrdA_iter = wrdA_obj.iterator();
                                    Iterator<String> wrdB_iter = wrdB_obj.iterator();
                                    while( wrdA_iter.hasNext() && wrdB_iter.hasNext() ){
                                        // Process event1 and event2 data and write to c.output
                                        String v1 = wrdA_iter.next();
                                        String v2 = wrdB_iter.next();
                                        avg_temp = (Float.parseFloat(v1) + Float.parseFloat(v2));
                                        array_lenght += 1;
                                    }
                                    if(avg_temp == 0.0f){
                                        System.out.println("Unable to join event1 and event2");
                                    }else{
                                        avg_temp = avg_temp/array_lenght;
                                        c.output(String.valueOf(avg_temp));
                                    }

                                }
                            }));
    //write to kafka topic
    final_data.apply(KafkaIO.<Void, String>write()
            .withBootstrapServers(IP)
            .withTopic("iotOut")
            .withValueSerializer( StringSerializer.class).values());
    //Here we are starting the pipeline
    pipeline.run();
}

pom.xml配置

<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
  xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
  <modelVersion>4.0.0</modelVersion>

  <groupId>org.example</groupId>
  <artifactId>beam-app</artifactId>
  <version>1.0-SNAPSHOT</version>
  <packaging>jar</packaging>

  <name>beam-app</name>
  <url>http://maven.apache.org</url>

  <properties>
    <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
    <maven.compiler.source>11</maven.compiler.source>
    <maven.compiler.target>11</maven.compiler.target>
  </properties>
  <profiles>
    <profile>
      <id>dataflow-runner</id>
      <!-- Makes the DataflowRunner available when running a pipeline. -->
      <dependencies>
        <dependency>
          <groupId>org.apache.beam</groupId>
          <artifactId>beam-runners-google-cloud-dataflow-java</artifactId>
          <version>2.41.0</version>
          <scope>runtime</scope>
        </dependency>
      </dependencies>
    </profile>
  </profiles>
  <dependencies>
    <dependency>
      <groupId>org.apache.beam</groupId>
      <artifactId>beam-runners-google-cloud-dataflow-java</artifactId>
      <version>2.41.0</version>
      <scope>runtime</scope>
    </dependency>
    <dependency>
      <groupId>org.apache.beam</groupId>
      <artifactId>beam-sdks-java-io-google-cloud-platform</artifactId>
      <version>2.41.0</version>
    </dependency>
    <!-- https://mvnrepository.com/artifact/org.apache.beam/beam-sdks-java-extensions-join-library -->
    <dependency>
      <groupId>org.apache.beam</groupId>
      <artifactId>beam-sdks-java-extensions-join-library</artifactId>
      <version>2.41.0</version>
    </dependency>

    <dependency>
      <groupId>org.apache.beam</groupId>
      <artifactId>beam-sdks-java-core</artifactId>
      <version>2.41.0</version>
    </dependency>
    <dependency>
      <groupId>org.apache.beam</groupId>
      <artifactId>beam-sdks-java-io-kafka</artifactId>
      <version>2.41.0</version>
    </dependency>
    <dependency>
      <groupId>org.apache.beam</groupId>
      <artifactId>beam-runners-direct-java</artifactId>
      <version>2.41.0</version>
    </dependency>
    <!-- https://mvnrepository.com/artifact/org.apache.kafka/kafka-clients -->
    <dependency>
      <groupId>org.apache.kafka</groupId>
      <artifactId>kafka-clients</artifactId>
      <version>3.2.2</version>
    </dependency>

    <dependency>
      <groupId>junit</groupId>
      <artifactId>junit</artifactId>
      <version>4.13.2</version>
      <scope>test</scope>
    </dependency>
  </dependencies>
</project>

执行的Maven命令

mvn compile exec:java \
  -Dexec.mainClass=org.example.App \
  -Dexec.args=" \
  --project=round...2 \
  --runner=DataflowRunner \
  --jobName=my-first-job \
  --region=southamerica-west1 \
  --streaming=true \
  --zone=southamerica-west1-a \
  --tempLocation=gs://beam-testing-app/dataflow/temp \
  --gcpTempLocation=gs://b..p/dataflow/temp \
  --stagingLocation=gs://b...p/dataflow/staging \
  " \
  -Pdataflow-runner

控制台输出

[INFO] Scanning for projects...
[INFO]
[INFO] ------------------------< org.example:beam-app >------------------------
[INFO] Building beam-app 1.0-SNAPSHOT
[INFO] --------------------------------[ jar ]---------------------------------
Downloading from central: https://repo.maven.apache.org/maven2/io/grpc/grpc-api/maven-metadata.xml
Downloaded from central: https://repo.maven.apache.org/maven2/io/grpc/grpc-api/maven-metadata.xml (3.0 kB at 6.0 kB/s)
Downloading from central: https://repo.maven.apache.org/maven2/io/grpc/grpc-core/maven-metadata.xml
Downloaded from central: https://repo.maven.apache.org/maven2/io/grpc/grpc-core/maven-metadata.xml (4.6 kB at 89 kB/s)
Downloading from central: https://repo.maven.apache.org/maven2/io/grpc/grpc-netty-shaded/maven-metadata.xml
Downloaded from central: https://repo.maven.apache.org/maven2/io/grpc/grpc-netty-shaded/maven-metadata.xml (3.7 kB at 89 kB/s)
Downloading from central: https://repo.maven.apache.org/maven2/io/grpc/grpc-alts/maven-metadata.xml
Downloaded from central: https://repo.maven.apache.org/maven2/io/grpc/grpc-alts/maven-metadata.xml (3.5 kB at 77 kB/s)
Downloading from central: https://repo.maven.apache.org/maven2/io/grpc/grpc-xds/maven-metadata.xml
Downloaded from central: https://repo.maven.apache.org/maven2/io/grpc/grpc-xds/maven-metadata.xml (2.5 kB at 50 kB/s)
[INFO]
[INFO] --- maven-resources-plugin:2.6:resources (default-resources) @ beam-app ---
[INFO] Using 'UTF-8' encoding to copy filtered resources.
[INFO] skip non existing resourceDirectory /home/user/beam-app/src/main/resources
[INFO]
[INFO] --- maven-compiler-plugin:3.1:compile (default-compile) @ beam-app ---
[INFO] Changes detected - recompiling the module!
[INFO] Compiling 2 source files to /home/user/beam-app/target/classes
[INFO]
[INFO] --- exec-maven-plugin:3.1.0:java (default-cli) @ beam-app ---
Initiating connection with iot topic
SLF4J: Failed to load class "org.slf4j.impl.StaticLoggerBinder".
SLF4J: Defaulting to no-operation (NOP) logger implementation
SLF4J: See http://www.slf4j.org/codes.html#StaticLoggerBinder for further details.

问题解决步骤

1. 修复管道初始化逻辑

当前代码未读取命令行参数,导致DataflowRunner配置未生效,默认使用DirectRunner本地执行。修改代码如下:

public static void main(String[] args) {
    // 解析命令行参数,包括Dataflow配置
    PipelineOptions options = PipelineOptionsFactory.fromArgs(args).withValidation().create();
    // 自定义参数接收Kafka地址
    CustomOptions customOptions = options.as(CustomOptions.class);
    String kafkaIp = customOptions.getKafkaBootstrapServers();

    Pipeline pipeline = Pipeline.create(options);
    // 这里放入原iot_topic_connection方法中的所有逻辑,使用kafkaIp变量
    // ... 原管道逻辑 ...

    // 启动并等待作业完成(流式作业会持续运行)
    pipeline.run().waitUntilFinish();
}

// 自定义Options接口,用于接收Kafka地址参数
public interface CustomOptions extends PipelineOptions {
    @Description("Kafka bootstrap servers (e.g. ip:9092)")
    @Required
    String getKafkaBootstrapServers();
    void setKafkaBootstrapServers(String value);
}

2. 更新Maven命令参数

添加Kafka地址参数,并确保GCS路径存在且有权限:

mvn compile exec:java \
  -Dexec.mainClass=org.example.App \
  -Dexec.args=" \
  --project=round...2 \
  --runner=DataflowRunner \
  --jobName=my-first-job \
  --region=southamerica-west1 \
  --streaming=true \
  --zone=southamerica-west1-a \
  --tempLocation=gs://beam-testing-app/dataflow/temp \
  --stagingLocation=gs://b...p/dataflow/staging \
  --kafkaBootstrapServers=你的KafkaIP:9092 \
  " \
  -Pdataflow-runner

3. 优化pom.xml依赖

移除全局的beam-runners-direct-java依赖,将其放入本地测试专用profile,避免与DataflowRunner冲突:

<profiles>
    <profile>
        <id>dataflow-runner</id>
        <dependencies>
            <dependency>
                <groupId>org.apache.beam</groupId>
                <artifactId>beam-runners-google-cloud-dataflow-java</artifactId>
                <version>2.41.0</version>
                <scope>runtime</scope>
            </dependency>
        </dependencies>
    </profile>
    <profile>
        <id>local-test</id>
        <dependencies>
            <dependency>
                <groupId>org.apache.beam</groupId>
                <artifactId>beam-runners-direct-java</artifactId>
                <version>2.41.0</version>
                <scope>runtime</scope>
            </dependency>
        </dependencies>
    </profile>
</profiles>

4. 关于灵活模板

灵活模板是用于重复部署的可选方案,并非当前问题的必需项。先完成上述修复,确保作业能正常提交到Dataflow后,再根据需求考虑模板化部署。

内容的提问来源于stack exchange,提问作者4 3 2

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 04:38:16