Spark SQL 查询引擎连接器

Kudu

师成师成· 更新于 2026-09-29· 阅读 20 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

什么是 Apache Kudu

Apache Kudu 是开源 Apache Hadoop 生态系统的新增组件,它补全了 Hadoop 的存储层,从而实现对快速数据的快速分析。

在阅读本文档时,我们假设你不一定熟悉 Apache Kudu。但至少你需要有一个可连接的正在运行的 Kudu 集群。如果你还能了解 Apache Kudu 的能力,那就更好了。

如果本页缺少关于 Apache Kudu 的背景知识,你可以参考其官方网站。

为什么在 Kudu 上使用 Kyuubi

基本上,Kyuubi 可以取代 HiveServer2,作为 Hadoop 上的多租户即席 SQL 解决方案,并凭借 Spark SQL 带来的速度与性能优势,你可以针对数据源和 Hive 表运行 SQL 查询,且只需消耗你被授权的计算资源即可访问受保护的数据。

Spark SQL 支持通过 DataFrame 接口操作多种数据源。DataFrame 可以使用关系转换进行操作,也可以用来创建临时视图。将 DataFrame 注册为临时视图后,你就可以在其数据上运行 SQL 查询。本节介绍使用 Spark 数据源加载和保存数据的通用方法,然后深入介绍内置数据源可用的具体选项。

在 Kyuubi 中,我们可以将 Kudu 表和其他数据源表注册为 Spark 临时视图,从而实现跨 Hive、Kudu 和其他数据源的联邦联合查询。

Kudu 与 Apache Spark 的集成

在将 Kyuubi 与 Kudu 集成之前,我们强烈建议你先完成 Spark 与 Kudu 的集成和测试。你可以参考 Kudu 的在线文档指南 —— Kudu 与 Spark 的集成

Kudu 与 Kyuubi 的集成

安装 Kudu Spark 依赖

确认你的 Kudu 集群版本,下载对应的 kudu spark 依赖库,例如 org.apache.kudu:kudu-spark3_2.12-1.14.0,并放置到 $SPARK_HOME/jars 目录下。

启动 Kyuubi

现在,你可以使用这个内嵌了 Kudu 的 Spark 发行版来启动 Kyuubi 服务。

启动 Beeline 或你偏好的其他客户端

bin/beeline -u 'jdbc:hive2://<host>:<port>/;principal=<if kerberized>;#spark.yarn.queue=kyuubi_test'

将 Kudu 表注册为 Spark 临时视图

CREATE TEMPORARY VIEW kudutest
USING kudu
options (
  kudu.master "ip1:port1,ip2:port2,...",
  kudu.table "kudu::test.testtbl")
0: jdbc:hive2://spark5.jd.163.org:10009/> show tables;
19/07/09 15:28:03 INFO ExecuteStatementInClientMode: Running query 'show tables' with 1104328b-515c-4f8b-8a68-1c0b202bc9ed
19/07/09 15:28:03 INFO KyuubiSparkUtil$: Application application_1560304876299_3805060 has been activated
19/07/09 15:28:03 INFO ExecuteStatementInClientMode: Executing query in incremental mode, running 1 jobs before optimization
19/07/09 15:28:03 INFO ExecuteStatementInClientMode: Executing query in incremental mode, running 1 jobs without optimization
19/07/09 15:28:03 INFO DAGScheduler: Asked to cancel job group 1104328b-515c-4f8b-8a68-1c0b202bc9ed
+-----------+-----------------------------+--------------+--+
| database  |          tableName          | isTemporary  |
+-----------+-----------------------------+--------------+--+
| kyuubi    | hive_tbl                    | false        |
|           | kudutest                    | true         |
+-----------+-----------------------------+--------------+--+
2 rows selected (0.29 seconds)

查询 Kudu 表

0: jdbc:hive2://spark5.jd.163.org:10009/> select * from kudutest;
19/07/09 15:25:17 INFO ExecuteStatementInClientMode: Running query 'select * from kudutest' with ac3e8553-0d79-4c57-add1-7d3ffe34ba16
19/07/09 15:25:17 INFO KyuubiSparkUtil$: Application application_1560304876299_3805060 has been activated
19/07/09 15:25:17 INFO ExecuteStatementInClientMode: Executing query in incremental mode, running 3 jobs before optimization
19/07/09 15:25:17 INFO ExecuteStatementInClientMode: Executing query in incremental mode, running 3 jobs without optimization
19/07/09 15:25:17 INFO DAGScheduler: Asked to cancel job group ac3e8553-0d79-4c57-add1-7d3ffe34ba16
+---------+---------------+----------------+--+
| userid  | sharesetting  | notifysetting  |
+---------+---------------+----------------+--+
| 1       | 1             | 1              |
| 5       | 5             | 5              |
| 2       | 2             | 2              |
| 3       | 3             | 3              |
| 4       | 4             | 4              |
+---------+---------------+----------------+--+
5 rows selected (1.083 seconds)

将 Kudu 表与 Hive 表进行 Join

