【发布时间】:2019-07-25 22:12:48
【问题描述】:
我是第一次将 Apache Flink 与 AWS Kineses 一起使用。基本上,我的目标是转换来自 Kinesis 流的传入数据,以便我可以执行简单的转换,例如过滤和聚合。
我正在使用以下内容添加源:
return env.addSource(new FlinkKinesisConsumer<>(inputStreamName, new SimpleStringSchema(), inputProperties));
最终,当我打印传入流时,我按预期获取了 json 数据:
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
DataStream<String> input = createSourceFromStaticConfig(env);
input.print();
这是打印的示例结果:
{"event_num": "5530", "timestmap": "2019-03-04 14:29:44.882376", “金额”:“80.4”,“类型”:“购买”} {“event_num”:“5531”,“timestmap”:“2019-03-04 14:29:44.881379”,“金额”:“11.98”,“类型”:“服务”}
谁能告诉我如何访问这些 json 元素,以便我可以执行简单的转换,例如只选择包含“服务”作为类型的记录?
【问题讨论】:
标签: json amazon-web-services apache-flink amazon-kinesis data-stream