- 大数据
- 流处理
- 批处理
- 数据工程
【免费下载链接】flink
Flink 自带一个集成的交互式 Python Shell(REPL),它既能运行在本地启动的 local 模式,也能运行在集群启动的 cluster 模式(remote / YARN),让开发者可以像使用 Jupyter 一样逐行验证 PyFlink Table API 逻辑。本文以 python_shell.md 为主体,结合仓库中的 pyflink-shell.sh、PythonShellParser.java 与 shell.py 源码,完整讲解安装、四种启动方式、预置环境变量、Table API 交互示例及全部命令行参数,读完即可上手用 Python Shell 做流批作业的原型验证。
环境准备与安装
Python Shell 本质上是一个包装了 PyFlink 的交互式 Python 进程,因此使用前必须保证本机的 Python 与 PyFlink 环境就绪。
Python 版本要求
PyFlink 需要 Python 3.7 以上版本(3.8、3.9 或 3.10),建议先确认版本:
$ python --version如果你的系统安装了多个 Python 版本,可以通过软链接将python指向python3,或者创建 Python 虚拟环境(venv)来隔离依赖。更详细的环境安装说明见 Python Table API 环境安装。另外,Shell 启动时会调用python命令,可用环境变量PYFLINK_PYTHON覆盖默认解释器。
安装 PyFlink
推荐通过 PyPi 安装 PyFlink,然后即可使用 Python Shell:
# 安装 PyFlink $ python -m pip install apache-flink # 启动 Python Shell(local 模式) $ pyflink-shell.sh local若需要与当前仓库版本严格对应,可以安装指定版本:python -m pip install apache-flink==<版本号>。也可以从源码构建 Flink 后,使用flink-python/bin目录下的 pyflink-shell.sh 脚本启动。
启动脚本的工作原理
从源码层面理解 Python Shell 有助于排查问题。pyflink-shell.sh 的启动流程如下:
- 通过
find-flink-home.sh与config.sh确定FLINK_HOME并构造 Flink 的 classpath,同时定位$FLINK_OPT_DIR/flink-python*.jar; - 将
$FLINK_OPT_DIR/python/下的pyflink.zip、py4j-*-src.zip、cloudpickle-*-src.zip加入PYTHONPATH,保证 Python 侧能导入 PyFlink 并通过 Py4J 与 JVM 通信; - 调用 JVM 端的参数解析器 PythonShellParser.java,把用户输入的 shell 参数翻译成 Flink 客户端可识别的提交参数(以
\0分隔输出,由 bash 解析回OPTIONS数组); - 最后以交互模式执行 Python 模块:
${PYFLINK_PYTHON} -i -m pyflink.shell,即进入 shell.py 定义好的 REPL 环境。
也就是说,Python Shell = JVM 参数翻译层 + Py4J 桥接 + 预置环境的交互式 Python 解释器。-i保证执行完初始化代码后停留在交互提示符,-m pyflink.shell负责加载预置环境与欢迎信息。
使用:预置的 Table Environment
Python Shell 启动后会自动加载 Table API 相关的全部导入与 Table Environment,无需手动创建:
st_env:StreamTableEnvironment,用于流处理 Table 程序;bt_env:BatchTableEnvironment,用于批处理 Table 程序;s_env:StreamExecutionEnvironment,流处理底层执行环境,可通过s_env.set_parallelism(n)设置并行度。
从 shell.py 源码可以看到,流式环境在模块加载时即完成初始化:s_env = StreamExecutionEnvironment.get_execution_environment(),st_env = StreamTableEnvironment.create(s_env)。同时模块顶部已from pyflink.common import *、from pyflink.table import *等批量导入,因此DataTypes、col、TableDescriptor、Schema、FormatDescriptor等符号开箱即用,欢迎信息中也会明确提示 "Use the prebound Table Environment to implement batch or streaming Table programs."
流处理 Table API 示例
下面是在 Python Shell 中逐行输入的流式示例:构造两张行的临时表,做一次a + 1的投影后写入文件系统 sink,并在 local 模式下读取结果文件验证输出。
>>> 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", sink_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())要点说明:
from_elements以 Python 列表 + 字段名列表快速构造内存表,非常适合 REPL 中的快速验证;TableDescriptor.for_connector("filesystem")声明文件系统连接器,配合FormatDescriptor.for_format("csv")声明 CSV 格式,option("path", ...)指定输出路径;execute_insert(...).wait()是阻塞式提交,确保作业执行完成后再读取结果文件;- 这里的流式写法与 shell.py 中的内置示例同构,后者使用
insert_into+st_env.execute("stream_job")亦可。
批处理 Table API 示例
批处理场景使用bt_env,逻辑与流式几乎一致,区别仅在于环境变量与 sink 表名:
>>> 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']) >>> bt_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", sink_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())启动方式详解
查看 Python Shell 提供的全部可选参数,先执行:
pyflink-shell.sh --help从 PythonShellParser.java 可以看到,脚本支持三种集群类型子命令:local、remote、yarn,以及顶层-h | --help。未指定集群类型时会直接报错退出。
Local 模式
Local 模式下,Python Shell 会在 JVM 内启动一个本地 Flink 集群(mini cluster)来执行作业,适合日常原型验证与教学:
pyflink-shell.sh local对应源码中 parseLocal 的实现:local 模式不附加任何额外提交参数,仅保留local关键字交给flink run使用。
Remote 模式
若已有独立部署的 Flink 集群(例如 Standalone 集群,部署方式见 本地安装),可以通过remote关键字指定 JobManager 的主机名与端口号:
pyflink-shell.sh remote <hostname> <portnumber>例如pyflink-shell.sh remote 10.0.0.1 8081。源码 parseRemote 会将remote <host> <port>翻译为-m <host>:<port>传给flink run,即提交到指定 JobManager;若未提供 host/port 会打印错误并提示用法。
YARN 集群模式(新建集群)
Python Shell 也可以运行在 YARN 之上:它会在 YARN 上部署一个新的 Flink 集群并自动连接,除指定 container 数量外,还可以指定 JobManager 内存、YARN 应用名、队列、slot 数等参数。例如在一个部署了两个 TaskManager 的 YARN 集群上运行:
pyflink-shell.sh yarn -n 2所有可选的 YARN 参数见下文"完整的参考"。源码 getYarnOptions 定义了-jm、-nm、-qu、-s、-tm五个专属选项,并在 parseYarn 中统一加上-m yarn-cluster目标,同时为每个选项添加y前缀(如-yn、-yjm)以对齐flink run的 YARN 提交参数。
YARN Session 模式(连接已有集群)
如果已经通过 Flink YARN Session 部署好一个 Flink 集群,则可以不带任何参数直接连接该集群:
pyflink-shell.sh yarn完整的参考:全部命令与参数
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 打印输出使用文档参数补充说明(依据 PythonShellParser.java 中的选项定义):
-jm / --jobManagerMemory:JobManager 容器内存,可带单位(如-jm 1024m),默认单位 MB;-n / --container:分配的 YARN container 数量,即 TaskManager 数量;-nm / --name:自定义 YARN Application 名称,便于在 ResourceManager 上区分任务;-qu / --queue:指定提交到的 YARN 队列;-s / --slots:每个 TaskManager 上的 slot 数量,影响单容器可运行的并行任务数;-tm / --taskManagerMemory:每个 TaskManager 容器内存,同样可带单位,默认单位 MB。
小结与进阶路径
Flink Python REPL 的价值在于"零工程化"地验证 Table API 逻辑:预置的s_env/st_env/bt_env免去了每次编写环境初始化样板代码,local / remote / yarn 三种集群目标让原型可以无缝从单机验证过渡到集群提交。仓库中还提供了对应的自动化验证用例 test_shell_example.py,它直接在测试中from pyflink.shell import s_env, st_env, DataTypes复用预置环境并跑通同样的文件系统 sink 流程,可作为理解 REPL 内部行为的参考。
如果你需要在 Shell 之外编写完整的 PyFlink 作业,可以继续阅读 Python Table API 环境安装;若想了解从源码构建 Flink 后如何获得pyflink-shell.sh,参见 从源码构建 Flink。
- 大数据
- 流处理
- 批处理
- 数据工程
【免费下载链接】flink
相关推荐
SharpXDecrypt:自动化Xshell全版本密码恢复技术方案
SharpXDecrypt:自动化Xshell全版本密码恢复技术方案 SharpXDecrypt是一款专业级Xshell密码恢复工具,专为技术运维和安全审计人员
大数据流处理批处理数据工程Flink PyFlink 配置指南:Python DataStream / Table API 的配置项设置与调优详解
Flink PyFlink 配置指南:Python DataStream / Table API 的配置项设置与调优详解 本文以 Flink 仓库中 PyFli
大数据流处理批处理数据工程从6K到12K:戴森球计划翘曲器蓝图选择完全指南
从6K到12K:戴森球计划翘曲器蓝图选择完全指南 想象一下,当你终于建好了星际物流网络,却发现翘曲器库存告急,舰队停滞在星海之间——这种星际交通的"燃料危机"是
游戏开发
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考