Apache Zeppelin Flink 解释器实战指南:从本地开发到 Yarn 集群的批流一体分析 数据分析数据可视化大数据后端前端任务调度【免费下载链接】zeppelinWeb-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 集群环境名称类说明%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 的 classpathlib 目录中。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 ShellFlink 解释器的入口负责创建 Flink 程序的全部入口对象如 ExecutionEnvironment、StreamExecutionEnvironment 和 TableEnvironment并负责编译运行 Scala 代码和 SQLPython ShellPyFlink 的入口负责编译运行 Python 代码。从源码看FlinkScalaInterpreter见 FlinkScalaInterpreter.scala在open()时依次完成初始化 Flink 配置 → 创建 FlinkILoopScala REPL→ 创建批/流 TableEnvironment → 绑定 ZeppelinContextz→ 注册 JobListener用于把作业与段落关联、上报进度→ 可选注册 Hive catalog → 自动加载 UDF jar。配置详解Flink 解释器通过 Zeppelin 提供的属性进行配置下表同时你也可以添加表外其他的 Flink 属性以table.exec、parallelism等 Flink 官方配置项为例参考 Flink 官方文档的 Available Properties 列表。从源码看解释器启动时会把解释器属性中所有条目写入 Flink 的Configurationproperties.asScala.foreach(entry configuration.setString(...))因此任何 Flink 配置项都能以解释器属性方式下发。环境类属性属性默认值说明FLINK_HOME必填Flink 安装位置。必须指定否则无法在 Zeppelin 中使用 FlinkHADOOP_CONF_DIR必填yarn 模式Hadoop 配置目录位置yarn 模式下必须设置HIVE_CONF_DIR可选Hive 配置目录位置需要连接 Hive metastore 时必须设置执行模式与集群资源类属性属性默认值说明flink.execution.modelocalFlink 执行模式local/remote/yarn/yarn-applicationflink.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_DIRisYarnApplicationMode分支因此无需再手动指定这些目录。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.0SQL 并发与 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是否启用 Hivezeppelin.flink.hive.version2.3.7要连接的 Hive 版本zeppelin.flink.module.enableHivefalse是否启用 Hive module若启用Hive UDF 优先于 Flink UDFzeppelin.flink.maxResult1000SQL 解释器返回的最大行数zeppelin.flink.job.check_interval1000检查 Flink 作业进度的间隔毫秒flink.interpreter.close.shutdown_clustertrue关闭解释器时是否关闭 Flink 集群zeppelin.interpreter.close.cancel_jobtrue关闭解释器时是否取消 Flink 作业源码佐证zeppelin.flink.maxResult直接用于构造FlinkZeppelinContextz.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批执行环境stenvStreamTableEnvironmentblink planner即新 plannerbtenvBatchTableEnvironmentblink planner即新 plannerzZeppelinContext源码中这些变量通过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 PlannerZeppelin 0.11 之后移除了对 flink planner旧 planner的支持——Flink 1.14 之后也移除了旧 planner。因此当前仓库仅使用 blink planner新 planner。流式 WordCount 示例在 Zeppelin 中可以直接编写任意 Scala 代码例如经典的流式 WordCount 示例%flink // 使用 senvStreamExecutionEnvironment编写流式作业代码补全在 Zeppelin 中按 Tab 键即可触发代码补全。从源码看补全能力来自FlinkILoop的Completion组件scalaCompletion。ZeppelinContextZeppelinContext提供了一些附加函数与工具详细内容可参考 Zeppelin-Context 文档。在 Flink 解释器中可以用z展示 Flink 的 Dataset/Table例如z.show(DataSet)展示批式 DataSetz.show(Batch Table)展示批式 Tablez.show(Stream Table)展示流式 Table。Flink SQLZeppelin 提供两类 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 的查询。PyFlinkPyFlink 是 Flink on Zeppelin 的 Python 入口。内部 Flink 解释器会创建 Python Shell并创建 Flink 的环境变量ExecutionEnvironment、StreamExecutionEnvironment 等。需要注意PyFlink 背后的 Java 环境是在 Scala Shell 中创建的即 Scala Shell 与 Python Shell 共享同一个底层环境。Python Shell 中创建的变量变量含义s_envStreamExecutionEnvironmentb_envExecutionEnvironmentst_envStreamTableEnvironmentblink planner即新 plannerbt_envBatchTableEnvironmentblink planner即新 planner配置 PyFlink要让 PyFlink 在 Zeppelin 中工作需要配置三件事安装 pyflink例如pip install apache-flink1.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.jarsflink.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} doesnt exist。Flink UDFZeppelin 中定义 UDF 有 4 种方式编写 Scala UDF编写 PyFlink UDF通过 SQL 创建 UDF通过flink.udf.jars配置 UDF jarScala 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、appendrefreshInterval3000用于%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组合。赞分享数据分析数据可视化大数据后端前端任务调度【免费下载链接】zeppelinWeb-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 语法开发流式数据处理应用并保持后端大数据流处理批处理上一篇RIOT OS 板级支持详解STM32 Nucleo-F746ZGARM Cortex-M7 开发板的资源配置、烧录与调试下一篇Topit窗口置顶工具5个实用技巧彻底改变你的macOS多任务体验创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考