- 大数据
- 流处理
- 批处理
- 数据工程
【免费下载链接】flink
本文以 Apache Flink 仓库中 PyFlink DataStream 参考文档 为主线,系统梳理 PyFlink DataStream API 的完整公开接口:从StreamExecutionEnvironment的作业构建与执行模式,到DataStream的转换算子、状态与定时器、窗口与检查点、侧输出,再到文件/Kafka/Pulsar/JDBC 等连接器与 Avro/CSV/JSON/Parquet 格式。读完本文,你将掌握 PyFlink DataStream 编程模型的全貌,并能在开发时按模块快速定位到具体的类与方法,配合源码理解底层行为。
文档定位:PyFlink DataStream API 的官方索引
flink-python/docs/reference/pyflink.datastream/目录是 PyFlink DataStream API 的官方参考文档集合,其中 index.rst 是总入口页,声明了该页面的定位——"This page gives an overview of all public PyFlink DataStream API",即覆盖所有公开的 PyFlink DataStream API。它通过 Sphinxtoctree指令组织出 10 个主题子文档:
| 子文档 | 主题 |
|---|---|
| stream_execution_environment.rst | 执行环境与作业控制 |
| datastream.rst | 数据流核心类型与算子 |
| functions.rst | 用户自定义函数 |
| state.rst | 状态与状态后端 |
| timer.rst | 定时服务与时间域 |
| window.rst | 窗口、触发器与窗口分配器 |
| checkpoint.rst | 检查点配置与存储 |
| sideoutput.rst | 侧输出标签 |
| connectors.rst | 各类外部系统连接器 |
| formats.rst | 序列化格式 |
下面按编程模型的生命周期顺序展开:先构建执行环境,再定义数据流算子,最后接入状态、时间与外部系统。
StreamExecutionEnvironment:流式程序的执行上下文
stream_execution_environment.rst 定义StreamExecutionEnvironment为"流式程序被执行的上下文":LocalStreamEnvironment会在附着的 JVM 内执行,RemoteStreamEnvironment会在远程集群上执行。它提供了两类能力:
- 控制作业执行:设置并行度、容错/检查点参数、缓冲超时、算子链、重启策略等;
- 与外部世界交互:接入数据源、注册 Python 依赖(文件、requirements、归档、解释器)、注册缓存文件等。
该环境暴露的方法覆盖了作业生命周期的每个阶段,按其职责可归纳为:
- 并行度与资源:
set_parallelism/get_parallelism、set_max_parallelism/get_max_parallelism、set_default_local_parallelism/get_default_local_parallelism、register_slot_sharing_group(配合SlotSharingGroup与MemorySize使用); - 执行模式与算子链:
set_runtime_mode、set_buffer_timeout/get_buffer_timeout、disable_operator_chaining/is_chaining_enabled; - 容错与状态:
enable_checkpointing/get_checkpoint_interval/get_checkpointing_mode、get_checkpoint_config、set_state_backend/get_state_backend、enable_changelog_state_backend/is_changelog_state_backend_enabled、set_default_savepoint_directory/get_default_savepoint_directory、set_restart_strategy/get_restart_strategy、is_unaligned_checkpoints_enabled/is_force_unaligned_checkpoints; - 类型与序列化:
add_default_kryo_serializer、register_type_with_kryo_serializer、register_type、set_stream_time_characteristic/get_stream_time_characteristic、configure; - Python 依赖注入:
add_python_file、set_python_requirements、add_python_archive、set_python_executable、add_jars、add_classpaths; - 数据接入:
create_input、add_source、from_source、read_text_file、from_collection、register_cached_file; - 作业提交:
execute、execute_async、get_execution_plan、close,以及便捷的get_execution_environment静态工厂。
其中enable_checkpointing、set_state_backend、set_restart_strategy等是每个生产作业的"标准动作",与后文 Checkpoint 小节 的配置项一一对应。
RuntimeExecutionMode:流批一体的执行模式
文档专门解释了RuntimeExecutionMode——它控制 DataStream 程序的运行时执行模式,影响任务调度、网络 Shuffle 行为和时间语义,部分算子也会依据执行模式改变记录发射行为。三种取值:
STREAMING:流水线以流式语义执行。所有任务在执行开始前全部部署,启用检查点,处理时间和事件时间都被完整支持;BATCH:流水线以批式语义执行。任务按所属调度区域渐进式调度,区域间的 Shuffle 是阻塞式的,watermark 被假定为"完美"(即不存在迟到数据),处理时间在执行期间被假定为不推进;AUTOMATIC:Flink 自动判定——若所有 Source 都是有界数据源则采用 BATCH,若存在至少一个无界数据源则采用 STREAMING。
这一枚举与 Java 侧RuntimeExecutionMode完全对齐,在 PyFlink 中通过StreamExecutionEnvironment.set_runtime_mode(...)设置,是"流批一体"能力在 DataStream API 上的直接体现。
DataStream 核心类型与算子
datastream.rst 是篇幅最大的子文档,它定义了 PyFlink DataStream API 的核心类型体系,其中DataStream被定义为"同类型元素的流"(a stream of elements of the same type)。其源码实现位于 pyflink/datastream/data_stream.py。
DataStream:基础流类型
DataStream的方法按用途分为三类:
元信息与算子属性(对应 Java DataStream 的算子描述能力):get_name/name(算子名)、uid/set_uid_hash(算子 UID,用于状态恢复时的稳定标识)、set_parallelism/set_max_parallelism、get_type、get_execution_environment、force_non_parallel、set_buffer_timeout、start_new_chain/disable_chaining(算子链控制)、slot_sharing_group、set_description。
转换算子(Transformation):map、flat_map、key_by、filter、window_all、union、connect、project、process、assign_timestamps_and_watermarks、add_sink、sink_to。这些算子构成流式处理的主体,例如key_by(data_stream.py#L357)接收一个 Callable 或KeySelector将流按 key 分组,返回KeyedStream;window_all(data_stream.py#L457)接收WindowAssigner返回AllWindowedStream;assign_timestamps_and_watermarks(data_stream.py#L666)接收WatermarkStrategy注入 watermark 生成逻辑。
数据重分区(Partitioning):shuffle、rescale、rebalance、forward、broadcast、partition_custom。这些算子只改变数据在并行子任务间的分发策略,不改变业务逻辑,是调整数据倾斜、控制网络传输方式的关键手段。
结果输出与旁路:execute_and_collect(data_stream.py#L839,支持job_execution_name与limit参数,用于本地调试时同步收集结果)、print、get_side_output、cache([data_stream.py#L899](https://link.gitcode.com/i/ba53962e20fcc4cc9b0bc95b26a87f11#L899),返回CachedDataStream)。
KeyedStream:按键分组的流
KeyedStream由key_by产生,文档描述为按 key 分组的流。它继承DataStream的全部能力,并新增了针对每个 key 独立计算的算子:reduce、sum、min、max、min_by、max_by(聚合)、window/count_window(窗口化)、process(KeyedProcessFunction,可访问 keyed state 与定时器)。由于数据已按 key 分区,这些算子的状态和窗口计算天然具备 per-key 隔离性。
WindowedStream 与 AllWindowedStream:窗口流
WindowedStream表示"元素按 key 分组,每个 key 的流再依据WindowAssigner切分为窗口"的流。文档特别强调两点:窗口按 key 分别求值(每个 key 的窗口可以在不同时刻触发);WindowedStream纯粹是 API 层面的构造,运行时会被折叠(collapse)成 KeyedStream 与窗口操作的一个单一算子。其方法包括trigger(自定义触发器)、allowed_lateness(允许的迟到时间)、side_output_late_data(迟到数据侧输出)、reduce、aggregate、apply、process。AllWindowedStream表示不按键分组的窗口流(全局窗口化)。文档额外说明:若指定了Evictor,它会在 Trigger 触发之后、窗口实际求值之前驱逐窗口内元素;使用 evictor 会显著降低窗口性能,因为无法使用窗口结果的预聚合。其方法集与WindowedStream一致(trigger、allowed_lateness、side_output_late_data、reduce、aggregate、apply、process),运行时会与窗口操作折叠为单一算子。
ConnectedStreams、BroadcastStream 与 BroadcastConnectedStream:多流关联
ConnectedStreams表示两条(可能类型不同)流的连接,常用于"一条流上的操作直接受另一条流影响"的场景,典型示例是将随时间变化的规则流作用到数据流上:一条流持有规则,另一条流持有待应用规则的元素,连接后的算子把当前规则集维护在状态中,收到规则更新则更新状态,收到数据元素则用状态中的规则处理它。概念上可看作Either类型的并集流。方法:key_by、map、flat_map、process。BroadcastStream用于广播状态的分发;BroadcastConnectedStream表示将 keyed 或非 keyed 流与携带BroadcastState的BroadcastStream连接的结果。文档给出的典型场景同样是动态规则:广播流持有规则并存入 broadcast state,从而在所有并行实例上可用,可以应用到另一条流的全部分区。其核心方法是process(配合BroadcastProcessFunction/KeyedBroadcastProcessFunction)。
DataStreamSink 与 CachedDataStream
DataStreamSink是"用于从流式拓扑发射元素的流式 Sink",提供name、uid、set_uid_hash、set_parallelism、set_description、disable_chaining、slot_sharing_group等算子属性设置方法;CachedDataStream表示"中间结果在首次计算时被缓存"的流,后续作业使用同一个CachedDataStream可复用缓存结果、避免重复计算(data_stream.py#L1774 的cache方法产生,data_stream.py#L1788 的invalidate方法使缓存失效)。它提供与DataStream几乎等价的算子集(map、flat_map、key_by、filter、window_all、union、connect、各类重分区、process、assign_timestamps_and_watermarks、add_sink、sink_to、execute_and_collect、print、get_side_output等),并额外提供invalidate。
Functions:用户自定义函数族
functions.rst 罗列了 DataStream API 中全部用户自定义函数(All user-defined functions)以及运行时上下文。
RuntimeContext是函数访问运行时信息的入口,提供:任务与子任务信息(get_task_name、get_number_of_parallel_subtasks、get_max_number_of_parallel_subtasks、get_index_of_this_subtask、get_attempt_number、get_task_name_with_subtasks)、作业参数get_job_parameter、指标组get_metrics_group,以及获取各类状态的方法:get_state(ValueState)、get_list_state、get_map_state、get_reducing_state、get_aggregating_state。
函数接口完整清单如下(与DataStream/KeyedStream的算子一一对应):
- 基础单流函数:
MapFunction、FlatMapFunction、FilterFunction、ReduceFunction、AggregateFunction、ProcessFunction、KeyedProcessFunction; - 双流连接函数:
CoMapFunction、CoFlatMapFunction、CoProcessFunction、KeyedCoProcessFunction; - 窗口函数:
WindowFunction、AllWindowFunction、ProcessWindowFunction、ProcessAllWindowFunction; - 广播函数:
BroadcastProcessFunction、KeyedBroadcastProcessFunction; - 辅助接口:
KeySelector(key 提取)、NullByteKeySelector、Partitioner(配合partition_custom)。
其中ProcessFunction系列是最灵活的底层抽象:它们能访问RuntimeContext、维护状态、注册定时器并向侧输出发射数据,是map/filter等简单函数无法覆盖复杂场景时的首选。
State:状态类型、TTL 与状态后端
state.rst 围绕"状态"展开三个层次:状态类型、状态描述符与 TTL 配置、状态后端。
状态类型与描述符
- 状态存储接口:
OperatorStateStore.get_broadcast_state(获取广播状态); - 状态值类型:
ValueState(单值)、AppendingState、MergingState、ReducingState(归约值)、AggregatingState(聚合值)、ListState(列表)、MapState(映射)、ReadOnlyBroadcastState/BroadcastState(广播状态); - 状态描述符:
ValueStateDescriptor、ListStateDescriptor、MapStateDescriptor、ReducingStateDescriptor、AggregatingStateDescriptor——描述符声明了状态的名字、类型与默认值,是RuntimeContext.get_state(...)等方法的入参。
StateTtlConfig:状态过期策略
StateTtlConfig是状态 TTL(Time To Live)的配置载体,文档通过autoclass展开了四个枚举/类:
UpdateType:状态过期时间的更新策略(创建与写入时更新update_ttl_on_create_and_write,或读写时都更新update_ttl_on_read_and_write);StateVisibility:过期状态对读的可见性(return_expired_if_not_cleaned_up返回未清理的过期数据,never_return_expired永不返回过期数据);TtlTimeCharacteristic:TTL 计时特征(默认使用处理时间use_processing_time);CleanupStrategies:过期数据清理策略——cleanup_full_snapshot(全量快照清理)、cleanup_incrementally(增量清理)、cleanup_in_rocksdb_compact_filter(RocksDB compaction filter 清理)、disable_cleanup_in_background(禁用后台清理)。
配置方式遵循 Builder 模式:new_builder创建构建器,通过set_update_type、set_state_visibility、set_ttl_time_characteristic、set_ttl等设置,最终build产出,查询方法包括get_update_type、get_state_visibility、get_ttl、get_ttl_time_characteristic、is_enabled、get_cleanup_strategies。
StateBackend:状态存储的底层机制
文档对StateBackend给出了较完整的原理阐述:
- 定义:状态后端决定流式应用的状态在集群内如何本地存储。
HashMapStateBackend把工作状态保存在 TaskManager 内存中,轻量且无额外依赖;EmbeddedRocksDBStateBackend把工作状态保存在内嵌 RocksDB 实例中,可扩展到 TB 级,仅受所有 TaskManager 可用磁盘空间限制; - 原始字节存储与后端:
StateBackend为原始字节存储(raw bytes storage)以及 keyed state、operator state 创建服务。由它创建的AbstractKeyedStateBackend与OperatorStateBackend决定 key 与算子的工作状态如何持有、如何通过CheckpointStreamFactory做检查点; - 可序列化性:状态后端需要实现
java.io.Serializable,因为要与流式应用代码一起分发到并行进程,因此实现被设计为轻量工厂(只含配置),由工厂创建真正的状态存储; - 线程安全:多个线程可能并发创建流与 keyed/operator 状态后端,因此实现必须线程安全。
PyFlink 暴露的后端类包括:HashMapStateBackend、EmbeddedRocksDBStateBackend,以及为兼容旧版本保留的MemoryStateBackend、FsStateBackend、RocksDBStateBackend、CustomStateBackend(自定义后端)和PredefinedOptions(预定义 RocksDB 选项)。
Timer:定时服务与时间域
timer.rst 定义了时间相关的三类 API:
- TimerService(定时服务):提供
current_processing_time(当前处理时间)、current_watermark(当前 watermark)、register_processing_time_timer(注册处理时间定时器)、register_event_time_timer(注册事件时间定时器)、delete_processing_time_timer/delete_event_time_timer(删除对应定时器)。定时器通常通过KeyedProcessFunction/ProcessFunction的上下文获取,是事件时间迟到处理、会话超时检测等场景的基石; - TimeCharacteristic(时间特征):对应流处理的时间语义设置;
- TimeDomain(时间域):区分处理时间(ProcessingTime)与事件时间(EventTime)两种定时器注册域。
Window:窗口、触发器与窗口分配器
window.rst 完整覆盖窗口机制的三大构件:
- Window 类型:
TimeWindow(时间窗口,含起止时间戳)、CountWindow(计数窗口)、GlobalWindow(全局窗口,永不自然触发); - Trigger(触发器):决定窗口何时求值/清理,包括
TriggerResult(触发结果枚举)、EventTimeTrigger/ContinuousEventTimeTrigger、ProcessingTimeTrigger/ContinuousProcessingTimeTrigger、PurgingTrigger(清理型触发器)、CountTrigger、NeverTrigger(永不触发); - WindowAssigner(窗口分配器):决定元素进入哪个窗口,包括
MergingWindowAssigner(可合并窗口分配器的基类)、计数类CountTumblingWindowAssigner/CountSlidingWindowAssigner、时间类TumblingProcessingTimeWindows/TumblingEventTimeWindows/SlidingProcessingTimeWindows/SlidingEventTimeWindows、会话类ProcessingTimeSessionWindows/EventTimeSessionWindows/DynamicProcessingTimeSessionWindows/DynamicEventTimeSessionWindows、以及GlobalWindows; - SessionWindowTimeGapExtractor:为动态会话窗口提供"会话间隙"提取器,可从元素中动态计算会话时间间隙。
这些类与KeyedStream.window(...)/DataStream.window_all(...)配合使用,例如滚动窗口TumblingEventTimeWindows.of(Time.minutes(5))、滑动窗口SlidingProcessingTimeWindows.of(Time.minutes(10), Time.minutes(1))、会话窗口EventTimeSessionWindows.with_gap(Time.minutes(5))。
Checkpoint:检查点配置与 CheckpointStorage
checkpoint.rst 定义检查点相关的配置与存储机制。
CheckpointConfig 默认值与配置方法
CheckpointConfig是"捕获所有检查点相关设置的配置",文档明确给出四个默认常量:
| 常量 | 含义 | 默认值 |
|---|---|---|
DEFAULT_MODE | 默认检查点模式 | 精确一次(exactly once) |
DEFAULT_TIMEOUT | 单次检查点尝试的超时 | 10 分钟 |
DEFAULT_MIN_PAUSE_BETWEEN_CHECKPOINTS | 两次检查点之间的最小暂停 | 无(none) |
DEFAULT_MAX_CONCURRENT_CHECKPOINTS | 并发检查点数量上限 | 1 |
配置方法覆盖完整生命周期:开关与模式(is_checkpointing_enabled、get/set_checkpointing_mode、get/set_checkpoint_interval)、超时与节奏(get/set_checkpoint_timeout、get/set_min_pause_between_checkpoints、get/set_max_concurrent_checkpoints)、错误容忍(is/set_fail_on_checkpointing_errors、get/set_tolerable_checkpoint_failure_number)、外部化检查点(enable_externalized_checkpoints、set_externalized_checkpoint_cleanup、is_externalized_checkpoints_enabled、get_externalized_checkpoint_cleanup,配合ExternalizedCheckpointCleanup/ExternalizedCheckpointRetention枚举)、非对齐检查点(is_unaligned_checkpoints_enabled、enable/disable_unaligned_checkpoints、set/get_alignment_timeout、set/is_force_unaligned_checkpoints)、存储设置(set/get_checkpoint_storage、set_checkpoint_storage_dir)。
CheckpointStorage:检查点持久化策略
文档说明检查点存储定义了StateBackend如何为流式应用存储容错状态:
JobManagerCheckpointStorage:把检查点保存在 JobManager 内存中,轻量、无额外依赖,但不可扩展、仅支持小状态,适合本地测试与开发;FileSystemCheckpointStorage:把检查点保存在文件系统(HDFS、NFS、S3、GCS 等),支持数 TB 级大状态并为流式应用提供高可用基础,推荐用于大多数生产部署。
与 StateBackend 类似,CheckpointStorage也遵循"轻量工厂"设计:实现需可序列化(java.io.Serializable)以随应用代码分发,需线程安全以支持并发创建流;原始字节存储服务(通过CheckpointStreamFactory)被 JobManager 用于存储检查点与恢复元数据,也通常被 keyed/operator 状态后端用于存储检查点状态。PyFlink 暴露JobManagerCheckpointStorage、FileSystemCheckpointStorage与CustomCheckpointStorage。
Side Outputs:侧输出标签
sideoutput.rst 讲解侧输出机制的核心——OutputTag,它被定义为"用于标记算子侧输出的、带类型与名字的标签"。文档给出了可直接运行的示例,展示了 PyFlink 的三种构造方式与一个约束:
# 显式指定输出类型 >>> info = OutputTag("late-data", Types.TUPLE([Types.STRING(), Types.LONG()])) # 隐式将 list 包装为 Types.ROW >>> info_row = OutputTag("row", [Types.STRING(), Types.LONG()]) # 隐式使用 pickle 序列化 >>> info_side = OutputTag("side") # ERROR: tag id 不能为空字符串(Python API 的额外要求) >>> info_error = OutputTag("")注意最后一行展示的额外约束:PyFlink 要求 tag id 不能为空字符串,这是 Python API 特有的校验,Java API 无此限制。侧输出典型用于处理迟到数据:在窗口算子中通过side_output_late_data(tag)把迟到元素发射到侧输出流,再通过DataStream.get_side_output(tag)取回单独处理。
Connectors:外部系统连接器
connectors.rst 覆盖 PyFlink DataStream API 支持的全部连接器,是接入外部数据源/目标的清单:
- File System:文件 Source(
FileSource/FileSourceBuilder,配套FileEnumeratorProvider、FileSplitAssignerProvider、StreamFormat、BulkFormat);文件 Sink(FileSink、StreamingFileSink,配套BucketAssigner(分桶策略)、RollingPolicy/DefaultRollingPolicy/OnCheckpointRollingPolicy(滚动策略)、OutputFileConfig、FileCompactStrategy/FileCompactor(文件压缩)); - Number Sequence:
NumberSequenceSource,用于生成有界数字序列,常见于压测与测试场景; - Kafka:两代 API 并存——旧版
FlinkKafkaConsumer/FlinkKafkaProducer(含Semantic语义枚举)与新版KafkaSource/KafkaSourceBuilder(配套KafkaTopicPartition、KafkaOffsetResetStrategy、KafkaOffsetsInitializer)以及KafkaSink/KafkaSinkBuilder(配套KafkaRecordSerializationSchema/KafkaRecordSerializationSchemaBuilder、KafkaTopicSelector); - Kinesis:Source(
FlinkKinesisConsumer,配套KinesisShardAssigner、KinesisDeserializationSchema、WatermarkTracker);Sink(KinesisStreamsSink/KinesisStreamsSinkBuilder、KinesisFirehoseSink/KinesisFirehoseSinkBuilder,配套PartitionKeyGenerator); - Pulsar:Source(
PulsarSource/PulsarSourceBuilder,配套StartCursor、StopCursor、RangeGenerator);Sink(PulsarSink/PulsarSinkBuilder,配套TopicRoutingMode、MessageDelayer); - JDBC:
JdbcSink,配套JdbcConnectionOptions(连接配置)与JdbcExecutionOptions(执行配置); - RabbitMQ(RMQ):
RMQSource/RMQSink,配套RMQConnectionConfig; - Elasticsearch:
ElasticsearchSink,配套Elasticsearch6SinkBuilder/Elasticsearch7SinkBuilder、ElasticsearchEmitter、FlushBackoffType(flush 退避策略); - Cassandra:
CassandraSink,配套ConsistencyLevel、MapperOptions、ClusterBuilder、CassandraCommitter、CassandraFailureHandler。
Formats:序列化格式
formats.rst 定义 DataStream API 可用的数据格式:
- Avro:
AvroSchema、GenericRecordAvroTypeInfo、AvroInputFormat、AvroBulkWriters、AvroRowDeserializationSchema/AvroRowSerializationSchema; - CSV:
CsvSchema/CsvSchemaBuilder、CsvReaderFormat、CsvBulkWriters、CsvRowDeserializationSchema/CsvRowSerializationSchema; - JSON:
JsonRowDeserializationSchema/JsonRowSerializationSchema; - ORC:
OrcBulkWriters(列式批量写入); - Parquet:
AvroParquetReaders/AvroParquetWriters、ParquetColumnarRowInputFormat、ParquetBulkWriters。
其中*RowDeserializationSchema/*RowSerializationSchema系列用于连接器(如 Kafka)的 record 编解码,*BulkWriters系列用于FileSink的列式/批量写出。
实践要点:如何基于该参考定位开发
综合上述参考文档与源码,给出几条可直接落地的使用建议:
- 按编程模型顺序组织代码:
StreamExecutionEnvironment.get_execution_environment()→ 通过from_source/from_collection/read_text_file接入数据 → 用key_by/map/process等转换 → 用window/assign_timestamps_and_watermarks处理时间与窗口 → 用sink_to/add_sink输出 → 最后execute提交。每个环节的 API 都能在本参考的对应子文档中找到。 - 状态与检查点配套使用:
set_state_backend(如HashMapStateBackend/EmbeddedRocksDBStateBackend)与set_checkpoint_storage(如FileSystemCheckpointStorage)分别决定工作状态与检查点落盘方式,生产环境推荐后者;TTL 需求用StateTtlConfig.new_builder()链式配置后传入StateDescriptor。 - 时间语义先行:事件时间处理必须先通过
assign_timestamps_and_watermarks注入WatermarkStrategy,窗口才具备正确的迟到处理能力;迟到的元素可用side_output_late_data+OutputTag单独捕获(tag id 不能为空)。 - 调试期快捷方式:本地验证用
execute_and_collect()同步收集结果、print()打印到标准输出、NumberSequenceSource生成测试数据;跨作业复用中间结果用cache()返回的CachedDataStream,并在不需要时调用invalidate()释放缓存。 - 流批模式按数据源自动切换:
set_runtime_mode(RuntimeExecutionMode.AUTOMATIC)可让 Flink 依据 Source 有界/无界自动选择 BATCH/STREAMING 语义,但需注意两种模式下时间与调度行为的差异(见 RuntimeExecutionMode 一节)。
本文所有 API 清单均可在仓库内核实:参考文档位于 flink-python/docs/reference/pyflink.datastream/,核心实现位于 flink-python/pyflink/datastream/data_stream.py,Python 侧还可在 flink-python/pyflink/datastream/ 下找到 functions、state、window、connectors 等对应模块;对应的 Java 侧实现与测试可分别在 flink-streaming-java/src/main 与 flink-python/src/test 中继续深入。
- 大数据
- 流处理
- 批处理
- 数据工程
【免费下载链接】flink
相关推荐
Flink DataStream API 编程指南:从执行环境到流式应用的完整实战
Flink DataStream API 编程指南:从执行环境到流式应用的完整实战 Flink DataStream API 是 Apache Flink 中面
大数据流处理批处理数据工程Flink DataStream API 编程指南:从执行环境到流式窗口 WordCount 的完整实战
Flink DataStream API 编程指南:从执行环境到流式窗口 WordCount 的完整实战 DataStream API 是 Apache Fli
大数据流处理批处理数据工程PyFlink Python DataStream API 入门指南:从环境构建到作业提交的完整实践
PyFlink Python DataStream API 入门指南:从环境构建到作业提交的完整实践 本文以 Apache Flink 仓库中的官方文档 doc
大数据流处理批处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考