news 2026/9/25 2:29:34

PyFlink StreamExecutionEnvironment 完全指南:从环境创建到作业提交的核心 API 实战

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
PyFlink StreamExecutionEnvironment 完全指南:从环境创建到作业提交的核心 API 实战
  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】flink

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

StreamExecutionEnvironment是 PyFlink DataStream 程序运行的上下文环境:它既决定了作业在本地 JVM 还是远程集群上执行,也提供了控制作业并行度、容错(Checkpoint)、运行时模式、时间语义以及与外部世界交互(数据接入、Python 依赖)的全部入口。本文以 PyFlink DataStream 参考文档 为核心骨架,结合仓库内 stream_execution_environment.py 的源码实现与其测试用例,系统讲解该环境的每一个核心 API:读完本文后,你将能独立完成一个 PyFlink DataStream 作业从环境初始化、并行度与执行模式配置、Checkpoint 与状态后端设置、Python 依赖打包,到最后执行与获取执行计划的全流程。

StreamExecutionEnvironment 是什么

在 PyFlink 中,StreamExecutionEnvironment是所有 DataStream 程序的上下文,相当于 Flink 作业的"总控制台"。源码类注释给出了最精炼的定义(见 stream_execution_environment.py):

"The StreamExecutionEnvironment is the context in which a streaming program is executed. ALocalStreamEnvironmentwill cause execution in the attached JVM, aRemoteStreamEnvironmentwill cause execution on a remote setup."

它有两个层面的职责:

  • 控制作业执行:设置并行度(parallelism)、最大并行度、故障恢复与 Checkpoint 参数、运行时执行模式、buffer 冲刷频率、算子链(operator chaining)等;
  • 与外部世界交互:添加数据源(source)、读取文件与集合、注册分布式缓存文件、添加 Python 依赖与 JAR 包等。

底层实现上,PyFlink 的StreamExecutionEnvironment是对 Java 侧org.apache.flink.streaming.api.environment.StreamExecutionEnvironment的 Py4J 封装:Python 对象持有_j_stream_execution_environment(JavaObject),所有方法最终都会调用对应的 Java API(如setParallelism、enableCheckpointing、getStreamGraph等)。

获取执行环境

与 Java 版 Flink 类似,PyFlink 推荐通过静态工厂方法创建环境:

from pyflink.datastream import StreamExecutionEnvironment env = StreamExecutionEnvironment.get_execution_environment()

get_execution_environment 支持传入一个可选的pyflink.common.Configuration:

  • 独立运行(standalone)时:返回本地执行环境(LocalStreamEnvironment),直接在附着 JVM 中执行;
  • 从命令行提交时:给定的configuration会叠加在来自config.yaml的全局配置之上,重复的配置项会被覆盖(即代码内配置优先级更高)。

该方法最终调用 Java 侧StreamExecutionEnvironment.getExecutionEnvironment(configuration)来获取对应环境。

并行度控制:set_parallelism 与 set_max_parallelism

并行度是 DataStream 作业最重要的调优参数之一。环境层提供了一组成对出现的读写 API:

API作用关键语义
set_parallelism(parallelism)/get_parallelism()设置/获取环境中所有算子(map、filter 等)默认运行的并行实例数会覆盖该环境的默认并行度;本地环境默认等于 CPU 核心数;通过命令行从 JAR 提交时,默认并行度取该运行环境的配置值(见 set_parallelism)
set_max_parallelism(max_parallelism)/get_max_parallelism()设置/获取程序定义的最大并行度上限为32768(2^15),需满足0 < max_parallelism <= 2^15;最大并行度是动态扩容的上限,同时决定了分区状态使用的key group 数量(见 set_max_parallelism)
set_default_local_parallelism(n)/get_default_local_parallelism()设置/获取本地执行环境使用的默认并行度仅对本地执行生效(见 set_default_local_parallelism)

测试用例 test_get_set_parallelism 验证了设置并行度 10 后get_parallelism()返回 10;test_get_set_max_parallelism 验证了最大并行度 12 的读写一致性。

