☰
PyFlink DataStream 基础算子实战:map、filter、key_by 与 sum 入门指南
2026/9/25 3:28:43 网站建设 项目流程
  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】flink

项目地址:https://gitcode.com/gh_mirrors/fli/flink
点击查看免费下载

本文基于 Apache Flink 仓库中 PyFlink DataStream 官方示例(basic_operations.py 及其文档 basic_operations.rst),系统讲解 PyFlink 流处理中最核心的四个基础算子:map(逐条映射转换)、filter(数据筛选)、key_by(按键分区)与sum(滚动聚合)。读完本文,你将掌握如何用 Python 定义执行环境、构造内存数据源,并把一条条原始 JSON 字符串流式转换为结构化数据后再完成过滤与分组求和,为后续学习窗口、状态、连接器打下基础。

示例概览:一条流水线看懂四个算子

示例程序演示了 PyFlink 最典型的"定义环境 → 构造数据源 → 链式调用算子 → 输出结果"开发范式。其核心思路是:构造 4 条形如(id, json字符串)的记录,先用map解析 JSON 并修改电话号码,再用filter只保留指定 id 的记录,最后用key_by+sum按国家分组累计电话数字段。

完整代码位于 basic_operations.py,它与仓库中其他示例(word_count、process_json_data、state_access、event_time_timer、windowing等)一同通过 index.rst 的 toctree 组织进文档体系,对应文档页即 basic_operations.rst。

程序入口如下:

if __name__ == '__main__': logging.basicConfig(stream=sys.stdout, level=logging.INFO, format="%(message)s") basic_operations()

示例在main块中配置了输出到标准输出的日志,便于直接观察算子执行结果。

第一步:创建执行环境并设置并行度

env = StreamExecutionEnvironment.get_execution_environment() env.set_parallelism(1)

StreamExecutionEnvironment是 PyFlink DataStream 程序的起点,所有数据流都从它派生。set_parallelism(1)将整个作业的并行度显式设为 1,这样输出顺序与输入顺序一致,便于理解算子的逐条处理语义。

从源码看,stream_execution_environment.py 中set_parallelism直接桥接到 Java 端的StreamExecutionEnvironment.setParallelism():

def set_parallelism(self, parallelism: int) -> 'StreamExecutionEnvironment': self._j_stream_execution_environment = \ self._j_stream_execution_environment.setParallelism(parallelism) return self

需要注意的是(依据该方法的 docstring):此处设置的是"此环境内所有算子"的默认并行度;LocalStreamEnvironment默认并行度等于硬件上下文数(CPU 核数/线程数),而通过命令行客户端提交 JAR 时,默认并行度则取自作业配置。示例将其设为 1,主要是为了方便演示和阅读输出。

第二步:用 from_collection 构造内存数据源

ds = env.from_collection( collection=[ (1, '{"name": "Flink", "tel": 123, "addr": {"country": "Germany", "city": "Berlin"}}'), (2, '{"name": "hello", "tel": 135, "addr": {"country": "China", "city": "Shanghai"}}'), (3, '{"name": "world", "tel": 124, "addr": {"country": "USA", "city": "NewYork"}}'), (4, '{"name": "PyFlink", "tel": 32, "addr": {"country": "China", "city": "Hangzhou"}}') ], type_info=Types.ROW_NAMED(["id", "info"], [Types.INT(), Types.STRING()]) )

这里用from_collection从内存集合创建 DataStream,每条元素是一个二元组:id(整数)和info(JSON 字符串)。由于显式指定了type_info=Types.ROW_NAMED(["id", "info"], [Types.INT(), Types.STRING()]),元素被建模为带字段名的 Row 类型,后续在 Python 函数中可以直接通过data.info访问第二个字段。

从源码看,from_collection 的实现有两个关键点:

  • 类型转换:若指定了type_info,会先调用type_info.to_internal_type(element)将每个元素转换为内部类型表示,再通过PythonBridgeUtils.readPythonObjects读取;若未指定类型,则退化为 pickle 序列化的字节数组(Types.PICKLED_BYTE_ARRAY),此时不提供字段名访问能力。
  • 非并行源:该方法 docstring 明确指出"This operation will result in a non-parallel data stream source, i.e. a data stream source with parallelism one"——from_collection产生的是一个并行度为 1 的源,元素经临时文件落盘后由InputFormatSourceFunction以BOUNDED有界流的形式读出,因此适合测试和演示。

若数据规模较大或来自外部系统,可改用from_source(versionadded 1.13.0,见 stream_execution_environment.py)接入 Kafka、Pulsar 等连接器,或read_text_file逐行读取文件。

