Flink Python REPL(pyflink-shell)实战指南:local / remote / YARN 模式与 Table API 交互式开发
后端大数据流处理批处理【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载Flink 仓库在部署章节REPLs中提供了一款集成的交互式 Python ShellPython REPL它既可以在本机以 local 模式启动也可以连接到远程或 YARN 集群运行让开发者无需编写完整工程即可逐行验证 Table API 逻辑。本文以该章节的官方文档为主体完整覆盖 Python Shell 的安装、使用、四种部署模式与全部命令行参数并结合仓库中的启动脚本与源码pyflink-shell.sh、PythonShellParser、shell.py剖析其底层工作原理读完即可在本地或集群上独立开展交互式开发与调试。什么是 Flink Python ShellFlink 附带了一个集成的交互式 Python Shell。它既能够运行在本地启动的 local 模式也能够运行在集群启动的 cluster 模式下。启动后Table Environment 的相关内容会被自动加载开发者可以直接在提示符下编写并执行 Table API 语句非常适合用于快速原型验证、教学演示和日常调试。当前 Python Shell 主要支持 Table API 的功能。启动之后Table Environment 的相关内容会被自动加载可以通过变量bt_env来使用 BatchTableEnvironment通过变量st_env来使用 StreamTableEnvironment。核心的预绑定逻辑位于 flink-python/pyflink/shell.py启动时会创建s_envStreamExecutionEnvironment与st_envStreamTableEnvironment。环境要求与安装Python Shell 会调用python命令因此需要预先配置好 Python 执行环境。关于 Python 执行环境的要求请参考 Python Table API 环境安装PyFlink 需要 Python 3.7 以上版本文档标注 3.8、3.9 或 3.10可运行python --version确认版本也可以参考 faq.md 中的虚拟环境方案或通过 python.client.executable 与 python.executable 配置指定 Python 解释器路径。本地安装 Flink 请参考 本地安装Standalone 资源提供者也可以从源码构建 Flink详见 从源码构建 Flink。安装好 PyFlink 之后即可直接使用 Python Shell# 安装 PyFlink $ python -m pip install apache-flink # 执行脚本local 模式 $ pyflink-shell.sh local关于如何在一个 Cluster 集群上运行 Python Shell可以参考下文启动章节的介绍。快速上手启动 Shell 与预绑定环境在安装好 PyFlink 的前提下执行$ pyflink-shell.sh local启动过程会依次完成 Flink 安装目录定位、classpath 与依赖 zip 的组装并最终以交互模式加载pyflink.shell模块随后打印欢迎信息其中包含类似下面的提示NOTE: Use the prebound Table Environment to implement batch or streaming Table programs. Streaming - Use s_env and st_env variables也就是说Shell 启动后已经为你准备好s_env/st_env等环境变量无需再手动创建 ExecutionEnvironment 或 TableEnvironment可以直接进入 Table API 编程环节。Table API 交互式编程示例下面是通过 Python Shell 运行的简单示例分别对应流stream与批batch两种场景。流式场景stream import tempfile import os import shutil sink_path tempfile.gettempdir() /streaming.csv if os.path.exists(sink_path): ... if os.path.isfile(sink_path): ... os.remove(sink_path) ... else: ... shutil.rmtree(sink_path) s_env.set_parallelism(1) t st_env.from_elements([(1, hi, hello), (2, hi, hello)], [a, b, c]) st_env.create_temporary_table(stream_sink, TableDescriptor.for_connector(filesystem) ... .schema(Schema.new_builder() ... .column(a, DataTypes.BIGINT()) ... .column(b, DataTypes.STRING()) ... .column(c, DataTypes.STRING()) ... .build()) ... .option(path, path) ... .format(FormatDescriptor.for_format(csv) ... .option(field-delimiter, ,) ... .build()) ... .build()) t.select(col(a) 1, col(b), col(c))\ ... .execute_insert(stream_sink).wait() # 如果作业运行在 local 模式, 你可以执行以下代码查看结果: with open(os.path.join(sink_path, os.listdir(sink_path)[0]), r) as f: ... print(f.read())批式场景batch import tempfile import os import shutil sink_path tempfile.gettempdir() /batch.csv if os.path.exists(sink_path): ... if os.path.isfile(sink_path): ... os.remove(sink_path) ... else: ... shutil.rmtree(sink_path) b_env.set_parallelism(1) t bt_env.from_elements([(1, hi, hello), (2, hi, hello)], [a, b, c]) st_env.create_temporary_table(batch_sink, TableDescriptor.for_connector(filesystem) ... .schema(Schema.new_builder() ... .column(a, DataTypes.BIGINT()) ... .column(b, DataTypes.STRING()) ... .column(c, DataTypes.STRING()) ... .build()) ... .option(path, path) ... .format(FormatDescriptor.for_format(csv) ... .option(field-delimiter, ,) ... .build()) ... .build()) t.select(col(a) 1, col(b), col(c))\ ... .execute_insert(batch_sink).wait() # 如果作业运行在 local 模式, 你可以执行以下代码查看结果: with open(os.path.join(sink_path, os.listdir(sink_path)[0]), r) as f: ... print(f.read())两点使用提示示例通过TableDescriptor.for_connector(filesystem)创建临时表并将结果以 CSV 格式FormatDescriptor.for_format(csv)field-delimiter为,写入临时目录在 local 模式下作业运行完毕后可直接用 Python 文件读取代码打印结果文件内容。官方文档示例中create_temporary_table(...).option(path, path)里的path应替换为实际定义的sink_path变量shell.py 内置的演示代码使用的正是sink_path读者按自己的变量名填写即可。启动部署模式详解查看 Python Shell 提供的全部可选参数可以使用pyflink-shell.sh --helpLocal 模式Python Shell 运行在 local 模式下只需要执行pyflink-shell.sh local该模式会在本机启动一个内嵌的 Flink mini clusterlocal cluster适合本地开发调试。Remote 模式Python Shell 运行在一个指定的 JobManager 上通过关键字remote和对应的 JobManager 的地址与端口号来进行指定pyflink-shell.sh remote hostname portnumber其中hostname为 JobManager 的主机名portnumber为 JobManager 的端口号。从源码看该模式会被转换为-m hostname:portnumber以连接远程 JobManager见 PythonShellParser.java。Yarn Python Shell cluster 模式Python Shell 可以运行在 YARN 集群之上Python Shell 会在 YARN 上部署一个新的 Flink 集群并进行连接。除了指定 container 数量你也可以指定 JobManager 的内存、YARN 应用的名字等参数。例如在一个部署了两个 TaskManager 的 YARN 集群上运行 Python Shellpyflink-shell.sh yarn -n 2关于所有可选的参数可以查看本文完整参考部分的说明。Yarn Session 模式如果你已经通过 Flink Yarn Session 部署了一个 Flink 集群能够通过以下的命令连接到这个集群pyflink-shell.sh yarn注意yarn模式在未提供额外参数时连接已存在的 Yarn Session而yarn -n 2这类携带资源的调用则会在 YARN 上新建一个专属于 Shell 的集群。完整参考命令行参数一览Flink Python Shell 使用: pyflink-shell.sh [local|remote|yarn] [options] args... 命令: local [选项] 启动一个部署在 local 的 Flink Python shell 使用: -h,--help 查看所有可选的参数 命令: remote [选项] host port 启动一个部署在 remote 集群的 Flink Python shell host JobManager 的主机名 port JobManager 的端口号 使用: -h,--help 查看所有可选的参数 命令: yarn [选项] 启动一个部署在 Yarn 集群的 Flink Python Shell 使用: -h,--help 查看所有可选的参数 -jm,--jobManagerMemory arg 具有可选单元的 JobManager 的 container 的内存默认值MB) -n,--container arg 需要分配的 YARN container 的 数量 (TaskManager 的数量) -nm,--name arg 自定义 YARN Application 的名字 -qu,--queue arg 指定 YARN 的 queue -s,--slots arg 每个 TaskManager 上 slots 的数量 -tm,--taskManagerMemory arg 具有可选单元的每个 TaskManager 的 container 的内存默认值MB -h | --help 打印输出使用文档各 YARN 参数含义归纳如下参数短选项说明--jobManagerMemory-jmJobManager Container 的内存可带单位默认单位为 MB--container-n需要分配的 YARN Container 数量等于 TaskManager 数量--name-nm自定义 YARN Application 的名字--queue-qu指定 YARN 队列--slots-s每个 TaskManager 上的 slot 数量--taskManagerMemory-tm每个 TaskManager Container 的内存可带单位默认单位为 MB--help-h打印使用文档源码剖析Python Shell 是如何工作的启动脚本 pyflink-shell.shPython Shell 的入口脚本位于 flink-python/bin/pyflink-shell.sh其主要执行流程如下通过 find-flink-home.sh 定位FLINK_HOMEpip 安装场景下会调用find_flink_home.py动态解析加载 config.sh 构造 Flink 运行时 classpath并定位flink-python*.jar与pyflink.zip、py4j-*-src.zip、cloudpickle-*-src.zip将其加入PYTHONPATH调用 Java 程序org.apache.flink.client.python.PythonShellParser解析命令行参数把解析结果以 NUL 分隔的方式写回 shell并导出为SUBMIT_ARGS最终执行${PYFLINK_PYTHON} -i -m pyflink.shell进入交互模式-i表示 interactive-m表示以模块方式执行 zip 包中的shell.py。其中PYFLINK_PYTHON环境变量默认为python可用它来指定 Python 解释器。参数解析器 PythonShellParserPythonShellParser.java 是命令行参数解析的核心它定义了三种集群类型常量local、remote、yarn以及-h/--help、-jm/--jobManagerMemory、-nm/--name、-qu/--queue、-s/--slots、-tm/--taskManagerMemory等选项。三种模式的解析与转换逻辑分别为local直接输出local让flink run使用本机 mini cluster 执行作业remote要求至少提供hostname与portnumber转换为-m hostname:portnumber连接远程 JobManageryarn转换为-m yarn-cluster并把 Python Shell 的 yarn 选项加上前缀y对齐flink run的 YARN 参数例如-jm 1024m转换为-yjm 1024m、-tm 4096m转换为-ytm 4096m。上述行为由单元测试 PythonShellParserTest.java 直接验证testParseLocalWithoutOptions、testParseRemoteWithoutOptions、testParseYarnWithoutOptions、testParseYarnWithOptions分别断言了三种模式及带参场景下转换出的命令选项。Shell 核心实现 shell.pyshell.py 是 Python Shell 运行时模块启动时会一次性导入pyflink.common、pyflink.datastream、pyflink.table、pyflink.table.catalog、pyflink.table.descriptors、pyflink.table.window、pyflink.metrics等子包并打印 Python 版本与 ASCII 欢迎横幅随后创建s_env StreamExecutionEnvironment.get_execution_environment() st_env StreamTableEnvironment.create(s_env)因此 Shell 内可直接使用s_env/st_env变量。文档示例中的bt_env/b_env对应历史版本中 BatchTableEnvironment 的预绑定变量实际使用时以当前环境中真实存在的预绑定变量为准。常见问题与注意事项Python 解释器要求Python Shell 会调用python命令请确保其版本满足 PyFlink 环境要求见 安装文档如需指定解释器可设置PYFLINK_PYTHON环境变量。结果查看local 模式下作业执行完毕后可以像示例那样读取输出目录中的结果文件进行验证若在 remote / YARN 模式下运行结果会落在对应集群的文件系统路径上。参数错误提示若未指定集群类型或传入了非法参数PythonShellParser会输出错误信息并提示合法的集群类型为local、remote hostname portnumber、yarn。功能范围当前 Python Shell 聚焦 Table API 场景适合交互式验证表操作、连接器配置与作业提交无需每次改动都重启一个完整工程。赞分享后端大数据流处理批处理【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Flink Python REPL 完全指南用 PyFlink Shell 交互式开发 Table API 作业Flink Python REPL 完全指南用 PyFlink Shell 交互式开发 Table API 作业 Flink 自带一个集成的交互式 Pytho后端大数据流处理批处理Flink Python REPL 完全指南用 pyflink-shell 交互式编写 Table API 程序Flink Python REPL 完全指南用 pyflink shell 交互式编写 Table API 程序 PyFlink 内置了一个开箱即用的交互式后端大数据流处理批处理Apache Flink PyFlink Table API 指南TableDescriptor、FormatDescriptor、Schema 与 ChangelogMode 详解Apache Flink PyFlink Table API 指南TableDescriptor、FormatDescriptor、Schema 与 Chan后端大数据流处理批处理上一篇WABT 上游提案测试目录 test/spec-new 解析以 wide-arithmetic 测试集为例下一篇Windows看不到iPhone照片免费HEIC缩略图插件三分钟安装告别灰色图标创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

