Flink Python REPL(pyflink-shell)实战指南:local / remote / YARN 模式与 Table API 交互式开发

发布时间:2026/10/10 2:30:07

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),仅供参考
延伸阅读

更多相关文章

2026/10/10 2:30:07

Web 扫描结果怎么定级:迹象别写成已验证

一条告警要先分成已验证、迹象、误报或未测。写不出对应响应的句子删掉,迹象不能升成已验证漏洞。 目录 前言 一、这条告警进哪一栏 二、置信度三档 三、已经看到的响应怎么组合 四、报告六字段和未测节 五、常见误升与停手

2026/10/10 2:25:05

Leetcode 字符串的排列【中等】

这个题用集合类去实现的时候,踩了好几个坑containsAll()方法是元素维度,跟个数无关{a,b}contailsAll{a,b,b}居然返回trueremove方法,指定元素类型删除,一次只删除一个元素。比如{a,b,b}执行一次remove(new Character(b))&#xff…

2026/10/10 3:45:11

顽固木马杀不死?内核级专杀工具与常规杀软的区别及实战

正在处理一份上周的文件,电脑忽然像被什么东西按住一样卡住不动,安全软件图标打不开,任务管理器里冒出几个乱码名字的进程。最气人的是,等我用常规杀毒软件全盘扫一遍,它提示“未发现威胁”,可重启之后症状…

2026/10/10 3:45:11

恶性木马专杀实战:内核级查杀与顽固病毒清理指南

电脑感染恶性木马,和普通病毒骚扰完全是两种体验。普通木马最多是弹窗、改首页、后台偷偷占资源,真正麻烦的是那种系统被恶意驱动接管、杀毒软件打不开、进程在任务管理器里看不到、重启之后病毒又原地复活的顽固感染。这种场景下,“火绒恶性…

2026/10/10 3:45:11

文本挖掘实战指南:从非结构化数据到决策信号

每次拿到一堆客服工单、用户反馈、社交媒体的评论,我都觉得头疼。这些文本数据又多又乱,但又确确实实藏着用户最真实的声音——满意度、痛点、产品缺陷、竞品动向,全在里面。问题是,它们都是非结构化数据,没法直接塞进…

2026/10/10 3:45:11

Win11 TPM2.0不可用?华硕主板fTPM开启全指南

1. 为什么Win11升级卡在“TPM 2.0不可用”?——不是硬件不支持,而是BIOS里藏着开关你点开Windows更新,看到那行加粗的红色提示:“此电脑无法运行Windows 11。缺少必需的安全功能:可信平台模块(TPM&#xff…

2026/10/10 3:45:11

少即是多:用减法重构程序员、PM与项目经理的生活系统

上周某晚,我在公司楼下等电梯,看到群里有人抛出一个问题:程序员、产品经理、项目经理,谁最不可能准时下班?评论区瞬间吵成一团,有人自嘲“三班倒”,有人说“谁有孩子谁先走”,还有人…

2026/10/10 3:40:10

RabbitMQ消费端可靠性实战:限流、超时与死信队列全解析

凌晨两点十七分,我手机上的告警通道开始连续发声。监控面板显示,order_process队列的 Ready 消息数在十分钟内从 200 冲到了 5000,而消费者进程明明还活着,日志里却在疯狂刷同一条消息的消费失败栈——又是那台订单处理服务。做 J…

2026/10/8 10:03:18

Jev+Agent接管浏览器:browser-use实战与jev-ultrafast性能优化

1. 从“Jev”说起:为什么我要把Agent接进浏览器“Jev”这个词最近在圈子里出现的频率越来越高,很多人第一次听到会以为是某个新模型的名字,其实它更像是一种思路——把Jev模型的能力当作底座,通过Agent的方式去接管浏览器&#xf…

2026/10/9 20:15:56

多智能体集群实战:DeepAgents编排、MCP与A2A协议及Skills体系

1. 从"单兵作战"到"集群协同":多智能体编排到底在解决什么问题如果你最近在折腾 Agent 相关的东西,大概率会有一种感觉:单个 Agent 能做的事情,其实很快就摸到天花板了。你给它一个提示词,挂几个工…

2026/10/8 6:05:44

无源低通滤波器设计实战:从RC到LC,手把手教你避开那些坑

/* 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 0:04:53

从逻辑门到计算机:数字电路核心原理与全加器搭建实战

如果你拆过一台旧电脑的主板,盯着那些黑乎乎的小芯片看上一会儿,可能会冒出同一个疑问:这堆引脚密集的元件,到底是怎么“变”出那么复杂的应用的?答案并不在某个神秘的部件里,而是在所有芯片内部都在反复使…

还想了解更多?直接咨询顾问

免费诊断 + 免费方案 + 透明报价。

全国咨询热线400-8866-253
免费获取方案
☎咨询二维码 ☎ ↑