ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

PyFlink DataStream Word Count 双模式实战:从批处理到流式统计的完整实现

PyFlink DataStream Word Count 双模式实战:从批处理到流式统计的完整实现 大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载导读本文围绕 Flink Python DataStream 示例文档 展开深入剖析 PyFlink 官方提供的两个 Word Count 示例批处理版 word_count.py 与流式版 streaming_word_count.py。读者将掌握如何使用StreamExecutionEnvironment构建数据流管道如何通过FileSource/from_collection/datagen三种方式定义数据源如何用flat_map、map、key_by、reduce完成分词与统计以及如何用FileSink或print()输出结果。同时结合 PyFlink 源码说明每个关键 API 的底层行为与适用场景帮助读者写出可直接运行、可扩展到真实业务的 PyFlink 作业。示例文档与代码定位word_count.rst是 PyFlink 官方文档 DataStream 示例索引 中的第一篇通过literalinclude指令直接嵌入两个示例的完整源码这意味着文档即代码仓库中 pyflink/examples/datastream 目录下的源文件是文档的唯一权威内容来源。示例文件处理模式数据来源核心知识点word_count.py批处理BATCH内存集合或文本文件FileSource、RuntimeExecutionMode.BATCH、FileSinkstreaming_word_count.py流处理STREAMINGdatagen连接器持续生成数据Table 与 DataStream 互转、无限流处理一、批处理 Word Count经典分词统计的完整实现批处理示例的核心代码位于 word_count.py主函数word_count(input_path, output_path)接受输入与输出路径两个可选参数完整构建了一条「读取 → 分词 → 计数 → 输出」的数据流管道。1.1 准备执行环境与运行模式env StreamExecutionEnvironment.get_execution_environment() env.set_runtime_mode(RuntimeExecutionMode.BATCH) # write all the data to one file env.set_parallelism(1)StreamExecutionEnvironment.get_execution_environment()是 PyFlink 所有 DataStream 作业的入口定义在 stream_execution_environment.py。RuntimeExecutionMode.BATCH表示以批语义执行任务全部部署完成后才开始执行事件时间与处理时间按批语义处理。枚举定义见 execution_mode.py。与之相对的是STREAMING模式——所有任务先部署、开启 checkpoint完整支持处理时间与事件时间。set_parallelism(1)将整个作业的并行度设为 1注释明确说明目的是把所有数据写入同一个文件保证输出结果集中、便于查看。1.2 三种输入路径文件读取与内存集合示例对input_path做了分支处理演示了两种数据源定义方式if input_path is not None: ds env.from_source( sourceFileSource.for_record_stream_format(StreamFormat.text_line_format(), input_path) .process_static_file_set().build(), watermark_strategyWatermarkStrategy.for_monotonous_timestamps(), source_namefile_source ) else: print(Executing word_count example with default input data set.) print(Use --input to specify file input.) ds env.from_collection(word_count_data)方式一FileSource读取文本文件。StreamFormat.text_line_format(charset_nameUTF-8)按行读取文件底层委托给 Java 的TextLineInputFormat使用java.io.InputStreamReader按指定字符集解码字节流参见 file_system.py。FileSource.for_record_stream_format(...)构建按记录流方式读取的源process_static_file_set()将其设置为有界批模式只处理启动时已存在的文件全部处理完后作业结束monitor_continuously(interval)则相反会持续监控新文件参见 file_system.py。WatermarkStrategy.for_monotonous_timestamps()为无时间戳的纯文本输入提供单调递增的水位线策略。批模式下该示例并没有真正使用事件时间这里主要演示 API 的完整拼装方式。方式二from_collection读取内存数据。当未指定--input时示例内置了莎士比亚《哈姆雷特》To be, or not to be 独白的 32 行文本作为默认数据word_count.py非常适合开箱即用地验证作业逻辑。1.3 核心转换链分词 → 映射 → 分组 → 归约def split(line): yield from line.split() # compute word count ds ds.flat_map(split) \ .map(lambda i: (i, 1), output_typeTypes.TUPLE([Types.STRING(), Types.INT()])) \ .key_by(lambda i: i[0]) \ .reduce(lambda i, j: (i[0], i[1] j[1]))这一链式调用是 Word Count 的精髓四个算子各司其职flat_map(split)split是生成器函数yield from line.split()将每一行按空白字符拆成单词并逐一发射实现一行 → 多词的展平。map(lambda i: (i, 1), output_type...)把每个单词映射为(单词, 1)二元组。output_type显式声明类型为Types.TUPLE([Types.STRING(), Types.INT()])即字符串与整数的元组这是 PyFlink 类型推断的重要环节——Lambda 表达式无法自动推断类型时必须显式指定。key_by(lambda i: i[0])按单词本身分组Flink 会依据 key 的哈希将相同单词路由到同一并行子任务这是后续增量累加的前提。reduce(lambda i, j: (i[0], i[1] j[1]))对同一 key 下的二元组两两归约把计数累加。由于reduce是基于 key 的有状态算子相同单词的计数会在其所在分区内持续累积最终得到每个单词的总出现次数。1.4 结果输出FileSink 与 stdout 双通道if output_path is not None: ds.sink_to( sinkFileSink.for_row_format( base_pathoutput_path, encoderEncoder.simple_string_encoder()) .with_output_file_config( OutputFileConfig.builder() .with_part_prefix(prefix) .with_part_suffix(.ext) .build()) .with_rolling_policy(RollingPolicy.default_rolling_policy()) .build() ) else: print(Printing result to stdout. Use --output to specify output path.) ds.print()FileSink.for_row_format(base_path, encoder)按行格式写出Encoder.simple_string_encoder()将每条记录转为字符串后写入参见 file_system.py。OutputFileConfig.builder().with_part_prefix(prefix).with_part_suffix(.ext)控制输出文件命名分片文件会以prefix为前缀、以.ext为后缀。RollingPolicy.default_rolling_policy()采用默认滚动策略决定文件何时滚动成新文件例如按文件大小或写入间隔参见 file_system.py。若未指定--output则调用ds.print()直接把结果打到标准输出便于本地快速调试。最后env.execute()提交作业执行。1.5 命令行入口if __name__ __main__: logging.basicConfig(streamsys.stdout, levellogging.INFO, format%(message)s) parser argparse.ArgumentParser() parser.add_argument(--input, destinput, requiredFalse, helpInput file to process.) parser.add_argument(--output, destoutput, requiredFalse, helpOutput file to write results to.) argv sys.argv[1:] known_args, _ parser.parse_known_args(argv) word_count(known_args.input, known_args.output)使用argparse解析--input与--output两个可选参数并用parse_known_args忽略 Flink 平台注入的其他参数。典型运行方式# 使用内置默认数据结果打印到 stdout python pyflink/examples/datastream/word_count.py # 指定输入文件结果写入输出目录 python pyflink/examples/datastream/word_count.py --input /path/to/input.txt --output /path/to/out二、Streaming Word Count无限数据流上的持续统计流式示例 streaming_word_count.py 展示了另一个维度数据源是一个持续产生数据的无限流因此统计永远不会结束这是与批处理版的本质区别。2.1 基于 datagen 连接器构建无限数据源words [flink, window, timer, event_time, processing_time, state, connector, pyflink, checkpoint, watermark, sideoutput, sql, datastream, broadcast, asyncio, catalog, batch, streaming] max_word_id len(words) - 1示例内置 18 个 Flink / PyFlink 领域词汇作为字典。数据源通过datagen 连接器构造env StreamExecutionEnvironment.get_execution_environment() t_env StreamTableEnvironment.create(stream_execution_environmentenv) # define the source # randomly select 5 words per second from a predefined list t_env.create_temporary_table( source, TableDescriptor.for_connector(datagen) .schema(Schema.new_builder() .column(word_id, DataTypes.INT()) .build()) .option(fields.word_id.kind, random) .option(fields.word_id.min, 0) .option(fields.word_id.max, str(max_word_id)) .option(rows-per-second, 5) .build()) table t_env.from_path(source) ds t_env.to_data_stream(table)这里出现了一个重要的跨 API 协作模式用StreamTableEnvironment.create(stream_execution_environmentenv)在已有的 DataStream 执行环境之上创建 Table 环境。通过TableDescriptor.for_connector(datagen)声明一个内置的 datagen 数据生成连接器表字段word_id为INT类型fields.word_id.kindrandom表示随机取值min0、maxmax_word_id限定取值范围rows-per-second5控制每秒生成 5 行数据即每秒从词汇表中随机挑 5 个词。最后t_env.to_data_stream(table)把 Table 转回 DataStream完成从 Table API 到 DataStream API 的桥接后续即可沿用 DataStream 的算子链做统计。2.2 把 word_id 翻译成单词并统计def id_to_word(r): # word_id is the first column of the input row return words[r[0]] # compute word count ds ds.map(id_to_word) \ .map(lambda i: (i, 1), output_typeTypes.TUPLE([Types.STRING(), Types.INT()])) \ .key_by(lambda i: i[0]) \ .reduce(lambda i, j: (i[0], i[1] j[1]))id_to_word取输入行的第一个字段即word_id作为下标把数字翻译回词汇表中的单词随后的map → key_by → reduce链与批处理版完全一致。由于数据源源源不断reduce会持续维护每个单词的累计计数任何时刻的状态都代表到目前为止的统计结果——这正是流式 Word Count 与批处理 Word Count 在语义上的关键差异。2.3 输出与入口if output_path is not None: ds.sink_to( sinkFileSink.for_row_format( base_pathoutput_path, encoderEncoder.simple_string_encoder()) .with_output_file_config( OutputFileConfig.builder() .with_part_prefix(prefix) .with_part_suffix(.ext) .build()) .with_rolling_policy(RollingPolicy.default_rolling_policy()) .build() ) else: print(Printing result to stdout. Use --output to specify output path.) ds.print()输出逻辑与批处理版共用同一套FileSink配置行格式编码、prefix/.ext文件命名、默认滚动策略。命令行仅提供--output一个可选参数# 默认打印到 stdout观察持续刷新的计数结果 python pyflink/examples/datastream/streaming_word_count.py # 写入文件 python pyflink/examples/datastream/streaming_word_count.py --output /path/to/out三、两个示例的对比与学习要点维度word_count.py批处理streaming_word_count.py流式运行模式RuntimeExecutionMode.BATCH默认流式语义数据源FileSource文本文件或from_collection内存datagen连接器无限生成数据是否有限有限处理完自动结束无限持续运行统计语义全量统计后输出一次增量累计状态持续演进涉及额外 APIWatermarkStrategy、StreamFormatStreamTableEnvironment、TableDescriptor、to_data_stream输出方式FileSink或print()与批处理完全相同两个示例共享的核心统计算子链map → key_by → reduce正是 DataStream 编程模型中最具代表性的模式key_by是状态与并行度的结合点相同 key 的数据被路由到同一个并行实例使reduce可以安全地在本地维护每个 key 的累加状态无需全局协调。显式output_type是 PyFlink 的实践要点Lambda 表达式的返回类型无法被自动推断必须用Types.TUPLE([Types.STRING(), Types.INT()])显式声明否则作业提交阶段会因类型信息缺失而失败。四、源码层面的纵深理解4.1 FileSource 的批/流双模式FileSourceBuilderfile_system.py提供两个互斥的读取模式方法process_static_file_set()有界模式处理启动时路径下已存在的文件全部完成后源即结束示例批处理版使用monitor_continuously(discovery_interval)无界模式按固定间隔扫描新文件并持续读取。此外StreamFormat.text_line_format()的源码注释揭示了两个底层事实file_system.py使用 Java 内置InputStreamReader按字符集解码默认 UTF-8不支持 checkpoint 优化恢复恢复时会重读并丢弃上次 checkpoint 之前已处理的行数因为字符集解码器的内部缓冲状态无法精确定位行偏移。而FileSource.for_record_stream_format还支持按文件扩展名自动解压.deflate、.xz、.bz2、.gz、.gzip等压缩格式file_system.py这让示例稍作改动即可直接处理压缩输入。4.2 RuntimeExecutionMode 对行为的影响execution_mode.py 对RuntimeExecutionMode的说明指出运行模式不仅影响任务调度方式还会影响网络 shuffle 行为、时间语义以及部分算子的记录发射行为。其中BATCH任务先全部部署再执行适合有界输入STREAMING任务边部署边执行开启 checkpoint完整支持处理时间与事件时间适合无界输入。这也解释了为何批处理示例会主动set_parallelism(1)来保证输出单一文件——在批模式下并行度会显著影响分片文件的生成数量。4.3 datagen 连接器的参数语义流式示例通过TableDescriptor为 datagen 表配置了四个关键 optionOption示例值含义fields.word_id.kindrandom字段生成方式为随机值fields.word_id.min0随机取值下限fields.word_id.max17len(words)-1随机取值上限rows-per-second5每秒生成的数据行数由于是随机取值每秒生成的 5 个词中必然存在重复key_by reduce的累计效果会随运行时间不断增长——这是演示流式增量统计最直观的方式。五、运行前置条件与延伸阅读运行这两个示例需要已安装 PyFlink 及其依赖Py4J、CloudPickle、python-dateutil、Apache Beam 等详见 flink-python/README.md。本地开发可通过如下方式验证环境# 在仓库根目录下构建并安装 PyFlink 后即可运行示例 python flink-python/pyflink/examples/datastream/word_count.py示例位于 flink-python/pyflink/examples/datastream 目录同目录下还有 basic_operations.pymap/filter/key_by 基础操作、process_json_data.pyJSON 处理、state_access.py状态访问、event_time_timer.py事件时间与定时器、windowing窗口等进阶示例对应文档索引见 DataStream 示例总览。这些示例与本文的 Word Count 共享同一套环境构建、算子链与 FileSink 输出模式是继续深入 PyFlink DataStream API 的理想起点。结语从 word_count.rst 出发本文完整还原了 PyFlink DataStream 的两个官方 Word Count 实现批处理版展示了FileSource有界读取、内存集合输入与全量统计流式版展示了 datagen 无限数据源、Table/DataStream 桥接与增量累计。两者共享的map → key_by → reduce算子链是理解 Flink 状态化流处理的核心范式。掌握了这两个示例你就拥有了构建更复杂 PyFlink 作业窗口聚合、状态管理、多源连接的坚实基础。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐PyFlink DataStream API 实战教程从零构建一个 Python 流式词频统计作业PyFlink DataStream API 实战教程从零构建一个 Python 流式词频统计作业 Apache Flink 的 DataStream API大数据流处理批处理数据工程从模糊建议到精确数值skills项目证据而非品味的设计工程审查哲学全解析从模糊建议到精确数值skills项目证据而非品味的设计工程审查哲学全解析 skills 是一个 AI 智能体技能Agent Skills集合项目为Flink DataStream API 编程指南从执行环境到流式应用的完整实战Flink DataStream API 编程指南从执行环境到流式应用的完整实战 Flink DataStream API 是 Apache Flink 中面大数据流处理批处理数据工程上一篇【72小时限时】6语言情感分析API化指南从BERT模型到生产级服务的零成本落地下一篇旧Mac免费装新系统OCLP 3 道准入线、5 步实操让十年老机器跑上最新 macOS创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表