为什么最大并行度重要?从源码注释可以提炼出两个核心点:一是它界定了动态扩缩容的上限(超过该值扩容将不可用);二是它直接决定 keyed state 的 key group 划分,因此一旦作业运行后,最大并行度不能随意修改,否则会影响状态恢复时的 key 分布。set_parallelism返回环境对象自身,支持链式调用。

RuntimeExecutionMode:流/批执行的运行时模式

RuntimeExecutionMode枚举定义了 DataStream 程序的运行时执行模式,它会影响任务调度方式、网络 shuffle 行为和(事件/处理)时间语义,部分算子还会根据执行模式改变记录发射行为。原文档对此给出了完整定义(同时可见于 execution_mode.py):

  • STREAMING:以流式语义执行。所有任务在执行开始前全部部署,Checkpoint 被启用,processing time 与 event time 均得到完整支持。
  • BATCH:以批式语义执行。任务按所属调度区域(scheduling region)渐进式调度,区域间的 shuffle 是blocking(阻塞式)的;watermark 被假定为"完美"的,即不存在迟到数据;processing time 被假定为执行期间不推进。
  • AUTOMATIC:由 Flink 自动判定——若所有 source 都是有界(bounded)的则按 BATCH 执行,若存在至少一个无界(unbounded)source 则按 STREAMING 执行。

设置运行时模式

两种等价方式:

# 方式一:API 直接设置(等价于配置项 execution.runtime-mode) from pyflink.datastream.execution_mode import RuntimeExecutionMode env.set_runtime_mode(RuntimeExecutionMode.BATCH) # 方式二:通过 configure 统一注入配置 from pyflink.common import Configuration config = Configuration() config.set_string("execution.runtime-mode", "BATCH") env.configure(config)

set_runtime_mode(1.13.0 新增)的源码注释特别建议:优先不要用 API 硬编码模式,而是在命令行提交作业时通过execution.runtime-mode配置项指定,这样同一份应用代码可以在任意执行模式下运行,保持代码与运行环境解耦。

测试 test_set_runtime_mode 验证了调用set_runtime_mode(RuntimeExecutionMode.BATCH)后,底层ExecutionOptions.RUNTIME_MODE配置值变为"BATCH"。

运行时执行行为:buffer 冲刷、算子链与槽位共享组

输出缓冲区冲刷:set_buffer_timeout

set_buffer_timeout 设置输出缓冲区冲刷的最大时间频率(毫秒),存在三种逻辑模式:

取值行为适用场景
正整数按该毫秒数周期性地冲刷缓冲区默认行为,兼顾低延迟与平滑开发体验
0每处理一条记录立即冲刷,延迟最小对延迟极度敏感的任务
-1仅在输出缓冲区满时才冲刷,吞吐最大追求吞吐量的批式/大流量场景

配套的get_buffer_timeout()用于读取当前配置。测试 test_get_set_buffer_timeout 验证了 12000 毫秒的读写一致。

算子链(Operator Chaining):disable_operator_chaining

Flink 默认会把不存在 shuffle 的相邻算子链接(chaining)到同一个线程中执行,从而完全避免序列化/反序列化开销。调用disable_operator_chaining()可以关闭这一优化(见 disable_operator_chaining),is_chaining_enabled()用于查询当前是否启用。

测试 test_operation_chaining 验证了默认is_chaining_enabled()为True,调用disable_operator_chaining()后变为False。此外环境还提供is_chaining_of_operators_with_different_max_parallelism_enabled()查询"不同最大并行度的算子是否允许链接"(见 is_chaining_of_operators_with_different_max_parallelism_enabled)。

槽位共享组:register_slot_sharing_group 与 SlotSharingGroup

register_slot_sharing_group 将一个带资源规格(resource spec)的槽位共享组注册到环境中。其语义要点(源码注释原文):

槽位共享组只是提示调度器:组内算子可以被部署进同一个共享槽位,但不保证调度器一定将组内算子部署在一起;如果组内算子被部署到不同的槽位,槽位资源将依据组规格推导得出。