相关新闻

为什么你的 Node 测试会卡死?skills 项目教你诊断 flaky 测试与悬挂进程的完整指南

为什么你的 Node 测试会卡死?skills 项目教你诊断 flaky 测试与悬挂进程的完整指南

为什么你的 Node 测试会卡死?skills 项目教你诊断 flaky 测试与悬挂进程的完整指南 【免费下载链接】skills My own collection of skills for modern Node.js development 项目地址: https://gitcode.com/gh_mirrors/skills15/skills skills 是一个面向现代…

2026/10/10 2:28:59 阅读更多 →
如何读懂 NPUSim 指令流水图:Perfetto 可视化操作与关键字段全解

如何读懂 NPUSim 指令流水图:Perfetto 可视化操作与关键字段全解

如何读懂 NPUSim 指令流水图:Perfetto 可视化操作与关键字段全解 【免费下载链接】npu-simulator NPUSim(全称NPU Simulator)是一款面向算子开发场景的SoC级芯片仿真工具,用于分析运行在AI仿真器上的AI任务在各阶段的精度和性能数…

2026/10/10 2:28:59 阅读更多 →
developer-roadmap 中的 Angular 路由转换动画(Route Transitions)实战指南:为页面切换注入流畅过渡

developer-roadmap 中的 Angular 路由转换动画(Route Transitions)实战指南:为页面切换注入流畅过渡

