请描述在 Apache Flink 中利用算子执行数据转换的常见方式,并列举几个典型的转换算子及其功能。
考察说明
考察对 Flink 数据转换算子的理解与应用能力。
回答思路
- 【回答框架 1】Flink 的 DataStream API 通过算子对数据流进行转换,常见的转换包括单流转换、多流转换和窗口操作等。单流转换算子如 map、flatMap、filter 等,其中 map 是一对一转换,flatMap 是一对多转换,filter 是条件过滤。
- 【回答框架 2】多流转换算子包括 connect 和 union。union 要求流类型一致且只能是两条以上流,而 connect 可以连接类型不同的流,并通过 CoProcessFunction 进行合并处理。
- 【回答框架 3】还有 keyBy、reduce、aggregate 等算子用于分组和聚合。keyBy 根据 key 分区,之后可以使用 reduce 进行增量聚合,窗口算子如 window、apply 等用于对有限的时间或数量范围内的数据进行处理。
- 【回答框架 4】使用算子的核心是明确输入输出类型和数据流语义,如 event time 处理时需设置 watermark 和窗口分配器,确保结果正确。
- 【关键点 1】map 是一对一转换,flatMap 是一对多。
- 【关键点 2】union 要求流类型一致,connect 可连接不同类型流。
- 【关键点 3】keyBy 之后才能进行 keyed 状态和窗口操作。
- 【关键点 4】窗口算子支持滚动、滑动和会话窗口。
- 【易错点 1】误用 union 连接类型不一致的流。
- 【易错点 2】忽略事件时间设置导致窗口结果不准确。
- 【易错点 3】窗口内使用无状态算子导致状态丢失。