【问题标题】:How to understand Window mechanism in Apache Flink如何理解 Apache Flink 中的 Window 机制
【发布时间】:2020-06-08 18:26:43
【问题描述】:

我正在学习如何使用 Flink 处理流数据。

据我了解,我可以多次使用函数map进行各种变换。

说数据源一直在向 Flink 发送字符串。所有字符串都是 JSON 格式的数据,如下所示:

{"name":"titi","age":18}
{"name":"toto","age":20}
...

这是我的代码:

final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
FlinkPravegaReader<String> source = FlinkPravegaReader.<String>builder()
    .withPravegaConfig(pravegaConfig)
    .forStream(stream)
    .withDeserializationSchema(new PravegaDeserializationSchema<>(String.class, new JavaSerializer<>()))
    .build();

// Convert String to Json Object
// MyJson is a POJO class, defined by me
DataStream<MyJson> jsonStream = env.addSource(source).name("Pravega Stream")
    .map(new MapFunction<String, MyJson>() {
    @Override
    public MyJson map(String s) throws Exception {
        MyJson myJson = JSON.parseObject(s, MyJson.class);
        return myJson;
        }
    });
// Convert MyJson Object to String and extract what I need
DataStream<String> valueInJson = jsonStream
    .map(new MapFunction<MyJson, String>() {
        @Override
        public String map(MyJson myJson) throws Exception {
            return myJson.getName().toString();
        }
    });
valueInJson.print();
env.execute("StreamingJob");

如您所见,我的示例非常简单: 获取和反序列化数据 ---> 将字符串转换为 Json 对象 ---> 将 Json 对象转换为字符串并获取我需要的内容(这里我只需要 name)。

目前看来,一切正常。我确实从日志文件中得到了预期的输出。

但是,我知道 Flink 为我们提供了一个强大的功能:Window。

我想知道如何在我的示例中使用这种机制。

例如,如果我想用一些 2 秒的窗口分割数据流,如何编码?

我试过这样:

DataStream<String> valueInJson = jsonStream
    .timeWindow(Time.seconds(2))
    .map(new MapFunction<MyJson, String>() {
        @Override
        public String map(MyJson myJson) throws Exception {
            return myJson.toString();
        }
    });
valueInJson.print();

但是,我得到了一个错误:

找不到符号
符号:方法
timeWindow(org.apache.flink.streaming.api.windowing.time.Time)
位置:变量 jsonStream 类型 org.apache.flink.streaming.api.datastream.DataStream

但是,我已经导入了:

import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
import org.apache.flink.streaming.api.windowing.time.Time;

为什么会出现这个错误?我是否错误地使用了 Windows?我错过了对 Flink 的了解吗?

【问题讨论】:

标签: apache-flink flink-streaming


【解决方案1】:

您遇到错误是因为timeWindow() 函数是在KeyedStream 中定义的,而不是在DataStream 中定义的,因为它是基于键的操作。在您的情况下,将timeWindow() 更改为timeWindowAll() 就足够了。

【讨论】:

  • 这还不够——windowedStream 上没有 map 方法。需要有类似应用程序或带有 WindowFunction 或 ProcessWindowFunction 的进程来处理窗口内容。
猜你喜欢
  • 1970-01-01
  • 2021-09-18
  • 2021-07-18
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多