SlotSharingGroup及其构建器Builder定义在 slot_sharing_group.py,支持通过链式构建指定资源规格:

from pyflink.datastream import SlotSharingGroup from pyflink.datastream.slot_sharing_group import MemorySize group = SlotSharingGroup.builder("my_group") \ .set_cpu_cores(2.0) \ .set_task_heap_memory(MemorySize.of_mebi_bytes(256)) \ .set_task_off_heap_memory_mb(64) \ .set_managed_memory_mb(128) \ .set_external_resource("gpu", 1.0) \ .build() env.register_slot_sharing_group(group)

Builder提供的资源维度包括:CPU 核数(set_cpu_cores)、task heap 内存(set_task_heap_memory/set_task_heap_memory_mb)、task off-heap 内存(set_task_off_heap_memory/set_task_off_heap_memory_mb)、受管内存(set_managed_memory/set_managed_memory_mb)以及自定义外部资源(set_external_resource(name, value),同名旧值会被替换)。SlotSharingGroup对象本身提供get_name、get_cpu_cores、get_managed_memory、get_task_heap_memory、get_task_off_heap_memory、get_external_resources等读取方法。

MemorySize(slot_sharing_group.py)是对字节数的抽象表示,可通过MemorySize.of_mebi_bytes(n)或直接以字节数构造,并支持按get_bytes、get_kibi_bytes、get_mebi_bytes、get_gibi_bytes、get_tebi_bytes等单位读取。

容错与状态管理:Checkpoint、状态后端与重启策略

开启 Checkpoint:enable_checkpointing

enable_checkpointing 为流式作业开启周期性的状态快照:流式数据流的分布式状态会按给定间隔周期性快照,发生故障时,作业会从最近一次已完成的 checkpoint重启。

from pyflink.datastream import CheckpointingMode # 每 300 秒一次,至少一次语义 env.enable_checkpointing(300000, CheckpointingMode.AT_LEAST_ONCE) # 或只给间隔,默认使用 exactly-once 语义 env.enable_checkpointing(300000)
  • interval:两次状态 checkpoint 之间的时间间隔,单位毫秒;
  • mode:CheckpointingMode.EXACTLY_ONCE(默认)或CheckpointingMode.AT_LEAST_ONCE。

源码注释同时给出一个重要限制:当前并不正确支持迭代式流数据流(iterative dataflow)的 checkpoint,因此开启 checkpoint 时迭代式作业不会被启动。

配套读取 API:

  • get_checkpoint_interval():返回 checkpoint 间隔,未开启时返回 -1(等价于get_checkpoint_config().get_checkpoint_interval()的简写,见 get_checkpoint_interval);
  • get_checkpointing_mode()/get_checkpointing_consistency_mode():返回当前模式(exactly-once / at-least-once);
  • is_unaligned_checkpoints_enabled()、is_force_unaligned_checkpoints():查询非对齐 checkpoint(unaligned checkpoints)的启用状态与强制启用状态(见 is_unaligned_checkpoints_enabled)。

测试 test_get_set_checkpoint_interval 与 test_get_set_checkpointing_mode 分别验证了间隔 30000ms 的读写与模式切换。

细粒度 Checkpoint 配置:get_checkpoint_config

get_checkpoint_config()返回CheckpointConfig对象(checkpoint_config.py),用于进一步定制 checkpoint 行为。其默认常量(源码中定义):

常量默认值含义
DEFAULT_MODEEXACTLY_ONCE默认 checkpoint 模式
DEFAULT_TIMEOUT10 * 60 * 1000(10 分钟)单次 checkpoint 尝试的超时时间
DEFAULT_MIN_PAUSE_BETWEEN_CHECKPOINTS0两次 checkpoint 之间的最小暂停间隔(默认无)
DEFAULT_MAX_CONCURRENT_CHECKPOINTS1并发进行的 checkpoint 数量上限(默认 1 个)

