PyFlink DataStream 实战:用 Python 优雅处理 JSON 流式数据
大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载本文以 Apache Flink 的 Python 接口 PyFlink 为例讲解如何在 DataStream API 中直接解析与加工 JSON 字符串流从from_collection构造内存数据源到map中完成json.loads解析与字段修改、再到filter中基于嵌套 JSON 路径做条件过滤最终通过print()输出到标准输出。读完本文你将掌握一条可直接复制运行的 JSON 处理流水线并理解其背后的类型推断、算子封装与作业提交机制可用于快速搭建 PyFlink 的 JSON 数据清洗、字段增强与规则过滤任务。示例文档与示例代码的对应关系在仓库中本主题由两个文件共同构成文档入口flink-python/docs/examples/datastream/process_json_data.rst它通过.. literalinclude:: ../../../pyflink/examples/datastream/process_json_data.py指令将下方示例脚本的源码全文嵌入文档正文示例脚本flink-python/pyflink/examples/datastream/process_json_data.py即文档展示的完整可运行程序。该示例同时被收录进 DataStream 示例目录的索引页 flink-python/docs/examples/datastream/index.rst与word_count、basic_operations、timer、state、window、connectors并列属于 PyFlink 官方示例集的一部分随 flink-python/setup.py 打包分发pyflink.examples包及*.py、*/*.py数据文件均会随 pip 安装一并发布。完整示例一条 JSON 解析与过滤流水线以下是文档对应的核心示例程序完整保留源码逻辑并补充关键注释import json import logging import sys from pyflink.datastream import StreamExecutionEnvironment def process_json_data(): env StreamExecutionEnvironment.get_execution_environment() # 定义数据源一个包含 4 条 (id, json字符串) 元组的内存集合 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}})] ) def update_tel(data): # 解析 JSON 字符串 json_data json.loads(data[1]) json_data[tel] 1 return data[0], json_data def filter_by_country(data): # 此处 data[1] 已是解析后的 dict无需再次 json.loads return China in data[1][addr][country] ds.map(update_tel).filter(filter_by_country).print() # 提交作业执行 env.execute() if __name__ __main__: logging.basicConfig(streamsys.stdout, levellogging.INFO, format%(message)s) process_json_data()程序逻辑分四步通过StreamExecutionEnvironment.get_execution_environment()创建流执行环境用env.from_collection()把内存中的 4 条(id, json字符串)元组变成有界 DataStream链式调用map(update_tel)解析 JSON 并将tel加 1与filter(filter_by_country)保留addr.country为China的记录用print()把结果输出到标准输出最后env.execute()触发作业执行。预期输出如下tel已被 1且只保留中国区记录(2, {name: hello, tel: 136, addr: {country: China, city: Shanghai}}) (4, {name: PyFlink, tel: 33, addr: {country: China, city: Hangzhou}})运行方式与仓库中其他示例如 word_count.py、streaming_word_count.py一致本示例采用if __name__ __main__:入口风格可以直接以 Python 脚本方式运行python process_json_data.py也可以从仓库源码目录直接运行python flink-python/pyflink/examples/datastream/process_json_data.py运行前提已安装 PyFlink可通过pip install apache-flink或从仓库 flink-python 目录按 setup.py 配置安装首次运行时 PyFlink 会自动拉起本地 MiniCluster 完成作业调度与执行无需额外启动集群脚本末尾的logging.basicConfig(streamsys.stdout, levellogging.INFO, format%(message)s)用于把日志输出到 stdout 并精简格式方便在终端直接观察运行日志。关键 API 逐层拆解StreamExecutionEnvironment.get_execution_environment()get_execution_environment()是 PyFlink 流式作业的入口其底层实现见 stream_execution_environment.py它通过get_gateway()获取 Java 网关调用 Java 侧org.apache.flink.streaming.api.environment.StreamExecutionEnvironment的静态方法getExecutionEnvironment()再包装成 Python 的StreamExecutionEnvironment对象。若以独立脚本运行它默认返回本地执行环境若从命令行客户端提交则会把传入的Configuration叠加到全局配置config.yaml之上形成优先级更高的作业配置。from_collection构造内存数据源from_collection(collection, type_infoNone)的签名与语义见 stream_execution_environment.py未指定type_info时元素类型不做显式声明该方法创建的集合数据源是非并行源parallelism 为 1源码中通过forceNonParallel()强制适用于本地测试与小数据量场景底层执行链路_from_collection会把集合用 Pickle 序列化到临时文件再经PythonBridgeUtils读出为字节数组列表包装成InputFormatSourceFunction有界源Boundedness.BOUNDED挂到执行环境上。在本示例中type_info未传因此流中每个元素是一个 Python 元组(id, json字符串)后续算子直接按元组下标data[0]、data[1]访问即可。map解析 JSON 并修改字段DataStream.map的实现在 data_stream.py。它把普通 Python 函数适配成ProcessFunctionMapProcessFunctionAdapter每个输入元素恰好产出一个输出元素yield self._map_func(value)并命名为Map算子。示例中的update_tel正是利用这一点def update_tel(data): json_data json.loads(data[1]) # 将 JSON 字符串解析为 Python dict json_data[tel] 1 # 原地修改 tel 字段 return data[0], json_data # 返回 (id, dict)dict 会作为流元素继续流转值得注意json.loads返回的 dict 会直接作为流元素在算子间传递因此后续filter中data[1]已经是解析后的 dict无需再次调用json.loads示例注释也明确强调了这一点。若未指定output_typemap 输出默认按 Pickle 字节数组序列化。filter基于嵌套 JSON 路径过滤DataStream.filter的实现在 data_stream.py。它同样把 Python 函数适配为ProcessFunctionFilterProcessFunctionAdapter仅当函数返回True时yield value保留元素否则丢弃算子命名为Filter。示例中的过滤条件为def filter_by_country(data): return China in data[1][addr][country]这里data[1]是上一步 map 产出的 dict通过[addr][country]直接沿着嵌套路径取到国家字段展示了在 PyFlink 中“JSON 已结构化”后的便捷访问方式。print输出到标准输出DataStream.print(sink_identifierNone)的实现在 data_stream.py对应 Java 侧的print()sink把每个元素的字符串表示写到 stdout。需要注意print()输出发生在执行该作业的 Flink worker 机器上且该 sink 不具备容错能力仅适合调试与示例验证生产场景应替换为 FileSink、Kafka Sink 等可靠 Sink可参考 streaming_word_count.py 中FileSink.for_row_format(...)的用法。env.execute提交执行env.execute(job_nameNone)的实现在 stream_execution_environment.py它会先生成 StreamGraph_generate_stream_graph再交给 Java 执行环境执行返回JobExecutionResult含运行耗时与累加器。PyFlink 还提供异步提交版本execute_async()返回JobClient便于在提交后继续与作业交互。若在本地运行作业会在脚本结束时随之退出因此示例使用同步execute()保证结果完整打印。延伸一显式声明元素类型Row 方式同目录下的 basic_operations.py 展示了与本例几乎相同的数据集在“显式类型声明”下的写法通过type_infoTypes.ROW_NAMED([id, info], [Types.INT(), Types.STRING()])把流元素声明为命名 Row之后便可用属性访问data.info替代下标data[1]例如from pyflink.common import Types ... ds env.from_collection( collection[(1, {name: Flink, tel: 123, ...}), ...], type_infoTypes.ROW_NAMED([id, info], [Types.INT(), Types.STRING()]) ) def update_tel(data): json_data json.loads(data.info) json_data[tel] 1 return data.id, json.dumps(json_data)可见是否显式指定type_info决定了后续访问元素的方式不指定时走 Pickle 元组指定后走结构化 Row。两种风格在本仓库示例中并存可按需选用。延伸二Table API 的 JSON 处理对照同一主题在 PyFlink Table API 中也有对应示例 table/process_json_data.py核心差异在于用声明式 SQL 表达式完成 JSON 取值table table.select(col(id), col(data).json_value($.addr.country, DataTypes.STRING()))它通过from_elements定义数据、TableDescriptor.for_connector(print)定义 sink并使用内建函数json_value配合 JSONPath 表达式$.addr.country提取嵌套字段最后table.execute_insert(sink).wait()提交。与 DataStream 的“Python 函数自由解析”相比Table API 更接近 SQL 语义适合声明式场景而本文的 DataStream 方案则保留了完整的 Python 编程自由度可任意组合json标准库逻辑两者互为补充。该组示例的文档索引见 flink-python/docs/examples/table/process_json_data.rst。小结与进一步阅读通过本文你已经掌握了一条可运行的 PyFlink JSON 处理流水线内存集合数据源 →map内json.loads解析并修改 →filter按嵌套 JSON 路径过滤 →print输出 →env.execute提交。示例完整源码位于 process_json_data.py底层 API 实现在 stream_execution_environment.py 与 data_stream.py 中可进一步查阅。若需继续深入建议依次阅读同目录下的其他示例word_count.pyflat_map、key_by、reduce与 FileSource/FileSink 读写basic_operations.pyTypes.ROW_NAMED显式类型、key_by后聚合streaming_word_count.pydatagen 数据源与文件 Sink 的完整配置table/process_json_data.pyTable API 的json_value声明式 JSON 解析。这些示例全部收录于 flink-python/docs/examples/datastream/index.rst 的文档目录中可作为从示例走向生产实践的阶梯。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐PyFlink DataStream Formats 完整指南Avro、CSV、JSON、ORC、Parquet 数据格式的读写实战PyFlink DataStream Formats 完整指南Avro、CSV、JSON、ORC、Parquet 数据格式的读写实战 PyFlink 的 py大数据流处理批处理数据工程PyFlink DataStream Word Count 双模式实战从批处理到流式统计的完整实现PyFlink DataStream Word Count 双模式实战从批处理到流式统计的完整实现 导读 本文围绕 Flink Python DataStre大数据流处理批处理数据工程PyFlink DataStream Side Outputs 完全指南使用 OutputTag 处理侧输出流PyFlink DataStream Side Outputs 完全指南使用 OutputTag 处理侧输出流 导读 在 PyFlink DataStream大数据流处理批处理数据工程上一篇DreamCraft3D延迟体积渲染技术深度解析下一篇如何快速上手Kudu5分钟搭建Git部署环境创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

相关新闻

google-images-download 工作流深度解析:从关键词输入到图片批量落盘的完整算法链路

google-images-download 工作流深度解析:从关键词输入到图片批量落盘的完整算法链路

网页爬虫CLI 【免费下载链接】google-images-download Python Script to download hundreds of images from Google Images. It is a ready-to-run code! 项目地址: https://gitcode.com/gh_mirrors/go/google-images-download 点击查看 免费下载 本文以仓库文档 d…

2026/9/25 2:44:18 阅读更多 →
swagger-codegen 生成的 Jersey1 Java 客户端 StoreApi 使用指南:Petstore 订单与库存接口实战

swagger-codegen 生成的 Jersey1 Java 客户端 StoreApi 使用指南:Petstore 订单与库存接口实战

开发工具代码生成API设计 【免费下载链接】swagger-codegen swagger-codegen contains a template-driven engine to generate documentation, API clients and server stubs in different languages by parsing your OpenAPI / Swagger definition. 项目地址: http…

2026/9/25 2:44:18 阅读更多 →
高通Camera PDAF调试:Type2到Type3迁移的五个关键环节

高通Camera PDAF调试:Type2到Type3迁移的五个关键环节

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

2026/9/25 2:43:18 阅读更多 →

最新新闻

Word表格自动上浮与跨页断行问题的根源与解决

Word表格自动上浮与跨页断行问题的根源与解决

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

2026/9/25 6:44:15 阅读更多 →
AfKayAs.2远控木马深度解析:从样本结构到检测规则

AfKayAs.2远控木马深度解析:从样本结构到检测规则

拿到这个样本的时候,我习惯性地先看了一眼文件哈希,然后在沙箱里丢了一把。AfKayAs.2这个名字,在威胁情报社区里其实不算陌生,它是某个远控木马家族的升级变种,前一代AfKayAs.1曾经在不少攻防演练和真实攻击场景里出现…

2026/9/25 6:44:15 阅读更多 →
Atlas 300V 24G部署YOLOv5实战:从模型转换到推理调优

Atlas 300V 24G部署YOLOv5实战:从模型转换到推理调优

1. 先说结论:Atlas 300V 24G到底是什么卡做AI应用这两年,总有人问我类似的选型问题:“预算有限,想上国产推理卡,Atlas 300V 24G能不能买?”“它到底算不算一张运算加速卡,还是只是个带显存的视频…

2026/9/25 6:44:15 阅读更多 →
Atlas 300V部署YOLOv5全流程实战:从环境搭建到性能调优

Atlas 300V部署YOLOv5全流程实战:从环境搭建到性能调优

1. 先搞清楚一件事:Atlas 300V到底是不是"运算加速卡"先回应那个热搜词——很多人拿到Atlas 300V,第一反应是"这玩意是不是类似一张NVIDIA显卡?能不能直接拿来跑CUDA?"答案是:能跑推理&#xff0c…

2026/9/25 6:44:15 阅读更多 →
FPGA配置Flash烧录与擦除实操指南:SPI协议、JEDEC命令与Vivado实战

FPGA配置Flash烧录与擦除实操指南:SPI协议、JEDEC命令与Vivado实战

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

2026/9/25 6:44:15 阅读更多 →
CTF入门实战复盘:从图片隐写到栈溢出的解题思路

CTF入门实战复盘:从图片隐写到栈溢出的解题思路

SUSCTF 2018那场比赛的周末,我是从一道Misc题开始的。当时刚入CTF圈不久,最大的感受是:题目不会按你“擅长”的来,但如果你能把每道题的思路记录下来,后面进步会很快。这篇做题记录不是完整题解,更像是我个…

2026/9/25 6:43:14 阅读更多 →

日新闻

AI元人文:从工具使用到思维重构的深度探索

AI元人文:从工具使用到思维重构的深度探索

最近半年我一直在琢磨一件事:AI元人文到底是什么?说白了,就是“用元视角重新审视人与AI的关系”,也在“探索AI如何反向逼着我们发现自己的思考边界”。标题里的“元探索”,在我看就是一层套一层的追问——当你用AI解决…

2026/9/25 0:00:41 阅读更多 →
Python+CNN车牌识别实战:从数据预处理到模型训练与部署

Python+CNN车牌识别实战:从数据预处理到模型训练与部署

简介:基于Python与卷积神经网络的车牌识别项目,面向计算机视觉初学者及智能交通开发者,目标是帮助用户掌握从数据预处理、模型构建到实际部署的完整流程。压缩包共25个文件,包含jpg/png图像样本、py训练脚本、md说明文档、dat数据…

2026/9/25 0:00:41 阅读更多 →
Vim基础操作全攻略:保存退出、模式切换与高频命令实战

Vim基础操作全攻略:保存退出、模式切换与高频命令实战

1. 项目概述1.1 核心需求解析今天聊聊Vim。写这个题目的原因是:几乎每个后端开发者、运维人员、数据工程师某天都会遇到一个场景——深夜加班,服务器登录界面只有黑底白字,编辑器只有vi/vim,你必须在五分钟内完成一次配置修改并保…

2026/9/25 0:00:41 阅读更多 →

周新闻

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

直接铺开项目本身吧。这几个月我一直在折腾一件事:用Flutter给OpenHarmony做一款游戏集合类的App,说白了就是把若干小游戏塞进一个壳里,用统一入口分发。这个方向本身不算新鲜,真正让我花了不少心思的,是首页那堆游戏卡…

2026/9/24 14:34:13 阅读更多 →
Word表格编号全攻略:从列表编号到题注交叉引用

Word表格编号全攻略:从列表编号到题注交叉引用

写Word文档,最让人头疼的往往是那些“看起来不起眼”的小问题。比如表格编号这事:今天在表后面多加了两个空白行,明天给客户交稿前发现整个章节的编号全部错位,光是挨个改序号就能耗掉大半个下午。我前阵子帮人整理一份上百页的技…

2026/9/24 9:10:42 阅读更多 →
从第一个站到第二个站:独立开发者的静态网站选型与落地实践

从第一个站到第二个站:独立开发者的静态网站选型与落地实践

1. 项目概述1.1 核心需求解析做独立开发者这几年,说实话,第一个网站上线的那天晚上我兴奋得没睡着。但等它跑了半年,流量惨淡、功能臃肿、代码自己都懒得看第二遍之后,我才慢慢琢磨明白一个道理:第一个网站是练手&…

2026/9/24 14:33:56 阅读更多 →

月新闻

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能分类:[AI/大模型]细分主题:AI 增强型 CI/CD 流水线自动化与 GitOps 实践:Agent 工作流、工具调用与任务拆解:从原型到生产的验收清单很多团队在尝试用大…

2026/9/24 12:50:34 阅读更多 →
容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场分类:[工程技术]细分主题:Kubernetes 生产环境运维与排障实战:可复制的项目复盘模板与决策记录大部分团队的事故复盘报告,最后都变成了躺在 Confluence 或钉…

2026/9/24 14:33:48 阅读更多 →
容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步分类:[工程技术]细分主题:Docker 容器化技术与镜像安全管理:核心链路的逐步实现与关键代码取舍面对一个积累了五六年历史包袱的单体架构应用(包含 Web 接口、后台…

2026/9/24 12:49:17 阅读更多 →