第三步:map 算子——逐条映射转换

def update_tel(data): # parse the json json_data = json.loads(data.info) json_data['tel'] += 1 return data.id, json.dumps(json_data) show(ds.map(update_tel), env)

map对数据流中每个元素调用一次函数,且每次调用恰好返回一个元素。这里的update_tel完成"JSON 解析 → 电话号 +1 → 重新序列化"的转换:

  • 输入:data是 ROW_NAMED 元素,data.info是 JSON 字符串;
  • 输出:(data.id, json.dumps(json_data))二元组。

在 data_stream.py 中,map的实现细节如下:

def map(self, func: Union[Callable, MapFunction], output_type: TypeInformation = None) -> 'DataStream': class MapProcessFunctionAdapter(ProcessFunction): def __init__(self, map_func): if isinstance(map_func, MapFunction): self._open_func = map_func.open self._close_func = map_func.close self._map_func = map_func.map else: self._open_func = None self._close_func = None self._map_func = map_func def process_element(self, value, ctx: 'ProcessFunction.Context'): yield self._map_func(value) return self.process(MapProcessFunctionAdapter(func), output_type).name("Map")

值得注意的实现事实:

  • 函数与类两种形式:map既接受普通 Python 可调用对象(callable),也接受继承MapFunction的类(此时会自动调用其open/close生命周期方法);
  • 内部桥接:普通函数会被包装进一个继承ProcessFunction的MapProcessFunctionAdapter,通过process_element的生成器语义逐条产出结果,最终以"Map"命名算子(对应 Web UI 中的算子名);
  • 输出类型推断:如果未显式传入output_type,输出数据将以 pickle 原始字节数组序列化(见该方法 docstring 中 "the output data will be serialized as pickle primitive byte array"),类型信息在算子链传递时可能受影响,生产环境中建议显式声明。

运行输出(tel 均 +1):

(1, '{"name": "Flink", "tel": 124, "addr": {"country": "Germany", "city": "Berlin"}}') (2, '{"name": "hello", "tel": 136, "addr": {"country": "China", "city": "Shanghai"}}') (3, '{"name": "world", "tel": 125, "addr": {"country": "USA", "city": "NewYork"}}') (4, '{"name": "PyFlink", "tel": 33, "addr": {"country": "China", "city": "Hangzhou"}}')

第四步:filter 算子——按条件保留数据

show(ds.filter(lambda data: data.id == 1).map(update_tel), env)

filter对每个元素执行谓词函数,仅保留返回True的元素,其余全部丢弃。这里用lambda data: data.id == 1只保留 id 为 1 的记录,再复用update_tel做映射。

从源码看(data_stream.py),filter与map的桥接模式一致:可调用对象或FilterFunction实例会被包装进FilterProcessFunctionAdapter,且关键区别在于——process_element只在谓词为真时才yield value,从而实现过滤语义:

def process_element(self, value, ctx: 'ProcessFunction.Context'): if self._filter_func(value): yield value

此外,filter的输出类型直接取自上游变换的getTransformation().getOutputType(),算子命名为"Filter"。运行输出只剩一条:

(1, '{"name": "Flink", "tel": 124, "addr": {"country": "Germany", "city": "Berlin"}}')

第五步:key_by + sum——按键分区与滚动求和

show(ds.map(lambda data: (json.loads(data.info)['addr']['country'], json.loads(data.info)['tel'])) .key_by(lambda data: data[0]).sum(1), env)

这是示例中最具实战价值的一段:先用map把每条记录提炼为(country, tel)二元组,再用key_by按国家名分区,最后sum(1)对第 1 个字段(tel)做按 key 独立维护的滚动求和。运行输出:

('Germany', 123) ('China', 135) ('USA', 124) ('China', 167)

注意('China', 167)是 135 + 32 的结果——两条中国记录被路由到同一个 key 分区,tel 字段逐条累加,这正是 key 聚合的核心语义。

key_by 的源码实现

key_by 内部先通过AddKey这个ProcessFunction为每条记录附加提取出的 key(组成 Row),再调用 Java 端keyBy完成物理分区:

class AddKey(ProcessFunction): def __init__(self, key_selector): if isinstance(key_selector, KeySelector): self._get_key_func = key_selector.get_key else: self._get_key_func = key_selector ... stream_with_key_info = self.process( AddKey(key_selector), output_type=Types.ROW([key_type, output_type_info])) stream_with_key_info.name(...STREAM_KEY_BY_MAP_OPERATOR_NAME) JKeyByKeySelector = gateway.jvm.KeyByKeySelector key_stream = KeyedStream( stream_with_key_info._j_data_stream.keyBy( JKeyByKeySelector(), Types.ROW([key_type]).get_java_type_info()), output_type_info, self)