可继续通过set_checkpointing_mode、set_checkpoint_interval、set_checkpoint_timeout、set_min_pause_between_checkpoints、set_max_concurrent_checkpoints等方法细化配置,也可以通过环境层configure()统一注入。

状态后端:set_state_backend / get_state_backend

状态后端决定了执行期间状态用何种数据结构存储(堆内存哈希表、RocksDB 或其他存储),以及checkpoint 数据持久化到哪里:

  • MemoryStateBackend:状态以对象形式保存在堆内存中,轻量、无额外依赖,但只能 checkpoint 较小状态(如计数器);
  • FsStateBackend:状态同样以堆对象维护,但 checkpoint 写入文件系统;配合复制型文件系统(HDFS、S3、Alluxio 等)可保证单节点故障不丢失状态,配合 HA 模式可实现高可用与强一致;
  • RocksDBStateBackend:使用 RocksDB 存储,适合大规模状态。
env.set_state_backend(EmbeddedRocksDBStateBackend())

注意:set_state_backend自 1.19 版本起标记为DeprecationWarning弃用,源码注释明确建议改用configure()方式设置状态后端(见 set_state_backend)。get_state_backend()返回当前状态后端(未设置时为None)。

Changelog 状态后端:enable_changelog_state_backend

enable_changelog_state_backend(1.14.0 新增)为当前状态后端启用变更日志(change log),其工作机制(源码注释 + FLIP-158):

  1. 有状态算子除了把状态变更应用到 RocksDB 状态表或内存哈希表外,还会将变更写入 changelog;
  2. 算子可以在变更日志到达持久化 checkpoint 存储的那一刻即确认 checkpoint,无需等待状态表落盘;
  3. 状态表独立于 checkpoint 周期性物化到 checkpoint 存储(materialization);
  4. 状态物化完成后,changelog 可被截断到对应位置。

这为跨状态后端大幅缩短流式应用的 checkpoint 间隔提供了途径。is_changelog_state_backend_enabled()返回Optional[bool]——若用户从未显式调用过该方法,返回None,此时 changelog 的启用与否由 job/local/cluster 各层配置决定。测试 test_enable_changelog_state_backend 验证了置True/False的读写行为。

默认 Savepoint 目录与重启策略

  • set_default_savepoint_directory(directory):设置默认 savepoint 目录——当触发 savepoint 时未显式提供路径,将写入该目录(get_default_savepoint_directory()可读取,未设置时返回None):
env.set_default_savepoint_directory("hdfs://savepoints")
  • set_restart_strategy(restart_strategy_configuration):设置作业失败后的重启策略配置:
from pyflink.common import RestartStrategies env.set_restart_strategy(RestartStrategies.no_restart())

get_restart_strategy()可读取当前重启策略。需要注意:set_restart_strategy与后面的序列化注册类方法一样,自 1.19 起被弃用,源码注释建议改用configure()注入配置。

时间语义:set_stream_time_characteristic

set_stream_time_characteristic(characteristic)为该环境创建的所有流设置时间特性,可选值为 time_characteristic.py 中定义的TimeCharacteristic枚举:

  • ProcessingTime:处理时间;
  • IngestionTime:摄入时间;
  • EventTime:事件时间(默认值,测试 test_get_set_stream_time_characteristic 中验证了默认即为EventTime)。
env.set_stream_time_characteristic(TimeCharacteristic.EventTime)

源码注释特别提示:若设置为IngestionTime或EventTime,会默认设置 200ms 的水位线更新间隔;如果你的应用不适用该默认值,应通过pyflink.common.ExecutionConfig.set_auto_watermark_interval修改。

Python 专属配置:依赖、解释器与资源文件

这部分是 PyFlink 相对 Java API 的差异化能力,全部通过操作底层PythonOptions配置键实现(源码中均使用PythonConfigUtil.getEnvironmentConfig获取环境配置)。

添加 Python 文件:add_python_file

env.add_python_file("my_udf.py") # 单文件 env.add_python_file("my_pkg/") # 目录