文档教程知识库 【免费下载链接】developer-roadmap Interactive roadmaps, guides and other educational content to help developers grow in their careers. 项目地址: https://gitcode.com/GitHub_Trending/de/developer-roadmap 点击查看 免费下载 当用户从一…

2026/10/10 2:28:59 阅读更多 →

最新新闻

SpringBoot3多数据源实战:从选型配置到避坑指南

SpringBoot3多数据源实战:从选型配置到避坑指南

做后端这些年,只要业务稍微复杂一点,“一个应用连一个库”的理想状态基本撑不住。用户数据放用户库、订单数据放订单库、日志又要独立一套,再加上读写分离和多租户隔离的需求,所有问题都指向同一个核心:一个SpringBoot…

2026/10/10 3:20:15 阅读更多 →
07_在k8s集群中安装ingress-nginx

07_在k8s集群中安装ingress-nginx

文章目录0 背景与架构0.1 为什么需要ingress-nginx0.2 环境架构0.3 镜像拉取方案1 安装前准备1.1 下载ingress-nginx部署文件1.2 替换镜像地址1.3 修复DNS解析问题1.3.1 问题现象1.3.2 修复方法2 安装ingress-nginx2.1 应用部署文件2.2 验证部署3 测试Ingress3.1 部署测试应用3…

