- 数据分析
- 数据可视化
- 大数据
- 后端
- 前端
- 任务调度
【免费下载链接】zeppelin
Web-based notebook that enables>项目地址:https://gitcode.com/gh_mirrors/zeppe/zeppelin
本文以 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 集群环境:
| 名称 | 类 | 说明 |
|---|---|---|
%flink | FlinkInterpreter | 创建 ExecutionEnvironment / StreamExecutionEnvironment / BatchTableEnvironment / StreamTableEnvironment,并提供 Scala 环境 |
%flink.pyflink | PyFlinkInterpreter | 提供 Python 环境 |
%flink.ipyflink | IPyFlinkInterpreter | 提供 IPython 环境 |
%flink.ssql | FlinkStreamSqlInterpreter | 提供流式 SQL 环境 |
%flink.bsql | FlinkBatchSqlInterpreter | 提供批式 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.mode | local | Flink 执行模式:local/remote/yarn/yarn-application |
flink.execution.remote.host | 无 | 运行中 JobManager 的主机名,仅 remote 模式使用 |
flink.execution.remote.port | 无 | 运行中 JobManager 的端口,仅 remote 模式使用 |
jobmanager.memory.process.size | 1024m | JobManager 总内存大小(官方 Flink 属性) |
taskmanager.memory.process.size | 1024m | TaskManager 总内存大小(官方 Flink 属性) |
taskmanager.numberOfTaskSlots | 1 | 每个 TaskManager 的 slot 数 |
local.number-taskmanager | 4 | 本地模式下 TaskManager 总数 |
yarn.application.name | Zeppelin Flink Session | Yarn 应用名称 |
yarn.application.queue | default | Yarn 应用的队列名 |
源码实现细节(见 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.asLoginUser | true | 是否以 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.max | 10 | %flink.bsql批式 SQL 的最大并发数 |
zeppelin.flink.concurrentStreamSql.max | 10 | %flink.ssql流式 SQL 的最大并发数 |
zeppelin.pyflink.python | python | PyFlink 使用的 Python 可执行文件 |
table.exec.resource.default-parallelism | 1 | Flink SQL 作业的默认并行度 |
显示、Hive 与作业生命周期类属性
| 属性 | 默认值 | 说明 |
|---|---|---|
zeppelin.flink.scala.color | true | 是否彩色显示 Scala Shell 输出 |
zeppelin.flink.scala.shell.tmp_dir | 无 | 存放 Scala Shell 编译 jar 的临时目录 |
zeppelin.flink.enableHive | false | 是否启用 Hive |
zeppelin.flink.hive.version | 2.3.7 | 要连接的 Hive 版本 |
zeppelin.flink.module.enableHive | false | 是否启用 Hive module;若启用,Hive UDF 优先于 Flink UDF |
zeppelin.flink.maxResult | 1000 | SQL 解释器返回的最大行数 |
zeppelin.flink.job.check_interval | 1000 | 检查 Flink 作业进度的间隔(毫秒) |
flink.interpreter.close.shutdown_cluster | true | 关闭解释器时是否关闭 Flink 集群 |
zeppelin.interpreter.close.cancel_job | true | 关闭解释器时是否取消 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 中创建的内置变量:
| 变量 | 含义 |
|---|---|
senv | StreamExecutionEnvironment(流执行环境) |
benv | ExecutionEnvironment(批执行环境) |
stenv | StreamTableEnvironment(blink planner,即新 planner) |
btenv | BatchTableEnvironment(blink planner,即新 planner) |
z | ZeppelinContext |
源码中这些变量通过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_env | StreamExecutionEnvironment |
b_env | ExecutionEnvironment |
st_env | StreamTableEnvironment(blink planner,即新 planner) |
bt_env | BatchTableEnvironment(blink planner,即新 planner) |
配置 PyFlink
要让 PyFlink 在 Zeppelin 中工作,需要配置三件事:
- 安装 pyflink:例如
pip install apache-flink==1.11.1。若需要使用 PyFlink UDF,则必须在所有 TaskManager 节点上安装 pyflink——也就是说如果使用 yarn,所有 yarn 节点都需要安装 pyflink; - 复制 Python 文件夹:将
${FLINK_HOME}/opt下的python文件夹复制到${FLINK_HOME}/lib; - 设置
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 protobufZeppelinContext在 PyFlink 中同样可用,使用方式与 Flink Scala 几乎相同。IPython 的更多特性可参考 Python 解释器文档。
第三方依赖管理
无论使用 Scala、Python 还是 SQL 编写 Flink 作业,都很常见需要第三方依赖。在 IDE 中很容易添加依赖(例如在 pom.xml 中),但在 Zeppelin 中主要通过两个设置添加第三方依赖:
flink.execution.packagesflink.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 种方式:
- 编写 Scala UDF
- 编写 PyFlink UDF
- 通过 SQL 创建 UDF
- 通过
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.jarflink.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.jarZeppelin 会扫描该 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–*.jarflink-hadoop-compatibility_2.11–*.jarhive-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) |
refreshInterval | 3000 | 用于%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 恢复 |
runAsOne | false | 若为 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
相关推荐
Apache Zeppelin 与 Apache Flink 快速上手:Notebook 交互式流批开发实战指南
Apache Zeppelin 与 Apache Flink 快速上手:Notebook 交互式流批开发实战指南 本篇快速入门指南面向希望使用 Zeppelin
数据分析数据可视化大数据后端前端任务调度OmniVoice 训练数据准备完整指南:JSONL 清单、音频 Token 提取与 WebDataset 分片
OmniVoice 训练数据准备完整指南:JSONL 清单、音频 Token 提取与 WebDataset 分片 OmniVoice 是一款面向 600+ 语言
数据分析数据可视化大数据后端前端任务调度Apache Flink SQL 入门实战:从本地集群搭建到流式连续查询开发
Apache Flink SQL 入门实战:从本地集群搭建到流式连续查询开发 Flink SQL 允许开发者使用标准 SQL 语法开发流式数据处理应用,并保持
后端大数据流处理批处理