- 大数据
- 流处理
- 批处理
- 数据工程
【免费下载链接】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_MODE | EXACTLY_ONCE | 默认 checkpoint 模式 |
DEFAULT_TIMEOUT | 10 * 60 * 1000(10 分钟) | 单次 checkpoint 尝试的超时时间 |
DEFAULT_MIN_PAUSE_BETWEEN_CHECKPOINTS | 0 | 两次 checkpoint 之间的最小暂停间隔(默认无) |
DEFAULT_MAX_CONCURRENT_CHECKPOINTS | 1 | 并发进行的 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):
- 有状态算子除了把状态变更应用到 RocksDB 状态表或内存哈希表外,还会将变更写入 changelog;
- 算子可以在变更日志到达持久化 checkpoint 存储的那一刻即确认 checkpoint,无需等待状态表落盘;
- 状态表独立于 checkpoint 周期性物化到 checkpoint 存储(materialization);
- 状态物化完成后,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
相关推荐
PyFlink Python DataStream API 入门指南:从环境构建到作业提交的完整实践
PyFlink Python DataStream API 入门指南:从环境构建到作业提交的完整实践 本文以 Apache Flink 仓库中的官方文档 doc
大数据流处理批处理数据工程Flink Python DataStream API 入门指南:从环境构建到作业提交的完整实战
Flink Python DataStream API 入门指南:从环境构建到作业提交的完整实战 导读 本文以 Apache Flink 官方仓库中 intro
大数据流处理批处理数据工程PyFlink TableEnvironment 全面指南:环境创建、API 总览与配置实战
PyFlink TableEnvironment 全面指南:环境创建、API 总览与配置实战 PyFlink 的 TableEnvironment 是 Tabl
大数据流处理批处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考