2026/10/10 3:20:15 阅读更多 →
基于Spring Boot的宠物用品商城系统开发与实践指南

基于Spring Boot的宠物用品商城系统开发与实践指南

毕业设计季总有人拿着“宠物用品系统”这个题目来找我聊,说白了这个题在Java Web方向的选题清单里几乎年年出现。还有一类人不是学生,是宠物店老板想把自己那套手工记账和微信群接单的流程搬到线上。这两类人虽然出发点不一样,但最终要的东西…

2026/10/10 3:20:15 阅读更多 →
河北省口碑好的母婴护理服务商盘点,资质齐全的正规家政机构不踩坑

河北省口碑好的母婴护理服务商盘点,资质齐全的正规家政机构不踩坑

在衡水,找月嫂、请育婴师,几乎是每个新手家庭的必经之路。可真到了挑选的时候,很多父母才发现,家政市场鱼龙混杂,有的机构只管牵线不管服务,有的月嫂短期培训就匆匆上岗,出了问题连人都找不到。…

2026/10/10 3:20:15 阅读更多 →
个人微信API二次开发:号掉线别急着新建,先救映射表

个人微信API二次开发:号掉线别急着新建,先救映射表

线上最危险的操作不是掉线本身,而是值班同学随手「再建一个号顶上」。半小时后你会看到:一半消息失败、一半回错人、客服觉得有两个机器人在抢答。 多设备与恢复路径见 GeWe 开放文档。 真正炸掉的是什么 不是微信,是你的 客户 → 主责设备…

