☰
Apache Zeppelin Flink 解释器实战指南:从本地开发到 Yarn 集群的批流一体分析
2026/10/7 9:38:30 网站建设 项目流程

本文以 Zeppelin 仓库中的 Flink 解释器官方文档 为主体,系统讲解 Zeppelin 中 Flink 解释器组的五个解释器(%flink、%flink.pyflink、%flink.ipyflink、%flink.ssql、%flink.bsql)的安装配置、四种执行模式、Scala/Python/SQL 多语言编程、SQL 增强特性、流式可视化、UDF 注册与 Hive 集成等完整实战方案。读完本文,你将掌握如何在 Zeppelin 中以交互式笔记本的方式运行 Flink 批流任务,并理解其底层 Scala Shell / Python Shell 双入口的架构原理。

概述:Flink 解释器组

Apache Flink 是一个面向无界与有界数据流的有状态计算框架与分布式处理引擎,设计上可在各类常见集群环境中运行,并以内存级速度和任意规模执行计算。Zeppelin 在 0.9 版本对 Flink 解释器进行了重构以支持最新版 Flink。需要注意:当前仅支持 Flink 1.20 及以上版本,旧版本 Flink 无法正常工作。

Flink 在 Zeppelin 中由一个解释器组(Flink interpreter group)支持,共包含五个解释器,全部位于同一 Flink session 中、共享同一个 Flink 集群环境:

名称类说明
%flinkFlinkInterpreter创建 ExecutionEnvironment / StreamExecutionEnvironment / BatchTableEnvironment / StreamTableEnvironment,并提供 Scala 环境
%flink.pyflinkPyFlinkInterpreter提供 Python 环境
%flink.ipyflinkIPyFlinkInterpreter提供 IPython 环境
%flink.ssqlFlinkStreamSqlInterpreter提供流式 SQL 环境
%flink.bsqlFlinkBatchSqlInterpreter提供批式 SQL 环境

从源码结构看,%flink是整个解释器组的入口:FlinkInterpreter(见 FlinkInterpreter.java)会根据运行时 Scala 版本(2.11/2.12)动态加载对应的FlinkScalaInterpreter实现类,内部实际委托给 Flink Scala Shell 执行。%flink.ssql与%flink.bsql继承自公共基类 FlinkSqlInterpreter.java,在同一 session 内复用FlinkInterpreter的运行时环境与 ZeppelinContext。

核心特性

特性说明
支持多版本 Flink可以在一个 Zeppelin 实例中运行不同版本的 Flink
支持多语言支持 Scala、Python、SQL,并且可以跨语言协作,例如编写 Scala UDF 后在 PyFlink 中使用
支持多种执行模式Local / Remote / Yarn / Yarn Application
支持 Hive支持 Hive catalog
交互式开发交互式开发体验提升生产力
Flink SQL 增强在一个笔记本中同时支持流式 SQL 与批式 SQL;支持单行注释(--)与多行注释(/* */);支持高级配置(jobName、parallelism);支持多条 insert 语句
多租户多个用户可在一个 Zeppelin 实例中工作而互不影响
Rest API 支持不仅可以通过 Zeppelin 笔记本 UI 提交 Flink 作业,还可以通过其 Rest API 提交(可将 Zeppelin 用作 Flink 作业服务器)

快速体验:在 Zeppelin Docker 中运行 Flink

对于初学者,推荐在 Zeppelin Docker 中体验 Flink。Zeppelin 发行版不内置 Flink 二进制包,因此需要先自行下载 Flink。例如将 Flink 1.12.2 下载到/mnt/disk1/flink-1.12.2,然后挂载进 Zeppelin 容器并启动:

docker run -u $(id -u) -p 8080:8080 -p 8081:8081 --rm -v /mnt/disk1/flink-1.12.2:/opt/flink -e FLINK_HOME=/opt/flink --name zeppelin apache/zeppelin:0.10.0

