- 大数据
- 流处理
- 批处理
- 数据工程
【免费下载链接】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 基础算子的通用实现模式:
- 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.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
相关推荐
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
大数据流处理批处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考