【问题标题】:How to parse json in a beam pipeline?如何在光束管道中解析 json?
【发布时间】:2019-04-18 02:18:55
【问题描述】:

我的数据是换行符分隔的 json 格式,如下所示。我正在从 Kafka 主题中读取此类数据。

{"sender":"S1","senderHost":"ip-10-20-30-40","timestamp":"2018-08-13T16:17:12.874Z","topic":"test","messageType":"type_1","data":{"name":"John Doe", "id":"12DROIY321"}}

我想构建一个 apache Beam 管道,它从 Kafka 读取此数据,解析此 json 格式以提供如下所示的输出:

S1,2018-08-13T16:17:12.874Z,type_1,12DROIY321

输出基本上是一个逗号分隔的字符串,由数据中的发送者、时间戳、messageType 和 id 组成。

到目前为止我的代码如下:

public class Pipeline1{
    public static void main(String[] args){
        PipelineOptions options = PipelineOptionsFactory.create();

        // Create the Pipeline object with the options we defined above.
        Pipeline p = Pipeline.create(options);

        p.apply(KafkaIO.<Long, String>read()
                .withBootstrapServers("localhost:9092")
                .withTopic("test")
                .withKeyDeserializer(LongDeserializer.class)
                .withValueDeserializer(StringDeserializer.class)

                .updateConsumerProperties(ImmutableMap.of("auto.offset.reset", (Object)"earliest"))

                // We're writing to a file, which does not support unbounded data sources. This line makes it bounded to
                // the first 35 records.
                // In reality, we would likely be writing to a data source that supports unbounded data, such as BigQuery.
                .withMaxNumRecords(35)

                .withoutMetadata() // PCollection<KV<Long, String>>
        )
                .apply(Values.<String>create())
                .apply(TextIO.write().to("test"));

        p.run().waitUntilFinish();

    }
}

我无法弄清楚如何解析 json 以在管道中获取所需的 csv 格式。使用上面的代码,我可以将相同的 json 行写入一个文件,并使用下面的代码,我可以解析 json,但是任何人都可以帮我弄清楚如何通过光束管道来完成这个附加步骤逻辑?

JSONParser parser = new JSONParser();
            Object obj = null;
            try {
                obj = parser.parse(strLine);
            } catch (ParseException e) {
                e.printStackTrace();
            }
            JSONObject jsonObject =  (JSONObject) obj;

            String sender = (String) jsonObject.get("sender");

            String messageType = (String) jsonObject.get("messageType");

            String timestamp = (String) jsonObject.get("timestamp");

            System.out.println(sender+","+timestamp+","+messageType);

【问题讨论】:

  • 这里有两件事,您可以使用jsonObject.getString("timestamp"); 而不是jsonObject.get("timestamp");.2) how to parse the json to get the required csv format within the pipeline 是什么意思所以您想将JSON 中的一些属性保存到csv 文件中?你能举个例子吗?
  • Kafka 有一个 JSONDeserializer,顺便说一句

标签: java json apache-kafka apache-beam


【解决方案1】:

根据文档,您需要编写转换(或找到与您的用例匹配的转换)。

https://beam.apache.org/documentation/programming-guide/#composite-transforms

该文档还提供了一个很好的示例。

应该产生你的输出的例子:

.apply(Values.<String>create())
.apply(
    "JSONtoData",                     // the transform name
    ParDo.of(new DoFn<String, String>() {    // a DoFn as an anonymous inner class instance
        @ProcessElement
        public void processElement(@Element String word, OutputReceiver<String> out) {
            JSONParser parser = new JSONParser();
            Object obj = null;
            try {
                obj = parser.parse(strLine);
            } catch (ParseException e) {
                e.printStackTrace();
            }
            JSONObject jsonObject =  (JSONObject) obj;

            String sender = (String) jsonObject.get("sender");
            String messageType = (String) jsonObject.get("messageType");
            String timestamp = (String) jsonObject.get("timestamp");

            out.output(sender+","+timestamp+","+messageType);
        }
   }));

要返回 CSV 值,只需将泛型更改为:

new DoFn<String, YourCSVClassHere>()
OutputReceiver<YourCSVClassHere> out

我没有测试此代码,使用风险自负。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-01-17
    • 2020-06-18
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-09-15
    • 1970-01-01
    相关资源
    最近更新 更多