Python 文件、Python 包或本地目录会被加入Python UDF worker 的 PYTHONPATH,确保依赖可以被import。实现上会写入PythonOptions.PYTHON_FILES配置键,多个文件以FILE_DELIMITER分隔(见 add_python_file)。

第三方依赖:set_python_requirements

指定requirements.txt定义第三方依赖,依赖会被安装到临时目录并加入 UDF worker 的 PYTHONPATH;对于集群无法访问外网的场景,可额外指定本地安装包目录实现离线安装:

# shell 中准备 $ echo numpy==1.16.5 > requirements.txt $ pip download -d cached_dir -r requirements.txt --no-binary :all:
env.set_python_requirements("requirements.txt", "cached_dir")

注意事项(源码注释原文):

  • 安装包必须与集群平台及所用 Python 版本匹配;
  • 依赖通过 pip 安装,要求pip 版本 >= 20.3、setuptools 版本 >= 37.0.0。

虚拟环境归档:add_python_archive 与 set_python_executable

当集群缺少 UDF 所需的特定 Python 版本时,可上传整个虚拟环境:

$ zip -r py_env.zip py_env # 假设解释器相对路径为 py_env/bin/python
# 不指定 target_dir:解压到与归档同名的目录 env.add_python_archive("py_env.zip") env.set_python_executable("py_env.zip/py_env/bin/python") # 指定 target_dir="myenv":解压到 myenv 目录 env.add_python_archive("py_env.zip", "myenv") env.set_python_executable("myenv/py_env/bin/python") # 归档内的文件可通过相对路径在 UDF 中访问 def my_udf(): with open("myenv/py_env/data/data.txt") as f: ...

约束条件(见 add_python_archive 与 set_python_executable 的源码注释):

  • 上传的 Python 环境必须与集群平台匹配,Python 版本要求3.7 或更高;
  • 仅支持zip 格式归档(zip、jar、whl、egg 等),不支持tar、tar.gz、7z、rar;
  • UDF worker 依赖 Apache Beam(版本 == 2.43.0),指定的解释器环境需满足该要求。

JAR 与 Classpath:add_jars / add_classpaths

