
大数据流处理批处理数据工程【免费下载链接】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),仅供参考