启动后打开http://localhost:8080即可在 Zeppelin 中体验 Flink。文档说明在 Docker 中只验证了 Flink local 模式,其他模式可能受网络问题影响。参数说明:

  • -p 8080:8080:暴露 Zeppelin Web UI;
  • -p 8081:8081:暴露 Flink Web UI,可通过http://localhost:8081访问;
  • -v /mnt/disk1/flink-1.12.2:/opt/flink:将本地 Flink 目录挂载为容器内/opt/flink,配合-e FLINK_HOME=/opt/flink供解释器定位 Flink 安装目录;
  • -u $(id -u):以当前用户运行,避免挂载目录权限问题。

你也可以挂载自己的笔记本目录以替换内置教程笔记本,例如克隆 flink-sql-cookbook 仓库后:

docker run -u $(id -u) -p 8080:8080 --rm -v /mnt/disk1/flink-sql-cookbook-on-zeppelin:/notebook -v /mnt/disk1/flink-1.12.2:/opt/flink -e FLINK_HOME=/opt/flink -e ZEPPELIN_NOTEBOOK_DIR='/notebook' --name zeppelin apache/zeppelin:0.10.0

其中ZEPPELIN_NOTEBOOK_DIR指定 Zeppelin 从挂载目录读取笔记本。

环境准备:下载与配置 Flink 发行版

下载 Flink 1.19 或 1.20。由于 Zeppelin 的 Flink 解释器依赖 Scala Table API bridge 等额外 jar,需要对 Flink 发行版做如下调整(将${FLINK_VERSION}替换为实际安装的 Flink 版本):

  • 将${FLINK_HOME}/opt/flink-table-planner_2.12-${FLINK_VERSION}.jar移动到${FLINK_HOME}/lib;
  • 将${FLINK_HOME}/lib/flink-table-planner-loader-${FLINK_VERSION}.jar移动到${FLINK_HOME}/opt(与上一步互换位置);
  • 下载flink-table-api-scala-bridge_2.12-${FLINK_VERSION}.jar与flink-table-api-scala_2.12-${FLINK_VERSION}.jar放入${FLINK_HOME}/lib;
  • 将${FLINK_HOME}/opt/flink-sql-client-${FLINK_VERSION}.jar移动到${FLINK_HOME}/lib。

这样做的原因是:Zeppelin 的 Flink 解释器内部以 Scala Shell 方式编程式地调用 Table API,而不是启动 SQL Client 进程,因此需要把 planner 与 Scala bridge 直接置于 Flink 的 classpath(lib 目录)中。

Flink on Zeppelin 架构

Flink on Zeppelin 的整体架构如下图所示:左侧的 Flink 解释器实际上是一个 Flink 客户端,负责编译并管理 Flink 作业的生命周期(提交、取消作业、监控作业进度等);右侧的 Flink 集群负责实际执行 Flink 作业。集群形态可以是 MiniCluster(本地模式)、Standalone 集群(远程模式)、Yarn session 集群(yarn 模式)或 Yarn application session 集群(yarn-application 模式)。

Flink 解释器内部有两大关键组件:Scala Shell与Python Shell。

  • Scala Shell:Flink 解释器的入口,负责创建 Flink 程序的全部入口对象,如 ExecutionEnvironment、StreamExecutionEnvironment 和 TableEnvironment,并负责编译运行 Scala 代码和 SQL;
  • Python Shell:PyFlink 的入口,负责编译运行 Python 代码。

从源码看,FlinkScalaInterpreter(见 FlinkScalaInterpreter.scala)在open()时依次完成:初始化 Flink 配置 → 创建 FlinkILoop(Scala REPL)→ 创建批/流 TableEnvironment → 绑定 ZeppelinContext(z)→ 注册 JobListener(用于把作业与段落关联、上报进度)→ 可选注册 Hive catalog → 自动加载 UDF jar。

配置详解

Flink 解释器通过 Zeppelin 提供的属性进行配置(下表),同时你也可以添加表外其他的 Flink 属性(以table.exec、parallelism等 Flink 官方配置项为例,参考 Flink 官方文档的 Available Properties 列表)。从源码看,解释器启动时会把解释器属性中所有条目写入 Flink 的Configuration(properties.asScala.foreach(entry => configuration.setString(...))),因此任何 Flink 配置项都能以解释器属性方式下发。

环境类属性