env.add_jars("file:///path/to/my.jar") env.add_classpaths("file:///path/to/dir/")
  • add_jars:上传 JAR 到集群并被作业引用,底层同时写入PipelineOptions.JARS配置键并加入 context class loader;
  • add_classpaths:添加 URL 到程序每个用户代码 classloader 的 classpath,路径必须指定协议(如file://)且在所有节点可访问,底层对应PipelineOptions.CLASSPATHS配置键。

数据接入:Source、集合与文件

经典 API:add_source

from pyflink.datastream.functions import SourceFunction ds = env.add_source(source_func, source_name='Custom Source', type_info=Types.ROW(...))

通过用户自定义的SourceFunction添加数据源,source_name与type_info均可选。

推荐 API:from_source(1.13.0+)

ds = env.from_source(source, watermark_strategy, source_name, type_info=None)

基于新 Source API 添加数据源。返回的 DataStream 是有界(可批处理)还是无界(必须流式处理),由 source 的 boundedness 属性决定。对于自身能描述产出类型的 source,不应再传type_info以避免冗余指定(见 from_source)。

测试/演示用数据:from_collection 与 read_text_file

# 不指定类型:元素以 pickled byte array 传输 ds = env.from_collection([(1, 'Hi', 'Hello'), (2, 'Hello', 'Hi')]) # 指定类型:结果以 Row 形式呈现 ds = env.from_collection(['Hi', 'Hello'], type_info=Types.STRING())

from_collection从集合创建 DataStream,其实现(见 from_collection / _from_collection)会把集合通过PickleSerializer序列化到临时文件,再经PythonBridgeUtils读取为字节数组或 Python 对象,最终构造一个forceNonParallel()(并行度为 1)的InputFormatSourceFunction有界 source。测试 test_from_collection_without_data_types 与 test_from_collection_with_data_types 分别验证了两种用法。

ds = env.read_text_file("file:///some/local/file", charset_name="UTF-8")

read_text_file逐行读取文件生成字符串 DataStream,支持file://、hdfs://等 URI。注意:该接口不具备故障容错,仅用于测试目的(见 read_text_file)。

其他数据接入能力

  • create_input(input_format, type_info=None)(1.16.0+):通过InputFormat创建输入流;当 input_format 需要明确定义类型信息(如 Avro 的 generic record)时,可显式传type_info,或使用实现了ResultTypeQueryable的 InputFormat 自动推断(见 create_input);
  • register_cached_file(file_path, name, executable=False)(1.16.0+):在分布式缓存中以指定名称注册文件,运行时任何 UDF 都可通过本地路径访问;文件可以是本地文件(经 BlobServer 分发)或分布式文件系统文件,运行时按需临时拷贝到本地缓存(见 register_cached_file)。

序列化注册:Kryo 序列化器与 POJO 类型

环境提供三个用于注册类型/序列化器的 API(均传入 Java 类的全限定名):

# 注册 Kryo 默认序列化器 env.add_default_kryo_serializer("com.aaa.bbb.TypeClass", "com.aaa.bbb.Serializer") # 为指定类型注册 Kryo 序列化器 env.register_type_with_kryo_serializer("com.aaa.bbb.TypeClass", "com.aaa.bbb.Serializer") # 注册类型:最终作为 POJO 序列化则注册 POJO 序列化器, # 若用 Kryo 序列化则注册到 Kryo 以只写入 tag env.register_type("com.aaa.bbb.TypeClass")

弃用说明:三个方法自 1.19 起均被标记弃用,源码注释给出的理由是"通过硬编码注册数据类型和序列化器,升级作业版本时需要修改代码",应改用配置项pipeline.serialization-config统一配置。测试 test_add_default_kryo_serializer、test_register_type_with_kryo_serializer、test_register_type 分别验证了这些注册会正确写入ExecutionConfig的对应结构中。

统一配置入口:configure(1.15.0+)

configure 是 1.15.0 引入的统一配置入口:读取传入Configuration中所有相关选项(如pipeline.time-characteristic),并同时重配置StreamExecutionEnvironment、ExecutionConfig与CheckpointConfig三层。其语义是:只有 configuration 中显式设置的键才会改变当前值,未出现的键保持原值不动。

测试 test_configure 展示了典型用法:

configuration = Configuration() configuration.set_string('pipeline.operator-chaining', 'false') configuration.set_string('pipeline.time-characteristic', 'IngestionTime') configuration.set_string('execution.buffer-timeout', '1 min') configuration.set_string('execution.checkpointing.timeout', '12000') env.configure(configuration) # 验证结果: # is_chaining_enabled() == False # get_stream_time_characteristic() == TimeCharacteristic.IngestionTime # get_buffer_timeout() == 60000("1 min" 被解析为 60000ms) # get_checkpoint_config().get_checkpoint_timeout() == 12000

这印证了 Flink 配置字符串的解析能力(如时长1 min→ 60000ms),也说明configure是替代上述多个弃用 API 的推荐路径。

作业执行:execute、execute_async 与 get_execution_plan

同步执行:execute

job_result = env.execute("my_flink_job") # job_name 可选 job_id = job_result.get_job_id()

execute 触发程序执行:环境会执行所有以sink 操作(打印结果、转发到消息队列等)结尾的部分,返回JobExecutionResult(含运行耗时与累加器)。实现上会先通过_generate_stream_graph生成 StreamGraph 再交给 Java 侧执行。

异步执行:execute_async

job_client = env.execute_async("Flink Streaming Job") # 默认作业名

execute_async 异步触发程序,立即返回JobClient(提交成功后即可用于与作业通信)。默认作业名为'Flink Streaming Job'。

执行计划:get_execution_plan

plan_json = env.get_execution_plan() # 必须在执行前调用

get_execution_plan 以JSON 字符串形式返回执行数据流图(StreamGraph)的计划,必须在计划执行之前调用;若无法实例化编译器或无法联系 master 获取执行规划所需信息,会抛出异常。

收尾:close(1.16.0+)

close()关闭并清理执行环境,物理释放所有缓存的中国结果(见 close)。

完整示例:组装一个可运行的 PyFlink 作业

将上述 API 串联成一个完整示例(数据来源为集合,便于本地运行验证):

from pyflink.common import Configuration, RestartStrategies, Types from pyflink.common.restart_strategy import RestartStrategies from pyflink.datastream import StreamExecutionEnvironment, CheckpointingMode from pyflink.datastream.execution_mode import RuntimeExecutionMode from pyflink.datastream.time_characteristic import TimeCharacteristic env = StreamExecutionEnvironment.get_execution_environment() # 1) 并行度 env.set_parallelism(4) env.set_max_parallelism(64) # 2) 执行模式(也可通过 execution.runtime-mode 配置项在提交时指定) env.set_runtime_mode(RuntimeExecutionMode.STREAMING) # 3) 时间语义与 buffer env.set_stream_time_characteristic(TimeCharacteristic.EventTime) env.set_buffer_timeout(100) # 4) 容错:每 5 秒一次 exactly-once checkpoint env.enable_checkpointing(5000, CheckpointingMode.EXACTLY_ONCE) env.set_default_savepoint_directory("hdfs:///flink/savepoints") env.set_restart_strategy(RestartStrategies.fixed_delay_restart(3, 10000)) # 5) Python 依赖 env.add_python_file("my_udf.py") env.set_python_requirements("requirements.txt", "cached_dir") env.set_python_executable("py_env.zip/py_env/bin/python") # 6) 数据接入 ds = env.from_collection([(1, 'Hello'), (2, 'Flink')], type_info=Types.ROW( [Types.INT(), Types.STRING()])) # 7) 执行(打印执行计划 / 提交作业) # print(env.get_execution_plan()) env.execute("pyflink_streaming_job")

结语

StreamExecutionEnvironment是 PyFlink DataStream 程序的起点与中枢。从本仓库的源码与测试可以看到,其 API 表面上是 Python 方法,底层全部映射到 Java 侧StreamExecutionEnvironment,因此 Java 生态的并行度上限(32768)、默认 checkpoint 语义(exactly-once)、配置项(execution.runtime-mode、pipeline.serialization-config)等约束在 PyFlink 中同样成立。理解并熟练运用并行度、RuntimeExecutionMode、Checkpoint/状态后端、configure统一配置、Python 依赖注入与作业执行这六类核心 API,即可掌控 PyFlink DataStream 作业从配置到提交的全生命周期。

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

【免费下载链接】flink

项目地址:https://gitcode.com/gh_mirrors/fli/flink
点击查看免费下载
上一篇:【亲测免费】 jVectorMap 项目推荐
下一篇:社交分享计数器:SocialCount 深度剖析

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

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/25 2:27:08

使用 Go 与 gqlgen 实现 GraphQL Mutation:createLink 解析器实战

【免费下载链接】howtographql The Fullstack Tutorial for GraphQL 项目地址&#xff1a; https://gitcode.com/gh_mirrors/ho/howtographql 点击查看 免费下载 本篇文章以 How to GraphQL 仓库中 GraphQL Go 后端教程的 Mutations 章节 为骨架&#xff0c;讲解 GraphQL Muta…

作者头像 李华
网站建设 2026/9/25 2:26:44

OpenMRS Standalone 2.3.1部署实战:从解压到MySQL切换与避坑

简介&#xff1a;OpenMRS 2.3.1 独立版是一份面向医疗机构、医疗信息化开发者及运维人员的开源电子病历系统部署包。它旨在帮助用户在资源有限的环境中快速搭建可定制的病历管理平台&#xff0c;涵盖患者登记、病史记录、诊断报告、处方与实验室结果管理等核心功能。压缩包内共…

作者头像 李华
网站建设 2026/9/25 2:25:26

Python跳动的爱心代码:tkinter动画循环与心形曲线实现

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华