ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Flink Python REPL 完全指南:用 PyFlink Shell 交互式开发 Table API 作业

Flink Python REPL 完全指南:用 PyFlink Shell 交互式开发 Table API 作业 大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载Flink 自带一个集成的交互式 Python ShellREPL它既能运行在本地启动的 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 Shelllocal 模式 $ 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 EnvironmentPython Shell 启动后会自动加载 Table API 相关的全部导入与 Table Environment无需手动创建st_envStreamTableEnvironment用于流处理 Table 程序bt_envBatchTableEnvironment用于批处理 Table 程序s_envStreamExecutionEnvironment流处理底层执行环境可通过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_intost_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 / --jobManagerMemoryJobManager 容器内存可带单位如-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项目地址https://gitcode.com/gh_mirrors/fli/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),仅供参考
返回列表