属性默认值说明
FLINK_HOME(必填)Flink 安装位置。必须指定,否则无法在 Zeppelin 中使用 Flink
HADOOP_CONF_DIR(必填,yarn 模式)Hadoop 配置目录位置,yarn 模式下必须设置
HIVE_CONF_DIR(可选)Hive 配置目录位置,需要连接 Hive metastore 时必须设置

执行模式与集群资源类属性

属性默认值说明
flink.execution.modelocalFlink 执行模式:local/remote/yarn/yarn-application
flink.execution.remote.host无运行中 JobManager 的主机名,仅 remote 模式使用
flink.execution.remote.port无运行中 JobManager 的端口,仅 remote 模式使用
jobmanager.memory.process.size1024mJobManager 总内存大小(官方 Flink 属性)
taskmanager.memory.process.size1024mTaskManager 总内存大小(官方 Flink 属性)
taskmanager.numberOfTaskSlots1每个 TaskManager 的 slot 数
local.number-taskmanager4本地模式下 TaskManager 总数
yarn.application.nameZeppelin Flink SessionYarn 应用名称
yarn.application.queuedefaultYarn 应用的队列名

源码实现细节(见 FlinkScalaInterpreter.scala):

  • 内存类属性同时兼容旧名flink.jm.memory/flink.tm.memory,新名优先;
  • taskmanager.numberOfTaskSlots也兼容旧名flink.tm.slot;
  • 应用名兼容旧名flink.yarn.appName,队列兼容旧名flink.yarn.queue;
  • 在 yarn-application 模式下,源码会以 Yarn 容器当前工作目录作为FLINK_HOME、FLINK_CONF_DIR与HIVE_CONF_DIR(isYarnApplicationMode分支),因此无需再手动指定这些目录。

UI 与安全类属性

属性默认值说明
zeppelin.flink.uiWebUrl无用户指定的 Flink JobManager URL。可用于 remote 模式下已启动的集群,也可作为 URL 模板,例如https://knox-server:8443/gateway/cluster-topo/yarn/proxy/{{applicationId}}/,其中{{applicationId}}是 Yarn 应用 ID 的占位符
zeppelin.flink.run.asLoginUsertrue是否以 Zeppelin 登录用户运行 Flink 作业,仅当在 Hadoop Yarn 集群上运行且启用 shiro 时生效

{{applicationId}}占位符替换逻辑在源码getDisplayedJMWebUrl中实现:若设置了zeppelin.flink.uiWebUrl,则将其中的{{applicationId}}替换为真实 Yarn 应用 ID(见 FlinkScalaInterpreter.scala)。

依赖与 UDF 类属性

属性默认值说明
flink.udf.jars无Flink UDF jar(逗号分隔)。Zeppelin 会自动为用户注册这些 jar 中的 UDF。jar 可以是本地文件,若安装了 Hadoop 也可以是 HDFS 文件。UDF 名称即类名
flink.udf.jars.packages无需要扫描的包(逗号分隔),限定flink.udf.jars中 UDF 的搜索范围。指定后可减少扫描类数,否则将扫描 jar 中全部类
flink.execution.jars无附加用户 jar(逗号分隔),可以是本地或 HDFS 文件。用于指定 Flink connector jar 或 UDF jar(但不具备flink.udf.jars那样的 UDF 自动注册功能)
flink.execution.packages无附加用户 Maven 包(逗号分隔),例如org.apache.flink:flink-json:1.10.0

SQL 并发与 Python 类属性

属性默认值说明
zeppelin.flink.concurrentBatchSql.max10%flink.bsql批式 SQL 的最大并发数
zeppelin.flink.concurrentStreamSql.max10%flink.ssql流式 SQL 的最大并发数
zeppelin.pyflink.pythonpythonPyFlink 使用的 Python 可执行文件
table.exec.resource.default-parallelism1Flink SQL 作业的默认并行度

显示、Hive 与作业生命周期类属性

