ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

PyFlink DataStream 基础算子实战:map、filter、key_by 与 sum 入门指南

PyFlink DataStream 基础算子实战:map、filter、key_by 与 sum 入门指南 大数据流处理批处理数据工程【免费下载链接】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_bysum按国家分组累计电话数字段。完整代码位于 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(streamsys.stdout, levellogging.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_infoTypes.ROW_NAMED([id, info], [Types.INT(), Types.STRING()]) )这里用from_collection从内存集合创建 DataStream每条元素是一个二元组id整数和infoJSON 字符串。由于显式指定了type_infoTypes.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_sourceversionadded 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.pyfilter与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_typeTypes.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_infoTypes.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的 sinkenv.execute()触发作业执行。PyFlink 采用惰性执行模型——所有算子调用只构建执行计划DataStream 变换图必须显式调用execute()才会真正提交运行并打印结果。深入理解PyFlink 算子的底层执行原理结合上文源码可以归纳 PyFlink 基础算子的通用实现模式Python 侧构建变换图map/filter/key_by等 API 方法在 Python 侧返回新的DataStream/KeyedStream对象算子被命名Map/Filter/STREAM_KEY_BY_MAP_OPERATOR_NAME等供 Web UI 展示ProcessFunction 适配器桥接普通 Python 函数被包装成继承自ProcessFunction的适配器类MapProcessFunctionAdapter、FilterProcessFunctionAdapter、AddKey通过生成器式process_element表达一进一出一进零出/一进一出等不同语义再经 PyFlink 的 Python 算子运行时flink-python/src/main下的 Java 侧实现执行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.rstprocess_json_data.py更复杂的 JSON 处理对应文档 process_json_data.rststate_access.py键控状态访问对应文档 state.rstevent_time_timer.py事件时间与定时器对应文档 timer.rstwindowing 目录滚动/滑动/会话窗口示例对应文档 window.rstconnectors 目录Kafka、Pulsar、Elasticsearch 等连接器示例对应文档 connectors.rst。建议按基础算子 → JSON 处理 → 窗口 → 状态 → 连接器的顺序循序渐进即可快速构建出完整的 PyFlink 流处理能力。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Elm控制语句精讲if、case-of和let-in的实战应用技巧Elm控制语句精讲if、case of和let in的实战应用技巧 Elm作为一门函数式前端编程语言其控制语句的设计体现了函数式编程的优雅与严谨。在Elm开大数据流处理批处理数据工程Timely Dataflow 算子入门从 Map、Filter 到 partition 与 exchange 的流式计算核心Pathway 底层引擎实战Timely Dataflow 算子入门从 Map、Filter 到 partition 与 exchange 的流式计算核心Pathway 底层引擎实战后端流处理实时分析数据工程人工智能RAGPyFlink DataStream API 完全指南从基础流转换到窗口、连接与广播流PyFlink DataStream API 完全指南从基础流转换到窗口、连接与广播流 本文基于 Apache Flink 官方仓库中 flink pytho大数据流处理批处理数据工程上一篇3步用magic.css实现炫酷菜单动画提升网站导航体验的终极指南下一篇douyin-downloader抖音无水印批量下载上手指南与技术拆解创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表