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

Apache Beam读取JSON并流式处理:代码正确性咨询及示例指导

Apache Beam JSON File Reading & Streaming Processing Guide

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:

  1. Parsing the raw JSON string into a structured object
  2. 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: TextIO gives 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 TextIO for simple text/JSON files, FileIO when 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 10:17:48