属性默认值说明
zeppelin.flink.scala.colortrue是否彩色显示 Scala Shell 输出
zeppelin.flink.scala.shell.tmp_dir无存放 Scala Shell 编译 jar 的临时目录
zeppelin.flink.enableHivefalse是否启用 Hive
zeppelin.flink.hive.version2.3.7要连接的 Hive 版本
zeppelin.flink.module.enableHivefalse是否启用 Hive module;若启用,Hive UDF 优先于 Flink UDF
zeppelin.flink.maxResult1000SQL 解释器返回的最大行数
zeppelin.flink.job.check_interval1000检查 Flink 作业进度的间隔(毫秒)
flink.interpreter.close.shutdown_clustertrue关闭解释器时是否关闭 Flink 集群
zeppelin.interpreter.close.cancel_jobtrue关闭解释器时是否取消 Flink 作业

源码佐证:zeppelin.flink.maxResult直接用于构造FlinkZeppelinContext(z.show输出的最大行数);flink.interpreter.close.shutdown_cluster决定close()时是否调用clusterClient.shutDownCluster(),同时 yarn 模式下会删除 Flink staging 目录(见 FlinkScalaInterpreter.scala)。

解释器绑定模式

默认的解释器绑定模式为globally shared(全局共享),意味着所有笔记本共享同一个 Flink 解释器,也就共享同一个 Flink 集群。实践中更推荐使用isolated per note(按笔记本隔离):每个笔记本拥有独立的 Flink 解释器和各自的 Flink 集群,互不影响。

四种执行模式

Flink in Zeppelin 支持四种执行模式(通过flink.execution.mode设置):Local、Remote、Yarn、Yarn Application。

Local 模式

本地模式会在本地 JVM 中启动一个 MiniCluster。默认本地 MiniCluster 使用 8081 端口,请确保该端口可用,否则可通过rest.port指定其他端口。可通过local.number-taskmanager与flink.tm.slot自定义 TaskManager 数量与每 TM 的 slot 数——默认只有 4 个 TM、每 TM 1 个 slot,某些场景下可能不够用。

Remote 模式

Remote 模式会连接一个已存在的 Flink 集群(Standalone 集群或 Yarn session 集群)。除将flink.execution.mode设为remote外,还需要设置flink.execution.remote.host与flink.execution.remote.port指向 Flink JobManager 的 Rest API 地址。源码中会校验这两个参数:若未指定 host 或 port,直接抛出InterpreterException(见 FlinkScalaInterpreter.scala)。

Yarn 模式

在 Yarn 模式运行 Flink 需要满足以下设置:

  • 将flink.execution.mode设为yarn;
  • 在 Flink 解释器设置或zeppelin-env.sh中设置HADOOP_CONF_DIR;
  • 确保hadoop命令在PATH中。因为内部 Flink 会调用命令hadoop classpath并将所有 Hadoop 相关 jar 加载进 Flink 解释器进程。

此模式下,Zeppelin 会为你启动一个 Flink Yarn session 集群,并在关闭 Flink 解释器时销毁它。源码中 yarn 模式会额外校验FlinkYarnSessionCli类是否可加载:找不到时抛出 "No hadoop jar found, make sure you have hadoop command in your PATH" 的明确错误。

Yarn Application 模式

上述 yarn 模式在 Zeppelin 服务器主机上存在一个独立的 Flink 解释器进程;当解释器进程过多时可能耗尽资源。因此实践上推荐:若使用 Flink 1.11 或更高版本(yarn application 模式仅在 Flink 1.11 之后支持),使用 yarn application 模式。该模式下 Flink 解释器运行在 Yarn 容器中的 JobManager 内。

运行条件与 yarn 模式类似:

  • 将flink.execution.mode设为yarn-application;
  • 在 Flink 解释器设置或zeppelin-env.sh中设置HADOOP_CONF_DIR;
  • 确保hadoop命令在PATH中。

源码实现上,yarn-application 模式通过环境变量_APP_ID获取 Yarn 应用 ID,并从本地localhost的 Rest 端口连接 JobManager(见 FlinkScalaInterpreter.scala)。

Flink Scala 开发

Scala 是 Zeppelin 上 Flink 的默认语言(%flink),也是 Flink 解释器的入口。解释器底层创建 Scala Shell,并预创建若干内置变量,包括 ExecutionEnvironment、StreamExecutionEnvironment 等——不要重复创建这些 Flink 环境变量,否则可能遇到诡异的问题。你在 Zeppelin 中编写的 Scala 代码会提交到这个 Scala Shell 执行。