0: jdbc:hive2://spark5.jd.163.org:10009/> select t1.*, t2.* from hive_tbl t1 join kudutest t2 on t1.userid=t2.userid+1;
19/07/09 15:31:01 INFO ExecuteStatementInClientMode: Running query 'select t1.*, t2.* from hive_tbl t1 join kudutest t2 on t1.userid=t2.userid+1' with 6982fa5c-29fa-49be-a5bf-54c935bbad18
19/07/09 15:31:01 INFO KyuubiSparkUtil$: Application application_1560304876299_3805060 has been activated
<omitted lines.... >
19/07/09 15:31:01 INFO DAGScheduler: Asked to cancel job group 6982fa5c-29fa-49be-a5bf-54c935bbad18
+---------+---------------+----------------+---------+---------------+----------------+--+
| userid  | sharesetting  | notifysetting  | userid  | sharesetting  | notifysetting  |
+---------+---------------+----------------+---------+---------------+----------------+--+
| 2       | 2             | 2              | 1       | 1             | 1              |
| 3       | 3             | 3              | 2       | 2             | 2              |
| 4       | 4             | 4              | 3       | 3             | 3              |
+---------+---------------+----------------+---------+---------------+----------------+--+
3 rows selected (1.63 seconds)

插入 Kudu 表

请注意,Kudu 仅支持 INSERT INTO,不支持 OVERWRITE 数据。

0: jdbc:hive2://spark5.jd.163.org:10009/> insert overwrite table kudutest select *  from hive_tbl;
19/07/09 15:35:29 INFO ExecuteStatementInClientMode: Running query 'insert overwrite table kudutest select *  from hive_tbl' with 1afdb791-1aa7-4ceb-8ba8-ff53c17615d1
19/07/09 15:35:29 INFO KyuubiSparkUtil$: Application application_1560304876299_3805060 has been activated
19/07/09 15:35:30 ERROR ExecuteStatementInClientMode:
Error executing query as bdms_hzyaoqin,
insert overwrite table kudutest select *  from hive_tbl
Current operation state RUNNING,
java.lang.UnsupportedOperationException: overwrite is not yet supported
    at org.apache.kudu.spark.kudu.KuduRelation.insert(DefaultSource.scala:424)
    at org.apache.spark.sql.execution.datasources.InsertIntoDataSourceCommand.run(InsertIntoDataSourceCommand.scala:42)
    at org.apache.spark.sql.execution.command.ExecutedCommandExec.sideEffectResult$lzycompute(commands.scala:70)
    at org.apache.spark.sql.execution.command.ExecutedCommandExec.sideEffectResult(commands.scala:68)
    at org.apache.spark.sql.execution.command.ExecutedCommandExec.executeCollect(commands.scala:79)
    at org.apache.spark.sql.Dataset$$anonfun$6.apply(Dataset.scala:190)
    at org.apache.spark.sql.Dataset$$anonfun$6.apply(Dataset.scala:190)
    at org.apache.spark.sql.Dataset$$anonfun$52.apply(Dataset.scala:3259)
    at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:77)
    at org.apache.spark.sql.Dataset.withAction(Dataset.scala:3258)
    at org.apache.spark.sql.Dataset.<init>(Dataset.scala:190)
    at org.apache.spark.sql.Dataset$.ofRows(Dataset.scala:75)
    at org.apache.spark.sql.SparkSQLUtils$.toDataFrame(SparkSQLUtils.scala:39)
    at org.apache.kyuubi.operation.statement.ExecuteStatementInClientMode.execute(ExecuteStatementInClientMode.scala:152)
    at org.apache.kyuubi.operation.statement.ExecuteStatementOperation$$anon$1$$anon$2.run(ExecuteStatementOperation.scala:74)
    at org.apache.kyuubi.operation.statement.ExecuteStatementOperation$$anon$1$$anon$2.run(ExecuteStatementOperation.scala:70)
    at java.security.AccessController.doPrivileged(Native Method)
    at javax.security.auth.Subject.doAs(Subject.java:422)
    at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1698)
    at org.apache.kyuubi.operation.statement.ExecuteStatementOperation$$anon$1.run(ExecuteStatementOperation.scala:70)
    at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
    at java.util.concurrent.FutureTask.run(FutureTask.java:266)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
    at java.lang.Thread.run(Thread.java:745)


19/07/09 15:35:30 INFO DAGScheduler: Asked to cancel job group 1afdb791-1aa7-4ceb-8ba8-ff53c17615d1
0: jdbc:hive2://spark5.jd.163.org:10009/> insert into table kudutest select * from hive_tbl;
19/07/09 15:36:26 INFO ExecuteStatementInClientMode: Running query 'insert into table kudutest select *  from hive_tbl' with f7460400-0564-4f98-93b6-ad76e579e7af
19/07/09 15:36:26 INFO KyuubiSparkUtil$: Application application_1560304876299_3805060 has been activated
<omitted lines ...>
19/07/09 15:36:27 INFO DAGScheduler: ResultStage 36 (foreachPartition at KuduContext.scala:332) finished in 0.322 s
19/07/09 15:36:27 INFO DAGScheduler: Job 36 finished: foreachPartition at KuduContext.scala:332, took 0.324586 s
19/07/09 15:36:27 INFO KuduContext: completed upsert ops: duration histogram: 33.333333333333336%: 2ms, 66.66666666666667%: 64ms, 100.0%: 102ms, 100.0%: 102ms
19/07/09 15:36:27 INFO ExecuteStatementInClientMode: Executing query in incremental mode, running 1 jobs before optimization
19/07/09 15:36:27 INFO ExecuteStatementInClientMode: Executing query in incremental mode, running 1 jobs without optimization
19/07/09 15:36:27 INFO DAGScheduler: Asked to cancel job group f7460400-0564-4f98-93b6-ad76e579e7af
+---------+--+
| Result  |
+---------+--+
+---------+--+
No rows selected (0.611 seconds)

参考资料

https://kudu.apache.org/ https://kudu.apache.org/docs/developing.html#_kudu_integration_with_spark apache/kyuubi https://spark.apache.org/docs/latest/sql-data-sources.html

评论

登录后参与评论

正在加载评论…