Apache Beam代码编译报错:public String apply处存在语法错误
Fixing Your Apache Beam Code Syntax & Logic Issues
Let's break down the issues in your code and fix them step by step:
Key Problems in Your Current Code
- Incorrect
SimpleFunctionMethod Signature:
TheSimpleFunctionclass requires you to override the single-parameterapply()method, which must match the generic type you declared (<ReadableFile, KV<String,String>>). Your code has anapply()method with two parameters and returns aString—this misalignment with the generic contract is the root cause of your syntax error. - Misplaced Business Logic:
You wrote acreateKV()method with the correct logic to convertReadableFiletoKV<String,String>, but you're not actually using it. Instead, your emptyapply()method is invalid and doesn't serve any purpose. - Incomplete
FileIO.write()Setup:FileIO.write()can't be called directly—it needs additional configuration like output path, filename patterns, and a way to format yourKVdata into writable content.
Fixed Code Example
import org.apache.beam.sdk.Pipeline; import org.apache.beam.sdk.io.FileIO; import org.apache.beam.sdk.io.fs.ReadableFile; import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.options.PipelineOptionsFactory; import org.apache.beam.sdk.transforms.MapElements; import org.apache.beam.sdk.transforms.SimpleFunction; import org.apache.beam.sdk.values.KV; import org.apache.beam.sdk.values.Contextful; import java.io.IOException; public class BeamFileProcessingTest { public static void main(String[] args) throws IOException { System.out.println("Test Log"); PipelineOptions options = PipelineOptionsFactory.create(); options.setRunner(org.apache.beam.runners.spark.SparkRunner.class); Pipeline p = Pipeline.create(options); p.apply(FileIO.match().filepattern("hdfs://path/to/*.gz")) // Beam auto-detects gzip compression from the filename, no need for explicit config .apply(FileIO.readMatches()) .apply(MapElements.via( new SimpleFunction<ReadableFile, KV<String, String>>() { private static final long serialVersionUID = -7867677L; @Override public KV<String, String> apply(ReadableFile element) { String fileContent = null; try { fileContent = element.readFullyAsUTF8String(); } catch (IOException e) { // Replace with proper logging in production code e.printStackTrace(); } return KV.of(element.getMetadata().resourceId().toString(), fileContent); } } )) // Configure FileIO to write KV pairs as formatted text lines .apply(FileIO.write() .to("hdfs://path/to/output-directory") .withPrefix("processed-file-") .withSuffix(".txt") .via(Contextful.fn(kv -> String.format("%s: %s", kv.getKey(), kv.getValue())))); // Wait for pipeline to finish execution p.run().waitUntilFinish(); } }
Additional Notes
- Error Handling: I added a basic exception print for IO errors, but in production, use a logging framework like SLF4J instead of console prints.
- Spark Runner Dependency: Ensure your build file (Maven/Gradle) includes the Apache Beam Spark Runner dependency to avoid runtime class-not-found errors.
- Imports: Double-check all required imports are present—missing imports can often cause hidden syntax issues that are hard to spot.
内容的提问来源于stack exchange,提问作者Stella
相关产品推荐
相关产品推荐

