华为云数据湖探索 DLI 对接使用全攻略:从零基础到生产实践
在当今数据驱动决策的时代,企业面临着海量数据处理的挑战。传统的数据仓库和自建大数据平台往往需要投入大量的人力物力进行基础设施的维护和扩容。华为云数据湖探索(Data Lake Insight,简称DLI)作为一款Serverless的大数据分析服务,以其免运维、弹性伸缩、兼容开源生态等特性,正在成为越来越多企业数据分析的首选平台。本文将从零开始,全面讲解华为云DLI的对接使用方法,涵盖从服务开通到生产实践的完整流程,并辅以大量代码示例,帮助读者快速掌握DLI的核心技能。
需要先登录华为云控制台,点击:华为云控制台,还没有账号,点击:注册并关联,已有账号点击:登录后关联
一、认识华为云数据湖探索 DLI
1.1 什么是DLI
数据湖探索(Data Lake Insight,DLI)是华为云提供的一款完全兼容Apache Spark、Apache Flink、openLooKeng(基于Apache Presto)生态的Serverless大数据分析服务。DLI的核心特点是用户不需要管理任何服务器,即开即用,支持标准SQL、Spark SQL、Flink SQL等多种语法。它提供了批处理、流处理、交互式分析三位一体的融合处理能力,能够满足不同场景下的数据分析需求。
DLI支持多种数据接入方式,可以无缝对接华为云上的CloudTable、RDS、DWS、CSS等云服务,也可以连接ECS自建数据库以及线下数据库,实现异构数据的融合分析。数据无需复杂的抽取、转换、加载过程,使用SQL或程序就可以直接对数据进行探索和分析。
1.2 DLI的核心优势
Serverless免运维是DLI最显著的优势之一。用户无需关心集群的部署、扩容、运维等问题,DLI服务会自动管理计算资源的生命周期。DLI支持弹性资源池,计算资源可以根据作业负载自动伸缩,既保证了作业的及时响应,又避免了资源的浪费。
DLI兼容开源生态,支持Spark和Flink两大数据处理引擎,用户可以将已有的Spark或Flink作业平滑迁移到DLI平台。DLI还提供了丰富的开发工具支持,包括JDBC接口、Python SDK、Java SDK以及API Explorer等,方便开发者以多种方式与DLI进行交互。
二、DLI对接前的准备工作
在正式开始使用DLI之前,需要进行一系列的准备工作,包括服务开通、权限配置、资源创建等。这些准备工作是后续所有作业开发的基础。
2.1 开通DLI服务
首先需要登录华为云管理控制台,在服务列表中找到"数据湖探索 DLI",点击进入后按照页面提示开通服务。首次使用DLI的用户需要根据控制台的引导更新DLI委托,用于将操作权限委托给DLI服务,让DLI服务以用户的身份使用其他云服务,代替用户进行一些资源运维工作。该委托包含获取IAM用户相关信息、跨源场景访问和使用VPC、子网、路由、对等连接的权限、作业执行失败需要通过SMN发送通知消息的权限。
在DLI管理控制台的左侧导航栏单击"全局配置 > 服务授权",在委托设置页面勾选基础使用、跨源场景、运维场景的委托权限后,单击"更新委托权限"即可完成委托配置。
2.2 创建IAM用户并授权
对于企业用户,DLI支持通过统一身份认证服务(Identity and Access Management,简称IAM)进行精细的权限管理。通过IAM,可以在华为云账号中给员工创建IAM用户,并使用策略来控制他们对华为云资源的访问范围。例如,可以创建只具有DLI使用权限但不具有删除DLI等高危操作权限的用户。
DLI的权限管理包括角色(粗粒度授权)和策略(细粒度授权)两种方式。权限大类涵盖队列权限、数据权限(数据库、表、列)、作业权限(Flink作业)、程序包权限、跨源认证权限等多个维度。在实际生产环境中,建议遵循最小权限原则,只为用户授予完成工作所必需的最小权限集合。
2.3 创建弹性资源池和队列
使用DLI提交作业前,需要先创建弹性资源池,并在弹性资源池中创建队列,为提交作业准备所需的计算资源。弹性资源池是DLI计算资源的容器,队列则是实际执行作业的计算单元。
创建弹性资源池的步骤如下:登录DLI管理控制台,在左侧导航栏单击"资源管理 > 弹性资源池",进入弹性资源池管理页面。单击界面右上角的"购买弹性资源池",填写具体的弹性资源池参数,包括区域、可用区、规格等。参数填写完成后单击"立即购买",确认配置后单击"提交"完成弹性资源池的创建。
在弹性资源池的列表页,选择要操作的弹性资源池,单击操作列的"添加队列"。配置队列的基础参数,包括队列名称、队列类型等。队列类型分为SQL队列和通用队列两种:SQL队列用于执行SQL作业,通用队列用于执行Flink或Spark作业。创建完成后,队列就可以用于提交各类作业了。
2.4 配置DLI作业桶
DLI作业桶用于存储作业结果以及使用DLI服务产生的临时数据如作业日志等。在提交SQL作业前,建议先配置DLI作业桶。在DLI管理控制台的"全局配置 > 工程配置"中完成作业桶的配置。
配置作业桶后,系统会在运行SQL作业时把结果直接写到指定的OBS桶里。通过配置OBS桶的生命周期规则,可以实现定时删除OBS桶中的对象或者定时转换对象的存储类别,从而有效控制存储成本。
三、DLI的多种对接方式
DLI提供了多种对接方式,以适应不同开发场景和开发者习惯的需求。主要包括JDBC连接、SDK调用、控制台操作以及API调用等方式。
3.1 使用JDBC连接DLI
JDBC(Java Database Connectivity)是DLI提供的一种标准数据库连接方式,支持在Linux或Windows环境下使用JDBC应用程序连接DLI服务端提交作业。使用JDBC连接DLI提交的作业仅支持运行在Spark引擎上。
JDBC连接的准备工作:
首先需要在使用的机器中安装JDK,JDK版本为1.7或以上版本,并配置环境变量。然后下载DLI JDBC驱动包"huaweicloud-dli-jdbc-<version>.zip",解压后获得"huaweicloud-dli-jdbc-<version>-jar-with-dependencies.jar"。将该JAR文件添加至Java工程的classpath路径下。
JDBC连接配置:
DLI JDBC提供两种身份认证模式:Token认证和AK/SK认证。推荐使用AK/SK认证方式。连接DLI需要获取以下信息:
- Endpoint:在地区和终端节点页面获取DLI对应的Endpoint
- 项目编号:在华为云页面上方菜单栏单击用户名,然后在"我的凭证"页面获取项目编号
- 队列名称:在DLI管理控制台的"资源管理 > 队列管理"中查看
JDBC连接URL的格式为:jdbc:dl://{Endpoint}/{ProjectId}。例如:jdbc:dl://dli.cn-north-1.myhuaweicloud.com/96a17d961b84434baec6a58b9e567908。
JDBC连接代码示例:
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.ResultSet;
import java.sql.Statement;
import java.util.Properties;
public class DLIJDBCDemo {
public static void main(String[] args) {
// 配置连接参数
String endpoint = "dli.cn-north-1.myhuaweicloud.com";
String projectId = "your-project-id";
String queueName = "your-queue-name";
String ak = "your-access-key";
String sk = "your-secret-key";
String jdbcUrl = "jdbc:dl://" + endpoint + "/" + projectId;
Properties connectionProps = new Properties();
connectionProps.setProperty("user", ak);
connectionProps.setProperty("password", sk);
connectionProps.setProperty("queue", queueName);
try {
// 加载驱动
Class.forName("com.huawei.dli.jdbc.DliDriver");
// 建立连接
Connection conn = DriverManager.getConnection(jdbcUrl, connectionProps);
// 创建Statement
Statement stmt = conn.createStatement();
// 执行SQL查询
String sql = "SELECT * FROM your_table LIMIT 10";
ResultSet rs = stmt.executeQuery(sql);
// 处理结果集
while (rs.next()) {
System.out.println(rs.getString(1));
}
// 关闭资源
rs.close();
stmt.close();
conn.close();
} catch (Exception e) {
e.printStackTrace();
}
}
}
使用注意事项:使用JDBC 2.X版本时,对于2024年5月之前开通并使用DLI服务的用户,查询结果最多只能返回1000条。如需查看更多的作业结果,推荐开启队列的"查询结果写入桶"功能,从作业桶中查询超过1000条的数据。2024年5月起的新用户可以直接使用该功能,无需申请。
3.2 使用Python SDK操作DLI
DLI提供了Python SDK,方便Python开发者通过编程方式访问DLI服务。使用Python SDK需要先初始化DLI客户端,可以使用AK/SK或Token两种认证方式初始化客户端。其中Token认证仅DLI SDK V1版本支持,推荐使用AK/SK认证方式。
Python SDK安装与初始化:
# 安装DLI Python SDK(示例)
# pip install huaweicloud-dli-python-sdk
from dli.client import DliClient
# 使用AK/SK初始化客户端
client = DliClient(
ak='your-access-key',
sk='your-secret-key',
region='cn-north-1',
project_id='your-project-id'
)
# 提交SQL作业
sql = "SELECT * FROM your_table LIMIT 10"
queue_name = "your-queue-name"
result = client.submit_sql_job(sql, queue_name)
print(result)
DLI SDK还支持资源包的上传和管理。通过SDK提供的接口,可以上传资源包到DLI,kind参数用于指定资源包类型。
3.3 使用控制台操作DLI
对于不熟悉编程的用户或需要进行快速数据探索的场景,DLI管理控制台提供了直观的操作界面。控制台支持SQL作业编辑器、作业管理、资源管理、数据管理等核心功能。
在SQL编辑器页面,用户可以直接编写SQL语句执行数据查询操作。DLI的SQL编辑器支持SQL2003标准,兼容SparkSQL,可以批量执行SQL语句,并且作业编辑窗口的常用语法采用不同颜色突出显示,提升了开发体验。
控制台还提供了作业模板功能,用户可以将常用的SQL语句保存为模板,后续无需重新编写SQL语句,通过模板即可直接执行SQL操作。系统还提供了多条标准的TPC-H查询语句模板,用户可以直接使用。
四、DLI SQL作业开发
SQL作业是DLI最常用的作业类型之一,适用于使用标准SQL语句进行数据查询和分析的场景。DLI的SQL作业支持Spark和HetuEngine两种执行引擎。Spark引擎适用于离线分析场景,HetuEngine引擎适用于交互式分析场景。
4.1 创建数据库和表
在执行SQL作业前,需要先定义数据结构。DLI元数据是SQL作业、Spark作业场景开发的基础。元数据管理包括数据目录(Catalog)、数据库(Database)和表(Table)三个层次。
在DLI控制台的SQL编辑器中,可以执行DDL语句创建数据库和表:
-- 创建数据库
CREATE DATABASE IF NOT EXISTS my_database;
-- 使用数据库
USE my_database;
-- 创建DLI托管表
CREATE TABLE IF NOT EXISTS user_behavior (
user_id STRING COMMENT '用户ID',
action_type STRING COMMENT '行为类型',
action_time TIMESTAMP COMMENT '行为时间',
product_id STRING COMMENT '产品ID',
amount DOUBLE COMMENT '金额'
)
COMMENT '用户行为表'
PARTITIONED BY (dt STRING)
STORED AS PARQUET;
-- 创建OBS外表(数据存储在OBS中)
CREATE TABLE IF NOT EXISTS obs_user_behavior (
user_id STRING,
action_type STRING,
action_time TIMESTAMP,
product_id STRING,
amount DOUBLE
)
USING PARQUET
LOCATION 'obs://my-bucket/data/user_behavior/';
DLI支持在不迁移数据的情况下,直接对OBS中存储的数据进行查询分析。创建OBS外表时,数据实际存储在OBS中,DLI只管理元数据。
4.2 数据导入
使用DLI查询数据前,需要将数据文件上传至OBS中。DLI支持通过多种方式将数据导入到表中:
-- 使用INSERT INTO语句导入数据
INSERT INTO user_behavior PARTITION(dt='2026-07-17')
SELECT 'user001', 'click', TIMESTAMP '2026-07-17 10:00:00', 'prod001', 99.9
UNION ALL
SELECT 'user002', 'purchase', TIMESTAMP '2026-07-17 10:05:00', 'prod002', 199.9;
-- 从OBS导入数据到DLI表
LOAD DATA INPATH 'obs://my-bucket/data/input/user_behavior.csv'
INTO TABLE user_behavior
PARTITION(dt='2026-07-17');
4.3 数据查询与分析
数据导入完成后,就可以执行各种查询分析操作了:
-- 基础查询
SELECT * FROM user_behavior WHERE dt = '2026-07-17' LIMIT 100;
-- 聚合统计
SELECT
action_type,
COUNT(*) AS action_count,
SUM(amount) AS total_amount
FROM user_behavior
WHERE dt = '2026-07-17'
GROUP BY action_type
ORDER BY total_amount DESC;
-- 窗口函数分析
SELECT
user_id,
action_type,
action_time,
amount,
ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY action_time DESC) AS rn
FROM user_behavior
WHERE dt = '2026-07-17';
-- 多表关联分析
SELECT
u.user_id,
u.action_type,
p.product_name,
u.amount
FROM user_behavior u
JOIN products p ON u.product_id = p.product_id
WHERE u.dt = '2026-07-17';
-- 使用CASE WHEN进行条件判断
SELECT
user_id,
action_type,
CASE
WHEN amount > 1000 THEN '高消费'
WHEN amount > 100 THEN '中等消费'
ELSE '低消费'
END AS consumption_level
FROM user_behavior
WHERE dt = '2026-07-17';
五、DLI Flink作业开发
Flink作业是DLI提供的流处理能力,适用于实时数据分析和处理场景。DLI支持Flink SQL作业和Flink Jar作业两种类型。
5.1 Flink Jar作业开发
Flink Jar作业适用于需要自定义流处理逻辑、复杂的状态管理或特定库集成的数据分析场景。用户需要自行编写并构建Jar作业程序包,在提交Flink Jar作业前,将Jar作业程序包上传至OBS。
DLI控制台不提供Jar包的开发能力,用户需要在线下完成Jar包的开发。Jar包的开发可以参考DLI提供的Flink Jar开发基础样例。
Flink Jar作业开发流程:
- 开发Flink Jar作业程序,编译并打包为JAR文件
- 将JAR文件上传到OBS指定目录
- 在DLI控制台创建Flink Jar作业,配置JAR包路径、主类、参数等信息
- 提交作业并运行
Flink Jar作业代码示例(Java):
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.windowing.time.Time;
public class DLIStreamingJob {
public static void main(String[] args) throws Exception {
// 创建流执行环境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 从Socket接收数据(生产环境可从Kafka等数据源读取)
DataStream<String> text = env.socketTextStream("localhost", 9999);
// 数据转换和统计
DataStream<WordCount> counts = text
.map(new Tokenizer())
.keyBy(value -> value.word)
.timeWindow(Time.seconds(5))
.sum("count");
// 输出结果
counts.print();
// 执行作业
env.execute("DLI Streaming WordCount Job");
}
public static class Tokenizer implements MapFunction<String, WordCount> {
@Override
public WordCount map(String value) throws Exception {
return new WordCount(value, 1);
}
}
public static class WordCount {
public String word;
public int count;
public WordCount() {}
public WordCount(String word, int count) {
this.word = word;
this.count = count;
}
}
}
5.2 跨源场景的Flink作业
执行跨源场景的Flink作业,不能使用系统已有的default队列,需要创建通用队列。跨源分析场景中,可以使用DEW(数据加密服务)管理数据源的访问凭证,需要创建允许DLI访问DEW的委托。
使用Flink 1.15及以上版本的引擎执行作业时,可以在作业配置中添加自定义委托信息。在Flink作业运行过程中,作业程序可以获取配置的自定义委托信息,使用该委托来访问其他云服务。建议对自定义委托权限进行最小化处理,只配置作业运行过程中需要的必要权限。
六、DLI Spark作业开发
Spark作业是DLI提供的批处理能力,适用于大规模数据的离线处理和分析。DLI支持用户编写代码创建Spark作业来创建数据库、创建DLI表或OBS表和插入表数据等操作。
6.1 Spark作业开发流程
Spark作业的开发流程如下:
- 创建DLI通用队列作为作业运行的计算资源
- 配置OBS桶,上传数据文件,配置元数据信息存储路径
- 新建Maven工程,配置pom文件,引入DLI相关依赖
- 编写程序代码,实现数据处理逻辑
- 调试并编译代码,导出JAR包
- 将JAR包上传到OBS和DLI程序包中
- 在DLI控制台创建Spark Jar作业并提交运行
- 查看作业运行状态和日志
6.2 Spark作业代码示例
以下是一个完整的Spark作业示例,演示如何创建数据库、创建表和插入数据:
import org.apache.spark.sql.SparkSession;
public class DLISparkJob {
public static void main(String[] args) {
// 创建SparkSession
SparkSession spark = SparkSession.builder()
.appName("DLI Spark Job Example")
.getOrCreate();
// 创建数据库
spark.sql("CREATE DATABASE IF NOT EXISTS spark_demo");
// 使用数据库
spark.sql("USE spark_demo");
// 创建DLI表
spark.sql(
"CREATE TABLE IF NOT EXISTS spark_demo.sales (" +
" sale_id STRING," +
" product_name STRING," +
" sale_amount DOUBLE," +
" sale_date STRING" +
") " +
"USING PARQUET " +
"PARTITIONED BY (sale_date)"
);
// 插入数据
spark.sql(
"INSERT INTO spark_demo.sales PARTITION(sale_date='2026-07-17') " +
"VALUES ('S001', 'ProductA', 199.9, '2026-07-17'), " +
" ('S002', 'ProductB', 299.9, '2026-07-17')"
);
// 查询数据
spark.sql("SELECT * FROM spark_demo.sales").show();
// 关闭SparkSession
spark.stop();
}
}
该功能目前处于受限使用阶段,如需使用需要提交工单申请开通"使用Spark作业访问DLI元数据"的使用权限。如果使用Spark 3.1访问元数据,则必须新建队列。使用Spark 3.3.1访问元数据需自定义委托凭证,并在作业配置中添加委托信息。
七、跨源分析实践
DLI支持对云上CloudTable、RDS、DWS、CSS、ECS自建数据库以及线下数据库的异构数据进行探索。跨源访问的必要条件包括"DLI与数据源网络连通"和"DLI可获取数据源的访问凭证"。
7.1 配置网络连通
DLI与数据源网络连通需要配置增强型跨源连接。具体操作包括创建跨源连接、配置路由、安全组等。配置完成后,DLI队列的弹性资源网段需要添加到数据源的安全组入方向规则中。
7.2 跨源查询示例
-- 创建跨源连接后,可以直接在SQL中查询RDS中的数据
-- 假设已配置RDS数据源连接
SELECT * FROM rds_table
WHERE create_time > '2026-01-01';
-- 关联分析:将DLI中的数据和RDS中的数据进行JOIN
SELECT
d.user_id,
d.action_type,
r.user_name,
r.department
FROM dli_user_behavior d
JOIN rds_user_info r ON d.user_id = r.user_id
WHERE d.dt = '2026-07-17';
八、DLI安全管理最佳实践
数据安全是大数据平台的重中之重。DLI提供了多层次的安全能力,包括权限管理、数据加密、审计日志等。
8.1 加强权限管理
DLI服务支持通过IAM对用户权限进行精细化管理,通过设置不同的企业组织和操作权限,达到对DLI访问权限的隔离。建议创建IAM用户并授予最小必要权限,避免使用主账号进行日常操作。
8.2 数据加密存储
DLI服务支持使用OBS加密桶作为数据表存储。针对敏感数据,建议在创建DLI表时使用OBS加密桶作为数据表的存储。
8.3 数据备份与恢复
DLI服务提供了数据导入和导出接口,可以针对关键数据定期导出到OBS进行数据备份。建议每日或每周执行数据导出备份,当数据损坏时可以通过备份文件恢复数据。
8.4 安全审计日志
DLI服务支持日志审计功能,可以实时记录用户对数据的所有相关操作。通过对用户访问数据行为的记录、分析和汇报,可以帮助用户事后生成合规报告、进行事故追根溯源,提高数据资产安全性。
九、DLI成本优化实践
DLI的计费项包括存储费用与计算费用两项,计费类型包括包周期(包年包月)、套餐包和按需计费三种。合理规划和使用DLI资源可以有效控制成本。
9.1 选择合适的计费模式
SQL作业的计费包括存储计费和计算计费。包年包月计费根据购买周期进行扣费,推荐使用包年包月模式,价格优惠且在周期内独享计算资源。按需计费以小时为单位进行扣费,分为按CU时计费和按扫描数据量计费两种。建议优先选择按CU时计费,可资源独享且成本核算清晰。
CU时资费 = CU数 × 使用时长 × 单价,使用时长按自然小时计费,不足一个小时按一个小时计费。扫描数据量资费 = 执行SQL时产生的扫描数据量 × 单价,如果计算任务超时或失败则本次计算不收取费用。
9.2 使用DLI分析账单消费数据
华为云提供了使用DLI分析账单消费数据的实践指南。通过将账户的实际消费数据导入DLI进行分析,可以找出费用优化的空间。
具体操作步骤如下:
- 获取消费数据:下载账户的实际消费明细数据
- 分析账户消费结构:在DLI上分析消费数据,找出开支较大的资源或用户
- 制定优化措施:根据分析结果,针对性地进行成本优化
在DLI上创建消费分析表的示例:
CREATE TABLE `spending` (
account_period STRING,
EnterpriseProject STRING,
EnterpriseProjectID STRING,
accountID STRING,
product_type_code STRING,
product_type STRING,
product_code STRING,
product_name STRING,
product_id STRING,
mode STRING,
time1 STRING,
use_start STRING,
use_end STRING,
orderid STRING,
ordertime STRING,
resource_type STRING,
resource_id STRING,
resouce_name STRING,
tag STRING,
skuid STRING
);
十、常见问题解答
问1:DLI和自建Hadoop/Spark集群相比有什么优势?
答:DLI是Serverless服务,无需管理和维护集群,即开即用。DLI支持弹性伸缩,计算资源可以根据负载自动调整,避免了资源浪费。DLI还提供了统一的开发平台和丰富的运维能力,大大降低了大数据平台的使用门槛和运维成本。
问2:DLI支持哪些数据源?
答:DLI支持OBS、CloudTable、RDS、DWS、CSS等华为云服务,也支持ECS自建数据库以及线下数据库。通过跨源连接功能,可以实现对异构数据源的统一查询和分析。
问3:DLI的作业结果最多能返回多少条数据?
答:使用JDBC 2.X版本时,对于2024年5月之前开通并使用DLI服务的用户,查询结果最多只能返回1000条。建议开启队列的"查询结果写入桶"功能,从作业桶中查询超过1000条的数据。2024年5月起的新用户可以直接使用该功能。
问4:DLI如何保证数据安全?
答:DLI从多个维度保障数据安全:通过IAM进行精细的权限管理;支持使用OBS加密桶存储数据;提供数据导入导出接口支持数据备份;支持日志审计功能记录所有操作;支持使用DEW管理访问凭证。
问5:DLI的成本主要由哪些部分构成?如何优化?
答:DLI的成本包括存储费用和计算费用。优化建议包括:选择合适的计费模式(包年包月比按需更优惠);优先使用按CU时计费;合理规划队列规格避免资源浪费;使用生命周期管理自动清理过期数据;通过分析账单消费数据找出优化空间。
问6:DLI的SQL作业支持哪些执行引擎?如何选择?
答:DLI的SQL作业支持Spark和HetuEngine两种执行引擎。Spark引擎适用于离线分析场景,适合处理大规模数据的批处理任务。HetuEngine引擎适用于交互式分析场景,适合需要快速响应的即席查询。用户可以根据实际业务需求选择合适的引擎。




