【问题标题】:Beam json parsingBeam json解析
【发布时间】:2018-05-31 15:48:25
【问题描述】:

我正在尝试读取和解析 Apache Beam 代码中的 JSON 文件。

PipelineOptions options = PipelineOptionsFactory.create();
options.setRunner(SparkRunner.class);

Pipeline p = Pipeline.create(options);

PCollection<String> lines = p.apply("ReadMyFile", TextIO.read().from("/Users/xyz/eclipse-workspace/beam-project/myfirst.json"));
System.out.println("lines: " + lines);

下面是我需要从中解析 testdata 的示例 JSON: 我的第一个.json

{  
   “testdata":{  
      “siteOwner”:”xxx”,
      “siteInfo”:{  
         “siteID”:”id_member",
         "siteplatform”:”web”, 
         "siteType”:”soap”,
         "siteURL”:”www”
      }
   }
}

有人可以指导如何解析testdata 并从上述 JSON 文件中获取内容,然后我需要使用 Beam 流式传输数据吗?

【问题讨论】:

  • 不过,我还是做不到。如果有人可以提供帮助,请分享您的想法
  • 我可以使用任何JSON库吗?
  • 我能够使用 JsonFactory 和 Jackson 库解析 JSON 内容。如何将其交给 Beam?
  • 好的,所以在我回答这个问题之前,您能解释一下您的输入和输出要求吗?据我所知,您的问题中有一个嵌套的 json。你希望输出的解析字符串是什么样子的?因为没有这个,你有太多可能的答案可能不符合你的要求
  • 输入是问题中提到的 JSON。我想将“testdata”值解析为字符串。

标签: java json apache-beam


【解决方案1】:

首先,我认为处理“漂亮打印”的 JSON 是不可能的(或至少是常见的)。相反,JSON 数据通常从 newline-delimited JSON 提取,因此您的输入文件应如下所示:

{"testdata":{"siteOwner":"xxx","siteInfo":{"siteID":"id_member","siteplatform":"web","siteType":"soap","siteURL":"www,}}}
{"testdata":{"siteOwner":"yyy","siteInfo":{"siteID":"id_member2","siteplatform":"web","siteType":"soap","siteURL":"www,}}}

之后,使用lines 中的代码,您将拥有“一行行”。接下来,您可以通过在ParDo 中应用 parse-function,将 map 这个“行流”转换为“JSON 流”:

static class ParseJsonFn extends DoFn<String, Json> {

  @ProcessElement
  public void processElement(ProcessContext c) {
    // element here is your line, you can whatever you want, parse, print, etc
    // this function will be simply applied to all elements in your stream
    c.output(parseJson(c.element()))
  }
}

PCollection<Json> jsons = lines.apply(ParDo.of(new ParseJsonFn()))  // now you have a "stream of JSONs"

【讨论】:

  • 好的。为什么 parseJson 对我来说是未定义的?我需要添加正确的 jar 还是导入?
  • @Stella 是的,您需要选择您的 JSON 库:javarevisited.blogspot.com/2016/09/…。它们都提供了parseJson 函数的一些变体。
  • 同样出现此错误:PCollection 类型中的方法 apply(PTransform super PCollection,OutputT>) 不适用于参数 (ApacheBeamPrototype.ParseJsonFn)
  • 好的,对于您的示例代码,我必须选择哪个 json 库?杰克逊还是?
  • parseJson(c.element() 是否正确?它仍然给出错误,因为“ApacheBeamPrototype.ParseJsonFn 类型的方法 parseJson(String) 未定义”
【解决方案2】:

是的,现代 JSON 库可以让您将完全任意的 JSON 和伪 JSON 流解析为对象流,而无需将整个文件加载到内存中。

没有特别需要将您的对象放在一行中。事实上,在设计处理大数据的软件时,避免为批处理数据预留大量内存是一种很好的设计实践,因为此时可以使用仅千字节的内存进行按需流式处理。

看看简短的 Baeldung 教程:http://www.baeldung.com/jackson-streaming-api

我将在此处包含 Baeldung 文章代码的核心部分,因为这是一个很好的做法,以防网站出现故障:

while (jParser.nextToken() != JsonToken.END_OBJECT) {
    String fieldname = jParser.getCurrentName();
    if ("name".equals(fieldname)) {
        jParser.nextToken();
        parsedName = jParser.getText();
    }

    if ("age".equals(fieldname)) {
        jParser.nextToken();
        parsedAge = jParser.getIntValue();
    }

    if ("address".equals(fieldname)) {
        jParser.nextToken();
        while (jParser.nextToken() != JsonToken.END_ARRAY) {
            addresses.add(jParser.getText());
        }
    }
}

在这种情况下,解析器从对象开始标记开始,然后继续解析该对象。在你的情况下,你会想要继续循环直到你的文件完成,所以在你退出这个while循环之后,你会继续前进,直到找到JsonToken.START_OBJECT,然后创建一个新对象,通过这个解析例程,最后将对象交给 Apache Beam。

【讨论】:

  • 另一点是,我想使用 Beam 流式传输解析的测试数据内容
  • 我能够使用 JsonFactory 和 Jackson 库解析 JSON 内容。如何将其交给 Beam?
猜你喜欢
  • 2021-06-18
  • 1970-01-01
  • 2021-11-10
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-10-08
  • 2018-07-27
  • 2011-04-03
相关资源
最近更新 更多