【发布时间】: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