Flink Scala Shell 中创建的内置变量:

变量含义
senvStreamExecutionEnvironment(流执行环境)
benvExecutionEnvironment(批执行环境)
stenvStreamTableEnvironment(blink planner,即新 planner)
btenvBatchTableEnvironment(blink planner,即新 planner)
zZeppelinContext

源码中这些变量通过flinkILoop.intp.bind(...)绑定到 REPL 命名空间,同时预导入大量 Flink API 包(org.apache.flink.api.scala._、org.apache.flink.streaming.api.scala._、org.apache.flink.table.api._、ScalarFunction / AggregateFunction / TableFunction / TableAggregateFunction 等),确保用户代码可直接使用(见 FlinkScalaInterpreter.scala)。

Blink/Flink Planner

Zeppelin 0.11 之后移除了对 flink planner(旧 planner)的支持——Flink 1.14 之后也移除了旧 planner。因此当前仓库仅使用 blink planner(新 planner)。

流式 WordCount 示例

在 Zeppelin 中可以直接编写任意 Scala 代码,例如经典的流式 WordCount 示例:

%flink // 使用 senv(StreamExecutionEnvironment)编写流式作业

代码补全

在 Zeppelin 中按 Tab 键即可触发代码补全。从源码看,补全能力来自FlinkILoop的Completion组件(scalaCompletion)。

ZeppelinContext

ZeppelinContext提供了一些附加函数与工具,详细内容可参考 Zeppelin-Context 文档。在 Flink 解释器中,可以用z展示 Flink 的 Dataset/Table,例如:

  • z.show(DataSet):展示批式 DataSet;
  • z.show(Batch Table):展示批式 Table;
  • z.show(Stream Table):展示流式 Table。

Flink SQL

Zeppelin 提供两类 Flink SQL 解释器:

  • %flink.ssql:流式 SQL 解释器,通过StreamTableEnvironment启动 Flink 流式作业;
  • %flink.bsql:批式 SQL 解释器,通过BatchTableEnvironment启动 Flink 批式作业。

Zeppelin 的 Flink SQL 解释器等同于 Flink SQL Client,并增加了许多增强特性。

SQL 增强特性

批式 SQL 与流式 SQL 并存

在 Flink SQL Client 中,一个 session 要么运行流式 SQL,要么运行批式 SQL,无法同时进行。但在 Zeppelin 中两者可以共存:%flink.ssql运行流式 SQL,%flink.bsql运行批式 SQL,且批/流 Flink 作业运行在同一个 Flink session 集群中。

支持多语句

一个段落内可编写多条 SQL 语句,每条语句以分号(;)分隔。

支持 SQL 注释

Zeppelin 支持两种 SQL 注释:

  • 单行注释:以--开头;
  • 多行注释:以/* */包裹。

作业并行度设置

通过段落本地属性parallelism设置 SQL 并行度。源码setParallelismIfNecessary会读取段落 local properties 中的parallelism,同时更新senv、benv的并行度以及 TableEnvironment 的table.exec.resource.default-parallelism配置;maxParallelism也会被设置到流环境(见 FlinkScalaInterpreter.scala)。

支持多条 insert

有时你有多条 insert 语句读取同一数据源、写入不同 sink。默认情况下每条 insert 语句启动一个独立的 Flink 作业;将段落本地属性runAsOne设为true可以让它们在一个 Flink 作业中运行。

设置作业名

通过段落本地属性jobName为 insert 语句设置 Flink 作业名。注意:只能为 insert 语句设置作业名,select 语句暂不支持;且该设置仅对单条 insert 语句生效,对上述多条 insert 合并场景不生效。

流式数据可视化

Zeppelin 可以对 Flink 流式作业的 select SQL 结果进行可视化,共支持 3 种模式:Single、Update、Append。这三种模式分别对应源码中的SingleRowStreamSqlJob、UpdateStreamSqlJob、AppendStreamSqlJob(见 FlinkStreamSqlInterpreter.java,通过段落本地属性type选择)。

Single 模式

