Apache Beam读取JSON并流式处理:代码正确性咨询及示例指导
Hey there! Let's break down your problem and get your JSON data flowing properly in Apache Beam.
Why Your First Code Isn't Showing Data
First off, the reason you're only seeing lines: ReadMyFile/Read.out [PCollection] is because of Beam's lazy execution model. When you create a PCollection, you're just defining a step in your pipeline—not actually running it yet. That System.out.println is just printing metadata about the PCollection, not the actual file content. You need to trigger pipeline execution and add steps to process/output the data.
Analysis of Your Two Code Snippets
Code 1: TextIO.read()
This is the right starting point for reading text-based files like JSON. TextIO.read().from(...) correctly reads the file into a PCollection<String> where each element is a line of the file. But you're missing two key pieces:
- Parsing the raw JSON string into a structured object
- Triggering the pipeline to run with
p.run().waitUntilFinish()
Code 2: FileIO.match()
This snippet is incomplete. FileIO.match() only finds files matching the pattern—it doesn't read their content. You'd need to chain it with FileIO.readMatches() to actually load the file data, like:
p.apply(FileIO.match().filepattern("/path/to/test.json")) .apply(FileIO.readMatches()) .apply(ParDo.of(new DoFn<ReadableFile, String>() { @ProcessElement public void processElement(ProcessContext c) throws IOException { String content = c.element().open().readFullyAsUTF8String(); c.output(content); } }));
This is more flexible for complex file handling (like reading multiple files, handling metadata), but for a single JSON file, TextIO is simpler.
Complete Working Example
Here's a full example that reads your JSON file, parses the testdata object, and prints the content. We'll use Jackson for JSON parsing (add the Jackson dependency to your project first).
import org.apache.beam.sdk.Pipeline; import org.apache.beam.sdk.io.TextIO; import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.options.PipelineOptionsFactory; import org.apache.beam.sdk.transforms.DoFn; import org.apache.beam.sdk.transforms.ParDo; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; public class JsonReaderExample { public static void main(String[] args) { // Set up pipeline options PipelineOptions options = PipelineOptionsFactory.create(); options.setRunner(org.apache.beam.runners.spark.SparkRunner.class); // Create pipeline Pipeline p = Pipeline.create(options); // Read JSON file, parse, and print testdata p.apply("ReadJSONFile", TextIO.read().from("/Users/xyz/eclipse-workspace/beam-prototype/test.json")) .apply("ParseJSON", ParDo.of(new DoFn<String, JsonNode>() { private final ObjectMapper mapper = new ObjectMapper(); @ProcessElement public void processElement(ProcessContext c) throws Exception { // Parse the entire JSON string JsonNode root = mapper.readTree(c.element()); // Extract the testdata node JsonNode testData = root.get("testdata"); c.output(testData); } })) .apply("PrintTestData", ParDo.of(new DoFn<JsonNode, Void>() { @ProcessElement public void processElement(ProcessContext c) { System.out.println("Extracted testdata: " + c.element().toPrettyString()); } })); // Run the pipeline (critical step!) p.run().waitUntilFinish(); } }
Key Steps to Remember
- Always run the pipeline: Without
p.run().waitUntilFinish(), none of your steps will execute—Beam just builds the pipeline graph but doesn't process data. - Parse JSON properly:
TextIOgives you raw strings; use a library like Jackson or Gson to convert them into structured objects for easier processing. - Choose the right IO for your use case: Use
TextIOfor simple text/JSON files,FileIOwhen you need more control over file metadata or complex file patterns.
If you're working with streaming data (even if your input is a single file, Beam can treat it as a bounded stream), this approach still works—just ensure your runner (SparkRunner in your case) is configured correctly for streaming if needed.
内容的提问来源于stack exchange,提问作者Stella

