本文以 Apache Zeppelin 仓库中的 Cassandra CQL Interpreter 官方文档 为骨架,结合cassandra模块源码(CassandraInterpreter.java、InterpreterLogic.scala、ParagraphParser.scala 等),系统讲解如何在 Zeppelin 笔记本中使用%cassandra解释器直接执行 CQL 语句、描述集群 Schema、注入运行时查询参数、复用 Prepared Statement,并通过动态表单交互式查询 Cassandra。读完本文,你将掌握 Cassandra 解释器的完整命令体系、全部配置参数的含义与默认值,以及底层执行原理,可直接上手在笔记本中开展交互式数据分析。
解释器概览:%cassandra
Cassandra CQL 解释器是 Apache Zeppelin 官方提供的一个解释器模块,用于在笔记本段落中直接执行 Apache Cassandra 的 CQL 查询语言。其核心信息如下:
| Name | Class | Description |
|---|---|---|
%cassandra | CassandraInterpreter | 为 Apache Cassandra CQL 查询语言提供解释执行能力 |
该解释器基于 DataStax 官方Cassandra Java Driver构建,模块描述为 "Zeppelin: Apache Cassandra interpreter",依赖cassandra-driver-core3.0.1(见 cassandra/pom.xml)。从源码结构看,解释器采用 "Java 入口 + Scala 逻辑" 的双语言实现:Java 类CassandraInterpreter负责生命周期管理(open/close/interpret/getScheduler),Scala 类InterpreterLogic负责真正的话语解析与执行,ParagraphParser基于 Scala 组合子解析器(scala.util.parsing.combinator)将段落文本切分为命令块。
启用 Cassandra 解释器
在笔记本中启用Cassandra解释器:点击笔记本右上角的Gear(齿轮)图标,在弹出的解释器绑定对话框中勾选Cassandra即可。
Cassandra 解释器在构建产物中以zeppelin-cassandra_2.10的形式存在(见 cassandra/pom.xml 中的interpreter.name=cassandra),因此该模块会被 Zeppelin 在启动时自动识别并注册到解释器列表中。
使用 Cassandra 解释器
在段落的输入区,使用%cassandra选择 Cassandra 解释器,然后输入命令。例如:
%cassandra SELECT * FROM users LIMIT 10;要获取交互式帮助菜单,直接在段落中输入:
HELP;从源码看,HELP;命令由ParagraphParser中的HELP_PATTERN((?i)\s*HELP;\s*$,大小写不敏感)识别,并由EnhancedSession.execute(helpCmd)以%html格式输出帮助内容(见 EnhancedSession.scala)。
解释器命令体系
Cassandra 解释器接受以下五类命令:
| 命令类型 | 命令名称 | 描述 |
|---|---|---|
| Help 命令 | HELP | 显示交互式帮助菜单 |
| Schema 命令 | DESCRIBE KEYSPACE、DESCRIBE CLUSTER、DESCRIBE TABLES... | 用于描述 Cassandra Schema 的自定义命令 |
| Option 命令 | @consistency、@retryPolicy、@fetchSize... | 为段落中所有语句注入运行时选项 |
| Prepared statement 命令 | @prepare、@bind、@remove_prepare | 注册预编译语句,并稍后通过注入绑定值复用 |
| 原生 CQL 语句 | 所有 CQL 兼容语句(SELECT、INSERT、CREATE...) | 所有 CQL 语句直接对 Cassandra 服务器执行 |
在 ParagraphParser.scala 中,这五类命令都有对应的正则模式:运行时参数使用@consistency\s*=\s*(...)一类模式,@prepare、@bind、@remove_prepare使用方括号命名模式,Schema 命令则支持DESCRIBE/DESC两种缩写且全部大小写不敏感。
CQL 语句
该解释器兼容 Cassandra 支持的任何 CQL 语句,例如:
INSERT INTO users(login,name) VALUES('jdoe','John DOE'); SELECT * FROM users WHERE login='jdoe';每条语句必须以分号(;)分隔,但下列特殊命令除外(不要加分号):
@prepare@bind@remove_prepare@consistency@serialConsistency@timestamp@retryPolicy@fetchSize@requestTimeOut
多行语句以及同一行内的多条语句均被支持,只要以分号分隔即可。例如:
USE spark_demo; SELECT * FROM albums_by_country LIMIT 1; SELECT * FROM countries LIMIT 1; SELECT * FROM artists WHERE login='jlennon';批量(Batch)语句支持跨多行,DDL 语句(CREATE/ALTER/DROP)同样支持:
BEGIN BATCH INSERT INTO users(login,name) VALUES('jdoe','John DOE'); INSERT INTO users_preferences(login,account_type) VALUES('jdoe','BASIC'); APPLY BATCH; CREATE TABLE IF NOT EXISTS test( key int PRIMARY KEY, value text );CQL 语句是大小写不敏感的(列名和值除外),因此以下语句等价且合法:
INSERT INTO users(login,name) VALUES('jdoe','John DOE'); Insert into users(login,name) vAlues('hsue','Helen SUE');解释器对 CQL 的兼容性覆盖 Cassandra 1.2 / 2.0 / 2.1 / 2.2 / 3.x 多个版本(完整 CQL 语句清单可查阅 DataStax 各版本对应的 CQL 参考文档)。从源码看,通用语句的前缀模式GENERIC_STATEMENT_PREFIX覆盖了INSERT、UPDATE、DELETE、SELECT、CREATE、ALTER、DROP、GRANT、REVOKE、TRUNCATE、LIST、USE等关键字,且以(?is)开启大小写不敏感与 DOTALL 模式;批量语句则通过BATCH_PATTERN支持BEGIN BATCH、BEGIN UNLOGGED BATCH与BEGIN COUNTER BATCH三种类型。
语句中的注释
可以在语句之间添加注释。单行注释以**井号(#)或双斜杠(//)**开头,多行注释包裹在/**与*/之间。例如:
#Single line comment style 1 INSERT INTO users(login,name) VALUES('jdoe','John DOE'); //Single line comment style 2 /** Multi line comments **/ Insert into users(login,name) vAlues('hsue','Helen SUE');语法验证
解释器内置了一个语法验证器,但它只检查基本的语法错误;所有与 CQL 相关的语法校验都直接委托给Cassandra服务器本身。
大多数时候,语法错误源于语句之间缺少分号或拼写错误。从 InterpreterLogic.scala 的parseInput实现可见,当ParagraphParser解析失败时,会抛出包含提示语 "Did you forget to add ; (semi-colon) at the end of each CQL statement ?" 的InterpreterException,并在段落中以 ERROR 状态呈现。
Schema 命令:交互式 Schema 发现
为了让 Schema 发现更简单、更交互化,解释器提供了以下命令:
| 命令 | 描述 |
|---|---|
DESCRIBE CLUSTER; | 显示当前集群名称及其分区器(partitioner) |
DESCRIBE KEYSPACES; | 列出集群中所有已存在的 keyspace 及其配置(复制因子、durable write 等) |
DESCRIBE TABLES; | 列出集群中所有已存在的 keyspace,以及每个 keyspace 下的所有表名 |
DESCRIBE TYPES; | 列出集群中所有已存在的 keyspace,以及每个 keyspace 下的所有用户自定义类型(UDT)名 |
DESCRIBE FUNCTIONS; | 列出集群中所有已存在的 keyspace,以及每个 keyspace 下的所有函数名 |
DESCRIBE AGGREGATES; | 列出集群中所有已存在的 keyspace,以及每个 keyspace 下的所有聚合函数名 |
DESCRIBE MATERIALIZED VIEWS; | 列出集群中所有已存在的 keyspace,以及每个 keyspace 下的所有物化视图名 |
DESCRIBE KEYSPACE <keyspace_name>; | 描述给定 keyspace 的配置及其全部表的细节(名称、列等) |
DESCRIBE TABLE (<keyspace_name>).<table_name>; | 描述给定表。若未提供 keyspace,则使用当前登录的 keyspace;若未登录任何 keyspace,则使用默认的 system keyspace;若未找到表,则抛出错误消息 |
DESCRIBE TYPE (<keyspace_name>).<type_name>; | 描述给定的类型(UDT)。keyspace 缺省规则同上;若未找到类型,则抛出错误消息 |
DESCRIBE FUNCTION (<keyspace_name>).<function_name>; | 描述给定函数。keyspace 缺省规则同上;若未找到函数,则抛出错误消息 |
DESCRIBE AGGREGATE (<keyspace_name>).<aggregate_name>; | 描述给定聚合函数。keyspace 缺省规则同上;若未找到聚合函数,则抛出错误消息 |
DESCRIBE MATERIALIZED VIEW (<keyspace_name>).<view_name>; | 描述给定物化视图。keyspace 缺省规则同上;若未找到视图,则抛出错误消息 |
Schema 对象(cluster、keyspace、table、type、function 和 aggregate)以表格形式展示。左上角有下拉菜单用于展开对象详情,右上角菜单展示图标图例。
在实现层面,上述命令由 EnhancedSession.scala 中的execute(describeXxx)系列方法完成:它们读取session.getCluster.getMetadata,将元数据(如KeyspaceMetadata、TableMetadata、MaterializedViewMetadata等)交给 DisplaySystem.scala 渲染为%html交互界面。元数据展示模型(列类型分区键/聚类列/静态列/普通列、索引、复制配置等)定义在 MetaDataHierarchy.scala 中。对于 "keyspace 缺省则回退到当前登录 keyspace、否则使用 system" 的规则,源码中的实现是describeTable.keyspace.orElse(Option(session.getLoggedKeyspace)).getOrElse("system")。
运行时参数:向语句注入查询选项
有时你需要向语句传递运行时的查询参数。这些参数不属于 CQL 规范,是解释器特有的扩展。完整参数列表如下:
| 参数 | 语法 | 描述 |
|---|---|---|
| Consistency Level | @consistency=<value> | 将给定的一致性级别应用到段落中所有查询 |
| Serial Consistency Level | @serialConsistency=<value> | 将给定的串行一致性级别应用到段落中所有查询 |
| Timestamp | @timestamp=<long value> | 将给定的时间戳应用到段落中所有查询。注意:直接在 CQL 语句中传递的 timestamp 值会覆盖此参数 |
| Retry Policy | @retryPolicy=<value> | 将给定的重试策略应用到段落中所有查询 |
| Fetch Size | @fetchSize=<integer value> | 将给定的 fetch size 应用到段落中所有查询 |
| Request Time Out | @requestTimeOut=<integer value> | 将给定的请求超时时间(毫秒)应用到段落中所有查询 |
部分参数只接受受限取值:
| 参数 | 可选值 |
|---|---|
| Consistency Level | ALL, ANY, ONE, TWO, THREE, QUORUM, LOCAL_ONE, LOCAL_QUORUM, EACH_QUORUM |
| Serial Consistency Level | SERIAL, LOCAL_SERIAL |
| Timestamp | 任意 long 值 |
| Retry Policy | DEFAULT, DOWNGRADING_CONSISTENCY, FALLTHROUGH, LOGGING_DEFAULT, LOGGING_DOWNGRADING, LOGGING_FALLTHROUGH |
| Fetch Size | 任意整数 |
请注意:不要在每个参数语句末尾添加分号(;)。
使用示例——通过@timestamp控制写入时间戳:
CREATE TABLE IF NOT EXISTS spark_demo.ts( key int PRIMARY KEY, value text ); TRUNCATE spark_demo.ts; // Timestamp in the past @timestamp=10 // Force timestamp directly in the first insert INSERT INTO spark_demo.ts(key,value) VALUES(1,'first insert') USING TIMESTAMP 100; // Select some data to make the clock turn SELECT * FROM spark_demo.albums LIMIT 100; // Now insert using the timestamp parameter set at the beginning(10) INSERT INTO spark_demo.ts(key,value) VALUES(1,'second insert'); // Check for the result. You should see 'first insert' SELECT value FROM spark_demo.ts WHERE key=1;关于查询参数,有以下几点重要说明:
- 同一个段落中可以设置多个查询参数;
- 如果同一个查询参数被多次设置且值不同,解释器只取第一个值;
- 每个查询参数作用于同一段落中的所有 CQL 语句,除非你用纯 CQL 文本覆盖该选项(例如用
USING子句强制指定 timestamp);- 查询参数与 CQL 语句之间的先后顺序无关紧要。
上述行为在源码中有精确对应:ParagraphParser为每个参数定义了严格的校验正则(如TIMESTAMP_PATTERN = """^\s*@timestamp\s*=\s*([0-9]+)\s*$"""),而 InterpreterLogic.scala 的extractQueryOptions使用.headOption只取同一类参数的第一个值,applyQueryOptions则把解析后的选项逐一应用到每条Statement(setConsistencyLevel、setSerialConsistencyLevel、setDefaultTimestamp、setRetryPolicy、setFetchSize、setReadTimeoutMillis)。其中重试策略在源码中映射为 DataStax 驱动策略:DEFAULT→Policies.defaultRetryPolicy()、DOWNGRADING_CONSISTENCY→DowngradingConsistencyRetryPolicy.INSTANCE、FALLTHROUGH→FallthroughRetryPolicy.INSTANCE,以及对应的三种 Logging 包装策略。
支持 Prepared Statement
出于性能考虑,最好预先准备(prepare)语句,之后通过提供绑定值来复用。解释器提供 3 条命令处理 prepared 与 bound 语句:
@prepare@bind@remove_prepare
基本语法:
@prepare[statement-name]=... @bind[statement-name]='text', 1223, '2015-07-30 12:00:01', null, true, ['list_item1', 'list_item2'] @bind[statement-name-with-no-bound-value] @remove_prepare[statement-name]@prepare
使用语法@prepare[statement-name]=SELECT...创建预编译语句。statement-name 是必需的:解释器使用 Java 驱动准备给定的语句,并将生成的PreparedStatement保存到内部哈希表中,以提供的statement-name作为查找键。
注意:由于 Cassandra 解释器只有一个实例,这个内部预编译语句表被所有笔记本和所有段落共享。
如果解释器遇到多个使用相同statement-name(键)的
@prepare,只有第一条语句会被采纳。
示例:
@prepare[select]=SELECT * FROM spark_demo.albums LIMIT ? @prepare[select]=SELECT * FROM spark_demo.artists LIMIT ?对于上述示例,预编译语句是SELECT * FROM spark_demo.albums LIMIT ?,而SELECT * FROM spark_demo.artists LIMIT ?被忽略,因为预编译语句表中已存在键为select的条目。
在Zeppelin的语境下,笔记本可能被调度为定期执行,因此避免反复预编译同一条语句(这被视为反模式)是必要的。从源码看,InterpreterLogic使用new ConcurrentHashMap[String,PreparedStatement]存放预编译语句,并通过preparedStatements.getOrElseUpdate(statement.name, session.prepare(statement.query))实现 "已存在则不重复 prepare" 的语义(见 InterpreterLogic.scala)。
@bind
一旦语句被预编译(可能在另一个笔记本/段落中),你就可以为它绑定值:
@bind[select_first]=10@bind语句并不强制要求提供绑定值。但如果提供,它们必须符合以下语法规则:
- 字符串值应包裹在单引号(')之间;
- 日期值应包裹在单引号(')之间,并遵守如下格式(完整列表见 DataStax 的 timestamp 类型参考文档):
yyyy-MM-dd HH:MM:ssyyyy-MM-dd HH:MM:ss.SSS
null原样解析;- 布尔值(
true/false)原样解析; - 集合值必须遵循标准 CQL 语法:
- list:
['list_item1', 'list_item2', ...] - set:
{'set_item1', 'set_item2', ...} - map:
{'key1': 'val1', 'key2': 'val2', ...}
- list:
- tuple值应包裹在圆括号之间(见 CQL 的 Tuple 语法):
('text', 123, true) - udt值应包裹在花括号之间(见 CQL 的 UDT 语法):
{stree_name: 'Beverly Hills', number: 104, zip_code: 90020, state: 'California', ...}
也可以在 batch 内部使用
@bind:BEGIN BATCH @bind[insert_user]='jdoe','John DOE' UPDATE users SET age = 27 WHERE login='hsue'; APPLY BATCH;
从源码看,绑定值的类型转换发生在InterpreterLogic.createBoundStatement:解释器根据 prepared statement 的变量元数据(ps.getVariables)逐个匹配类型,通过驱动CodecRegistry解析文本值——TEXT/VARCHAR/ASCII去除引号、INT/VARINT转Int、BIGINT/COUNTER转Long、BLOB转ByteBuffer、BOOLEAN转Boolean、DECIMAL/DOUBLE/FLOAT转对应数字类型、INET转InetAddress、TIMESTAMP按上述两种格式解析、UUID/TIMEUUID转java.util.UUID、LIST/SET/MAP/UDT/TUPLE则交给 BoundValuesParser.scala 解析后再由 codec 转换。若绑定值的数量与 prepared statement 的变量数量不一致,会抛出明确错误。
@remove_prepare
为避免预编译语句永远驻留在预编译语句表中,可以使用@remove_prepare[statement-name]语法将其移除。移除一个不存在的预编译语句不会产生错误。
使用动态表单
与其硬编码 CQL 查询,不如使用Zeppelin 动态表单语法来注入简单值或多项选择表单。
旧式 mustache 语法({{ }})用于绑定输入文本和下拉选择表单仍然受支持,但已弃用,将在未来版本中移除。
Legacy(旧式语法)简单参数的语法是:
{{input_Label=default value}}。默认值是必需的,因为段落第一次执行时,我们在渲染表单之前就发起 CQL 查询,所以至少需要一个值。多项选择参数的语法是:
{{input_Label=value1 | value2 | ... | valueN}}。默认情况下,段落第一次执行时使用第一个选项执行 CQL 查询。
示例:
#Secondary index on performer style SELECT name, country, performer FROM spark_demo.performers WHERE name='${performer=Sheryl Crow|Doof|Fanfarlo|Los Paranoia}' AND styles CONTAINS '${style=Rock}';在上述示例中,第一次 CQL 查询将以performer='Sheryl Crow' AND style='Rock'执行。后续查询时,你可以直接用表单修改值。
请注意:我们把
${ }块包裹在单引号(')之间,因为 Cassandra 在这里期望一个字符串。也可以使用${style='Rock'}语法,但此时表单上显示的值是'Rock'而不是Rock。
动态表单同样可用于预编译语句:
@bind[select]=='${performer=Sheryl Crow|Doof|Fanfarlo|Los Paranoia}', '${style=Rock}'在底层,InterpreterLogic.scala 的maybeExtractVariables实现动态表单替换:先到 Zeppelin 的AngularObjectRegistry中按变量名(noteId+paragraphId作用域)查找已有值,若找到则直接替换;否则通过context.getGui.select(...)或context.getGui.input(...)动态创建表单控件。这也解释了为何文档强调 FormType.SIMPLE 下的表单是程序化动态生成的。CassandraInterpreter.getFormType()返回FormType.SIMPLE,正是动态表单的入口。
共享状态与并行执行
可以并行执行多个段落。但在后端,解释器仍然使用同步查询;只有能在InterpreterResult中返回Future值时才可能实现异步执行——这是 Zeppelin 项目的一个潜在改进点。
最近,Zeppelin允许你为解释器选择隔离级别(参见 解释器绑定模式)。简单来说,你有 3 种可用绑定:
- shared(共享):所有笔记本使用同一个 JVM和同一个解释器实例;
- scoped(作用域):同一个 JVM,但为每个笔记本创建不同的解释器实例;
- isolated(隔离):每个笔记本运行独立的 JVM和单个解释器实例。
使用shared绑定时,所有笔记本和段落共享同一个com.datastax.driver.core.Session对象。因此,如果你使用USE keyspace_name;语句登录某个 keyspace,它将改变所有当前使用该解释器的用户的 keyspace,因为每个 Cassandra 解释器实例只创建一个Session对象。
同样的提醒也适用于预编译语句哈希表:它被所有使用同一个 Cassandra 解释器实例的用户共享。
使用scoped绑定时,在同一个 JVM中Zeppelin会创建多个 Cassandra 解释器实例,从而创建多个Session对象。注意使用此绑定时的资源和内存开销!
isolated模式最极端:每个不同的笔记本都会创建一个 JVM 与Session对象。
从源码看,CassandraInterpreter.getScheduler()使用SchedulerFactory创建一个并行调度器,并发度由配置项cassandra.interpreter.parallelism(默认 10)控制——这就是 "多个段落可并行执行" 的实现基础。
解释器配置
要配置Cassandra解释器,进入Interpreter菜单并向下滚动修改参数。Cassandra 解释器使用官方Cassandra Java Driver,大多数参数都用于配置该 Java 驱动。
配置参数及其默认值如下:
| Property Name | Description | Default Value |
|---|---|---|
cassandra.cluster | 要连接的 Cassandra 集群名称 | Test Cluster |
cassandra.compression.protocol | 线上传输压缩协议。可选值:NONE、SNAPPY、LZ4 | NONE |
cassandra.credentials.username | 若启用了安全认证,提供用户名 | none |
cassandra.credentials.password | 若启用了安全认证,提供密码 | none |
cassandra.hosts | 逗号分隔的 Cassandra 主机(DNS 名称或 IP 地址)。例如:192.168.0.12,node2,node3 | localhost |
cassandra.interpreter.parallelism | 可并发执行的段落(查询块)数量 | 10 |
cassandra.keyspace | 要连接的默认 keyspace。强烈建议保留默认值,并在所有查询中给表名加上实际 keyspace 前缀 | system |
cassandra.load.balancing.policy | 负载均衡策略。默认 =new TokenAwarePolicy(new DCAwareRoundRobinPolicy())。要指定自定义策略,提供策略的_全限定类名(FQCN)_,运行时会通过Class.forName(FQCN)实例化 | DEFAULT |
cassandra.max.schema.agreement.wait.second | Cassandra 最大 schema 一致性等待时间(秒) | 10 |
cassandra.pooling.core.connection.per.host.local | 协议 V2 及以下默认 = 2;协议 V3 及以上默认 = 1 | 2 |
cassandra.pooling.core.connection.per.host.remote | 协议 V2 及以下默认 = 1;协议 V3 及以上默认 = 1 | 1 |
cassandra.pooling.heartbeat.interval.seconds | Cassandra 连接池心跳间隔(秒) | 30 |
cassandra.pooling.idle.timeout.seconds | Cassandra 空闲超时时间(秒) | 120 |
cassandra.pooling.max.connection.per.host.local | 协议 V2 及以下默认 = 8;协议 V3 及以上默认 = 1 | 8 |
cassandra.pooling.max.connection.per.host.remote | 协议 V2 及以下默认 = 2;协议 V3 及以上默认 = 1 | 2 |
cassandra.pooling.max.request.per.connection.local | 协议 V2 及以下默认 = 128;协议 V3 及以上默认 = 1024 | 128 |
cassandra.pooling.max.request.per.connection.remote | 协议 V2 及以下默认 = 128;协议 V3 及以上默认 = 256 | 128 |
cassandra.pooling.new.connection.threshold.local | 协议 V2 及以下默认 = 100;协议 V3 及以上默认 = 800 | 100 |
cassandra.pooling.new.connection.threshold.remote | 协议 V2 及以下默认 = 100;协议 V3 及以上默认 = 200 | 100 |
cassandra.pooling.pool.timeout.millisecs | Cassandra 连接池获取连接超时时间(毫秒) | 5000 |
cassandra.protocol.version | Cassandra 二进制协议版本 | 4 |
cassandra.query.default.consistency | Cassandra 查询默认一致性级别。可选值:ONE、TWO、THREE、QUORUM、LOCAL_ONE、LOCAL_QUORUM、EACH_QUORUM、ALL | ONE |
cassandra.query.default.fetchSize | Cassandra 查询默认 fetch size | 5000 |
cassandra.query.default.serial.consistency | Cassandra 查询默认串行一致性级别。可选值:SERIAL、LOCAL_SERIAL | SERIAL |
cassandra.reconnection.policy | Cassandra 重连策略。默认 =new ExponentialReconnectionPolicy(1000, 10 * 60 * 1000)。要指定自定义策略,提供 FQCN,运行时会通过Class.forName(FQCN)实例化 | DEFAULT |
cassandra.retry.policy | Cassandra 重试策略。默认 =DefaultRetryPolicy.INSTANCE。要指定自定义策略,提供 FQCN,运行时会通过Class.forName(FQCN)实例化 | DEFAULT |
cassandra.socket.connection.timeout.millisecs | Cassandra socket 默认连接超时时间(毫秒) | 500 |
cassandra.socket.read.timeout.millisecs | Cassandra socket 读取超时时间(毫秒) | 12000 |
cassandra.socket.tcp.no_delay | Cassandra socket TCP no delay | true |
cassandra.speculative.execution.policy | Cassandra 推测执行策略。默认 =NoSpeculativeExecutionPolicy.INSTANCE。要指定自定义策略,提供 FQCN,运行时会通过Class.forName(FQCN)实例化 | DEFAULT |
cassandra.ssl.enabled | 是否启用对配置了 SSL 的 Cassandra 的连接支持。要连接配置了 SSL 的 Cassandra,请使用true并提供如下 truststore 文件与密码选项 | false |
cassandra.ssl.truststore.path | 用于连接 SSL 版 Cassandra 的 truststore 文件路径 | (空) |
cassandra.ssl.truststore.password | 用于连接 SSL 版 Cassandra 的 truststore 文件密码 | (空) |
在上述文档列出的参数之外,从 CassandraInterpreter.java 的常量定义还可以看到几个额外的可配置项,说明文档表格并非配置全集:
cassandra.native.port:Cassandra 原生协议端口,默认9042(DEFAULT_PORT = "9042")。源码在open()中通过parseInt(getProperty(CASSANDRA_PORT))读取并传给Cluster.builder().withPort(port);cassandra.query.default.idempotence:查询默认幂等性,由JavaDriverConfig.getQueryOptions读取并设置;cassandra.socket.keep.alive、cassandra.socket.received.buffer.size.bytes、cassandra.socket.send.buffer.size.bytes、cassandra.socket.reuse.address、cassandra.socket.soLinger:额外的 socket 选项,仅在配置值非空时生效(见 JavaDriverConfig.scala 中的isNotBlank判断)。
SSL 连接的实现细节也值得注意:当cassandra.ssl.enabled=true时,CassandraInterpreter.open()会以 JKS 格式加载cassandra.ssl.truststore.path指定的 truststore,初始化TrustManagerFactory,构建 TLS 的SSLContext,并通过JdkSSLOptions.builder().withSSLContext(...)挂接到Cluster.Builder。
此外,连接池默认值会随协议版本自动调整:JavaDriverConfig.getProtocolVersion根据cassandra.protocol.version(1/2/3/4)分别设置对应版本的连接池默认值(例如 V3/V4 协议下 local 主机最大连接数为 1、单连接最大请求数为 1024),这与文档表格中 "协议 V2 及以下 / V3 及以上" 的说明完全吻合。
变更日志
3.0(Zeppelin 0.8.x):
- 更新文档与交互式文档
- 支持二进制协议V4
- 实现新的
@requestTimeOut运行时选项 - 升级 Java 驱动版本至3.0.1
- 允许解释器在使用 FormType.SIMPLE 时程序化添加动态表单
- 允许动态表单使用默认的 Zeppelin 语法
- 修复 FallThroughPolicy 拼写错误
- 在创建动态表单前先查找 AngularObjectRegistry 中的数据
- 补充对
ALTER语句的支持
2.0(Zeppelin 0.7.x):
- 更新帮助菜单并添加 changelog
- 支持用户自定义函数(UDF)、用户自定义聚合(UDA)与物化视图(Materialized Views)
- 升级 Java 驱动版本至3.0.0-rc1
1.0(Zeppelin 0.5.5-incubating):
- 初始版本
源码阅读指引
如果你希望深入理解本文涉及的所有机制,可以直接阅读cassandra模块下的以下关键文件:
- CassandraInterpreter.java:解释器生命周期、Cluster/Session 构建、SSL 支持、并行调度器;
- InterpreterLogic.scala:段落解析入口、查询选项提取、prepared statement 管理与绑定值类型转换、动态表单替换;
- ParagraphParser.scala:基于正则与组合子解析的完整命令语法定义;
- EnhancedSession.scala:DESCRIBE 系列 Schema 命令的执行与 HTML 渲染;
- JavaDriverConfig.scala:Socket/Query/Pooling 选项与策略的装配;
- BoundValuesParser.scala:
@bind绑定值(list/set/map/tuple/udt)的语法解析; - ParagraphParserTest.scala:覆盖混合语句、注释、一致性级别解析等核心场景的测试用例,是理解命令语法行为的最佳佐证。
如果你在使用该解释器时遇到 bug,可以到 Apache Zeppelin 的 JIRA 系统提交 ticket,并附上复现场景,以便社区跟进修复。
- 后端
- 前端
- 大数据
- 数据分析
【免费下载链接】zeppelin
Web-based notebook that enables>项目地址:https://gitcode.com/gh_mirrors/zeppelin2/zeppelin
相关推荐
Apache Zeppelin Cassandra CQL 解释器完全指南:连接配置、Schema 描述、Prepared Statement 与动态表单实战
Apache Zeppelin Cassandra CQL 解释器完全指南:连接配置、Schema 描述、Prepared Statement 与动态表单实战
数据分析数据可视化大数据后端Apache Zeppelin Cassandra CQL 解释器完全指南:从连接配置到预编译语句与格式化输出
Apache Zeppelin Cassandra CQL 解释器完全指南:从连接配置到预编译语句与格式化输出 Apache Zeppelin 内置的 %cas
数据分析数据可视化大数据后端前端任务调度OpenProject 平板分屏导航实战:iPad 与 Android 大屏上的高效工作包审阅工作流
OpenProject 平板分屏导航实战:iPad 与 Android 大屏上的高效工作包审阅工作流 OpenProject 移动应用(Beta)在 iPad
后端前端大数据数据分析