Single 模式适用于 SQL 语句结果始终只有一行的场景。输出格式为 HTML,可通过段落本地属性template指定最终输出内容模板,使用{i}作为结果第 i 列的占位符。

Update 模式

Update 模式适用于输出多行且持续更新的场景,例如使用 group by 的查询。

Append 模式

Append 模式适用于输出数据持续追加的场景,例如使用 tumble window 的查询。

PyFlink

PyFlink 是 Flink on Zeppelin 的 Python 入口。内部 Flink 解释器会创建 Python Shell,并创建 Flink 的环境变量(ExecutionEnvironment、StreamExecutionEnvironment 等)。需要注意:PyFlink 背后的 Java 环境是在 Scala Shell 中创建的,即 Scala Shell 与 Python Shell 共享同一个底层环境。

Python Shell 中创建的变量:

变量含义
s_envStreamExecutionEnvironment
b_envExecutionEnvironment
st_envStreamTableEnvironment(blink planner,即新 planner)
bt_envBatchTableEnvironment(blink planner,即新 planner)

配置 PyFlink

要让 PyFlink 在 Zeppelin 中工作,需要配置三件事:

  1. 安装 pyflink:例如pip install apache-flink==1.11.1。若需要使用 PyFlink UDF,则必须在所有 TaskManager 节点上安装 pyflink——也就是说如果使用 yarn,所有 yarn 节点都需要安装 pyflink;
  2. 复制 Python 文件夹:将${FLINK_HOME}/opt下的python文件夹复制到${FLINK_HOME}/lib;
  3. 设置zeppelin.pyflink.python:默认使用PATH中的 python。若安装了多个 Python 版本,需要将zeppelin.pyflink.python配置为要使用的 Python 版本。

源码佐证:PyFlinkInterpreter.open()会将zeppelin.pyflink.python映射为 Python 解释器的zeppelin.python属性,并将zeppelin.pyflink.useIPython映射为zeppelin.python.useIPython,随后启动 Python 进程与 JVM gateway(见 PyFlinkInterpreter.java)。

使用 PyFlink 的两种方式

  • %flink.pyflink:简单易用,除上述设置外无需其他操作,但功能也有限;
  • %flink.ipyflink:提供与 Jupyter 几乎一致的用户体验,官方建议使用此方式。

配置 IPyFlink

如果没有安装 anaconda,需要安装以下 3 个库:

pip install jupyter pip install grpcio pip install protobuf

如果已安装 anaconda,只需安装以下 2 个库:

pip install grpcio pip install protobuf

ZeppelinContext在 PyFlink 中同样可用,使用方式与 Flink Scala 几乎相同。IPython 的更多特性可参考 Python 解释器文档。

第三方依赖管理

无论使用 Scala、Python 还是 SQL 编写 Flink 作业,都很常见需要第三方依赖。在 IDE 中很容易添加依赖(例如在 pom.xml 中),但在 Zeppelin 中主要通过两个设置添加第三方依赖:

  • flink.execution.packages
  • flink.execution.jars

flink.execution.packages

这是推荐的依赖添加方式,其实现与在pom.xml中添加依赖相同:底层会从 Maven 仓库下载所有包及其传递依赖,然后放入 classpath。以下示例通过内联配置添加 Flink 1.10 的 Kafka connector:

%flink.conf flink.execution.packages org.apache.flink:flink-connector-kafka_2.11:1.10.0,org.apache.flink:flink-connector-kafka-base_2.11:1.10.0,org.apache.flink:flink-json:1.10.0

格式为artifactGroup:artifactId:version,多个包用逗号分隔。flink.execution.packages需要能访问互联网;如果无法访问互联网,则需要改用flink.execution.jars。

源码中该功能由DependencyResolver实现,将坐标解析下载到本地 Maven 仓库目录(默认~/.m2/repository),并支持通过zeppelin.proxy.url等属性配置代理(见 FlinkScalaInterpreter.scala)。

flink.execution.jars

如果 Zeppelin 机器无法访问互联网,或依赖未部署到 Maven 仓库,则用flink.execution.jars指定所依赖的 jar 文件(每个 jar 用逗号分隔)。例如添加 kafka 依赖(含 kafka connector 及其传递依赖):