它返回一个KeyedStream。从实现可推断:key_by生成的KeyedStream是后续所有键控状态(keyed state)与窗口聚合的基础——同一 key 的所有记录都会进入同一分区。

sum 的源码实现

sum 是KeyedStream上提供的滚动聚合之一,其 docstring 明确:它对指定位置做滚动求和(rolling sum),每个 key 维护一个独立累加器。参数position_to_sum既可以是表示列索引的整数,也可以是表示字段名的字符串(该方法自 1.16.0 起支持):

# Tuple 数据按索引聚合 >>> ds = env.from_collection([('a', 1), ('a', 2), ('b', 1), ('b', 5)]) >>> ds.key_by(lambda x: x[0]).sum(1) # Row 数据按字段名聚合 >>> ds = env.from_collection( ... [('a', 1), ('a', 2), ('a', 3), ('b', 1), ('b', 2)], ... type_info=Types.ROW_NAMED(["key", "value"], [Types.STRING(), Types.INT()])) >>> ds.key_by(lambda x: x[0]).sum("value")

其内部委托给_accumulate(position_to_sum, KeyedStream.AccumulateType.SUM),累加逻辑由AccumulateReduceFunction承担(见 data_stream.py)。同类算子还包括min(当前最小值,同样支持索引或字段名,data_stream.py),可参照使用。

输出与执行:print + execute

示例统一封装了输出辅助函数:

def show(ds, env): ds.print() env.execute()
  • ds.print():把结果流输出到标准输出(对应算子名为Print的 sink);
  • env.execute():触发作业执行。PyFlink 采用惰性执行模型——所有算子调用只构建执行计划(DataStream 变换图),必须显式调用execute()才会真正提交运行并打印结果。

深入理解:PyFlink 算子的底层执行原理

结合上文源码,可以归纳 PyFlink 基础算子的通用实现模式:

  1. Python 侧构建变换图:map/filter/key_by等 API 方法在 Python 侧返回新的DataStream/KeyedStream对象,算子被命名("Map"/"Filter"/STREAM_KEY_BY_MAP_OPERATOR_NAME等),供 Web UI 展示;
  2. ProcessFunction 适配器桥接:普通 Python 函数被包装成继承自ProcessFunction的适配器类(MapProcessFunctionAdapter、FilterProcessFunctionAdapter、AddKey),通过生成器式process_element表达"一进一出""一进零出/一进一出"等不同语义,再经 PyFlink 的 Python 算子运行时(flink-python/src/main下的 Java 侧实现)执行;
  3. Java 算子链复用:像keyBy、sum这类需要 Flink 运行时分区与状态支持的操作,最终通过gateway.jvm调用 Java 端 API(如StreamExecutionEnvironment、KeyedStream.keyBy),与 Java DataStream API 共享同一套执行引擎。

因此,PyFlink 基础算子的行为语义与 Java DataStream API 完全一致:map一对一、filter真值保留、key_by同 key 同分区、sum按 key 滚动累加。

延伸学习:PyFlink DataStream 示例体系

basic_operations.py是 PyFlink DataStream 示例的入门篇,仓库中同一目录下(flink-python/pyflink/examples/datastream)还有更多进阶示例可供串联学习:

  • word_count.py / streaming_word_count.py:经典词频统计,覆盖map/flat_map/key_by/sum组合(对应文档 word_count.rst);
  • process_json_data.py:更复杂的 JSON 处理(对应文档 process_json_data.rst);
  • state_access.py:键控状态访问(对应文档 state.rst);
  • event_time_timer.py:事件时间与定时器(对应文档 timer.rst);
  • windowing 目录:滚动/滑动/会话窗口示例(对应文档 window.rst);
  • connectors 目录:Kafka、Pulsar、Elasticsearch 等连接器示例(对应文档 connectors.rst)。

建议按"基础算子 → JSON 处理 → 窗口 → 状态 → 连接器"的顺序循序渐进,即可快速构建出完整的 PyFlink 流处理能力。

  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】flink

项目地址:https://gitcode.com/gh_mirrors/fli/flink
点击查看免费下载

相关推荐

上一篇:3步用magic.css实现炫酷菜单动画:提升网站导航体验的终极指南
下一篇:douyin-downloader:抖音无水印批量下载上手指南与技术拆解

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询