如何将本地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
相关产品推荐
相关产品推荐