%flink.conf flink.execution.jars /usr/lib/flink-kafka/target/flink-kafka-1.0-SNAPSHOT.jar

源码中 jar 路径支持本地文件,也支持包含://的 URL(通过HadoopUtils.downloadJar下载),本地文件不存在时会明确报错jar file: ${jar} doesn't exist。

Flink UDF

Zeppelin 中定义 UDF 有 4 种方式:

  1. 编写 Scala UDF
  2. 编写 PyFlink UDF
  3. 通过 SQL 创建 UDF
  4. 通过flink.udf.jars配置 UDF jar

Scala UDF

%flink class ScalaUpper extends ScalarFunction { def eval(str: String) = str.toUpperCase } btenv.registerFunction("scala_upper", new ScalaUpper())

定义 Scala UDF 的方式与在 IDE 中几乎相同。创建 UDF 类后,通过btenv注册;也可以通过stenv注册,它与btenv共享同一个 Catalog。

Python UDF

%flink.pyflink class PythonUpper(ScalarFunction): def eval(self, s): return s.upper() bt_env.register_function("python_upper", udf(PythonUpper(), DataTypes.STRING(), DataTypes.STRING()))

Python UDF 的定义同样与 IDE 中几乎相同。创建 UDF 类后,通过bt_env注册;也可以通过st_env注册,它与bt_env共享同一个 Catalog。

通过 SQL 创建 UDF

一些简单 UDF 可以直接在 Zeppelin 中编写,但如果 UDF 逻辑非常复杂,最好在 IDE 中编写,然后在 Zeppelin 中按如下方式注册:

%flink.ssql CREATE FUNCTION myupper AS 'org.apache.zeppelin.flink.udf.JavaUpper';

这种方式要求 UDF jar 必须在CLASSPATH上,因此需要配置flink.execution.jars将 UDF jar 加入 classpath:

%flink.conf flink.execution.jars /usr/lib/flink-udf-1.0-SNAPSHOT.jar

flink.udf.jars

上述 3 种方式都有局限:

  • 在 Zeppelin 中适合编写简单 Scala UDF 或 Python UDF,但不适合编写非常复杂的 UDF——因为笔记本相比 IDE 缺少高级特性,如包管理、代码导航等;
  • UDF 难以在笔记本或用户间共享,每次都要在每个 Flink 解释器中运行定义 UDF 的段落。

因此当 UDF 数量很多、或 UDF 逻辑很复杂、且不想每次都手动注册时,可以使用flink.udf.jars:

  • 步骤 1:在 IDE 中创建 UDF 项目并编写 UDF;
  • 步骤 2:将flink.udf.jars指向从 UDF 项目构建出的 jar。

例如:

%flink.conf flink.execution.jars /usr/lib/flink-udf-1.0-SNAPSHOT.jar

Zeppelin 会扫描该 jar,找出所有 UDF 类并自动注册,UDF 名称即类名。默认情况下 Zeppelin 会扫描 jar 中所有类,若 jar 很大(尤其是 UDF jar 还包含其他依赖时)扫描会相当慢,此时建议指定flink.udf.jars.packages限定扫描包范围,可显著减少扫描类数、加快 UDF 识别。

源码中loadUDFJar的实现印证了这一机制:遍历 jar 中所有.class文件(跳过含$的内部类),实例化后按ScalarFunction、TableFunction、AggregateFunction、TableAggregateFunction四种类型分别注册到btenv,并以类的简单名作为函数名;当设置了flink.udf.jars.packages时,只扫描匹配前缀的类(见 FlinkScalaInterpreter.scala)。

如何集成 Hive

在 Flink 中使用 Hive 需要做以下设置:

  • 将zeppelin.flink.enableHive设为true;
  • 将zeppelin.flink.hive.version设为你使用的 Hive 版本;
  • 将HIVE_CONF_DIR设为hive-site.xml所在位置。确保 Hive metastore 已启动,并在hive-site.xml中配置了hive.metastore.uris;
  • 将以下依赖复制到 Flink 安装目录的 lib 文件夹:
    • flink-connector-hive_2.11–*.jar
    • flink-hadoop-compatibility_2.11–*.jar
    • hive-exec-2.x.jar(对于 hive 1.x,需要复制hive-exec-1.x.jar、hive-metastore-1.x.jar、libfb303–0.9.2.jar和libthrift-0.9.2.jar)