2026/10/10 3:20:15 阅读更多 →
Docker安装报错全解析:从daemon权限到内核模块的排查指南

Docker安装报错全解析:从daemon权限到内核模块的排查指南

你有没有遇到过这样的场景:费了好大劲把 Docker 装上,兴冲冲地敲下docker ps,结果屏幕上一行红字:permission denied while trying to connect to the Docker daemon socket at unix:///var/run/docker.sock这种感觉就像门锁装好了…

2026/10/10 3:19:15 阅读更多 →

日新闻

卫星轨道分类全解析:从LEO到GEO的选型逻辑与工程实践

卫星轨道分类全解析:从LEO到GEO的选型逻辑与工程实践

1. 从“卫星轨道分类”这个标题说起:为什么值得花时间搞懂第一次接触“卫星轨道分类”这个概念,很多人会觉得它离自己很远——不就是天上的星星怎么转吗?但如果你正在做航天任务规划、遥感数据接收、星座设计,甚至只是准备一场航天…

2026/10/10 0:00:39 阅读更多 →
Spring AOP 核心原理与实战:从概念到日志切面落地

Spring AOP 核心原理与实战:从概念到日志切面落地

1. 从一个真实痛点说起:为什么你的代码里到处都是重复逻辑刚入行那会儿,我写过一个用户管理模块,注册、登录、改密码、注销四个接口。每个接口里都塞了几乎一样的日志打印、参数校验、事务开启和提交。当时觉得没什么,能跑就行。直…

2026/10/10 0:00:40 阅读更多 →
Python招聘数据采集与分析可视化:从采集清洗到薪资技能城市可视化全链路

Python招聘数据采集与分析可视化:从采集清洗到薪资技能城市可视化全链路

简介:这是一套面向计算机相关专业学生与项目实战学习者的Python数据采集与分析可视化完整项目,以Boss直聘岗位数据为对象,适合用作毕业设计、课程设计或期末大作业。资源包共38个文件,约246KB,以13个py源码文件为核心&…

2026/10/10 0:00:40 阅读更多 →

周新闻

KT148A语音芯片外挂8002D功放的工程实践指南

KT148A语音芯片外挂8002D功放的工程实践指南

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

2026/10/8 15:26:32 阅读更多 →
LLC谐振变换器增益公式推导:从FHA等效到完整归一化表达式

LLC谐振变换器增益公式推导:从FHA等效到完整归一化表达式

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

2026/10/10 1:36:08 阅读更多 →
ARM架构深度解析:从RISC设计理念到交叉编译实战

ARM架构深度解析:从RISC设计理念到交叉编译实战

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

2026/10/9 10:11:06 阅读更多 →

月新闻

我发现了一个新思路:用 Remotion + Claude Code 像写代码一样自动化生成短视频

我发现了一个新思路:用 Remotion + Claude Code 像写代码一样自动化生成短视频

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

2026/10/8 21:13:17 阅读更多 →
Windows下 Codex 中 Chrome 和 Computer Use 插件不可用问题排查及解决参考方式:TaoToken 统一 Key 配置与验证

Windows下 Codex 中 Chrome 和 Computer Use 插件不可用问题排查及解决参考方式:TaoToken 统一 Key 配置与验证

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

2026/10/9 21:32:20 阅读更多 →
黑夜航拍船只数据集训练YOLOV5模型全流程解析

黑夜航拍船只数据集训练YOLOV5模型全流程解析

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

2026/10/9 6:17:20 阅读更多 →