【问题标题】:Parsing json from incoming datastream to perform simple transformations in Flink从传入的数据流中解析 json 以在 Flink 中执行简单的转换
【发布时间】: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


    【解决方案1】:

    当您使用SimpleStringSchema 时,结果事件流的类型为String。因此,您需要先解析字符串,然后才能应用过滤器等。

    您可能想看看JsonNodeDeserializationSchema,它将产生ObjectNode

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-10-31
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多