源码佐证:registerHiveCatalog()会创建名为hive的HiveCatalog并设为默认 catalog 与 database(默认default),当zeppelin.flink.module.enableHive为 true 时还会加载HiveModule以支持 Hive 内置函数;未指定HIVE_CONF_DIR时会抛出明确异常(见 FlinkScalaInterpreter.scala)。

Paragraph 本地属性

在"流式数据可视化"一节中,我们通过段落本地属性type演示了不同可视化类型。本节完整列出 Flink 解释器支持的所有段落本地属性:

属性默认值说明
type无用于%flink.ssql,指定流式可视化类型(single、update、append)
refreshInterval3000用于%flink.ssql,指定流式数据可视化的前端刷新间隔
template{0}用于%flink.ssql,为single类型流式可视化指定 HTML 模板,可用{i}作为结果第 i 列的占位符
parallelism无用于%flink.ssql与%flink.bsql,指定 Flink SQL 作业并行度
maxParallelism无用于%flink.ssql与%flink.bsql,指定 Flink SQL 作业最大并行度(便于之后修改并行度)
savepointDir无若指定,在 Zeppelin 中取消 Flink 作业时会同时做 savepoint 并将状态存于此目录;恢复作业时从此 savepoint 恢复
execution.savepoint.path无恢复作业时从此 savepoint 路径恢复
resumeFromSavepoint无若指定savepointDir,则从 savepoint 恢复 Flink 作业
resumeFromLatestCheckpoint无若启用 checkpoint,则从最新 checkpoint 恢复
runAsOnefalse若为 true,所有 insert into SQL 在单个 Flink 作业中运行

关于 savepoint 恢复,源码setSavepointPathIfNecessary定义了明确的优先级:段落配置中记录的 savepoint 路径(段落被取消时由 Zeppelin 记录)→ 段落配置中的 checkpoint 路径(由作业进度轮询记录)→ 用户设置的本地属性execution.savepoint.path→ 否则移除execution.savepoint.path(见 FlinkScalaInterpreter.scala)。

教程笔记本

Zeppelin 内置了多个 Flink 教程笔记本(Flink Tutorial),包括 Flink Basics、Three Essential Steps for Building Flink Job、Flink Job Control Tutorial、Streaming ETL、Streaming Data Analytics、Batch ETL、Batch Data Analytics、Logistic Regression (Alink) 等,位于 notebook/Flink Tutorial 目录下,可作为深入学习更多特性的参考。

实践建议小结

  • 优先使用isolated per note的解释器绑定模式,避免多笔记本互相干扰;
  • 生产环境优先使用 yarn-application 模式(Flink 1.11+),避免过多解释器进程耗尽 Zeppelin 服务器资源;
  • 流式 SQL 结果可视化按场景选择 single(单行结果)、update(持续更新)、append(持续追加)三种模式;
  • 复杂 UDF 建议在 IDE 中编写并通过flink.udf.jars(配合flink.udf.jars.packages加速扫描)自动注册,避免重复手工注册;
  • 添加第三方依赖首选flink.execution.packages(需联网),离线环境改用flink.execution.jars;
  • 涉及 savepoint/checkpoint 恢复时,善用段落本地属性savepointDir、execution.savepoint.path、resumeFromSavepoint、resumeFromLatestCheckpoint组合。
  • 数据分析
  • 数据可视化
  • 大数据
  • 后端
  • 前端
  • 任务调度

【免费下载链接】zeppelin

Web-based notebook that enables>项目地址:https://gitcode.com/gh_mirrors/zeppe/zeppelin

点击查看免费下载

相关推荐

上一篇:RIOT OS 板级支持详解:STM32 Nucleo-F746ZG(ARM Cortex-M7 开发板)的资源配置、烧录与调试
下一篇:Topit窗口置顶工具:5个实用技巧彻底改变你的macOS多任务体验

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询