Ste*_*lla 8 java json apache-beam
我试图在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);
Run Code Online (Sandbox Code Playgroud)
下面是我需要解析的示例JSON testdata:myfirst.json
{
“testdata":{
“siteOwner”:”xxx”,
“siteInfo”:{
“siteID”:”id_member",
"siteplatform”:”web”,
"siteType”:”soap”,
"siteURL”:”www”
}
}
}
Run Code Online (Sandbox Code Playgroud)
有人可以指导如何testdata从上面的JSON文件解析和获取内容,然后我需要使用Beam流式传输数据?
首先,我认为处理“漂亮打印”的 JSON 是不可能的(或者至少是常见的)。相反,JSON 数据通常从换行符分隔的 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,}}}
Run Code Online (Sandbox Code Playgroud)
之后,lines你的代码就变成了“一行行”。接下来,您可以map通过在以下位置应用解析函数,将“行流”转换为“JSON 流” ParDo:
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"
Run Code Online (Sandbox Code Playgroud)