第 8 章 广电用户数据存储与处理的程序开发 完整版
一、章节整体概述
1. 章节定位与作用
本章是 Hive 学习从手工命令行操作向企业级程序自动化处理的核心转折点。前 7 章均通过 Hive CLI 手动输入 HQL 完成建表、查询、清洗,存在效率低、无法复用、难以集成到业务系统的局限。
本章通过 Java 语言基于 JDBC 接口实现 Hive 远程调用,将手工操作封装为可复用、可自动化执行的程序,完成广电数据全流程自动化存储、查询与处理,是衔接大数据平台操作与 Java 后端开发的关键章节,也是后续大数据项目开发的核心基础。
2. 课时安排(共 10 学时)
| 课时类型 | 时长 | 核心内容 |
|---|---|---|
| 理论课 | 2 学时 | HiveServer2 原理、JDBC 接口体系、工具类封装设计思想、客户端配置机制 |
| 实验课 1 | 3 学时 | Hive 远程服务配置、IDEA 开发环境搭建、配置文件部署、连接测试 |
| 实验课 2 | 3 学时 | 广电数据存储程序开发、数据查询与清洗程序开发 |
| 调试与答疑 | 2 学时 | 程序调试方法、问题排查、作业验收、优化拓展 |
3. 三维教学目标
| 目标类型 | 具体要求 |
|---|---|
| 知识目标 | 1. 理解 HiveServer2 的作用、架构与优势2. 掌握 JDBC 核心接口的功能与使用场景3. 理解工具类封装的设计思想与资源管理原则4. 区分服务端与客户端配置,理解resources目录存放配置的原理 |
| 能力目标 | 1. 能独立完成 Hive 远程服务的配置与启动2. 能在 IDEA 中搭建 Hive Java 开发环境,正确部署客户端配置文件3. 能通过 JDBC 实现 Hive 建库、建表、数据加载、查询、清洗全流程4. 能使用 IDEA 调试工具定位并解决程序问题5. 能识别并修复基础的资源泄漏、语法错误问题 |
| 素养目标 | 1. 建立大数据程序开发的工程化思维2. 培养代码规范与资源管理意识3. 提升问题排查与自主调试能力 |
4. 前置环境准备
(1)集群服务端准备(Linux)
- Hadoop 3.1.4 集群搭建完成并正常运行
- MySQL 服务正常运行,作为 Hive 元数据存储
- Hive 3.1.2 安装配置完成,基础 CLI 操作可用
- 广电业务 5 张原始表的 CSV 数据已上传至
/opt/data目录 - 集群所有节点关闭防火墙,或开放 9083、10000 等关键端口
(2)本地开发端准备(Windows)
- Windows 系统安装 JDK 1.8 并配置环境变量
- 安装 IntelliJ IDEA 2021.3 及以上版本
- 本地能 ping 通集群 master 节点,已配置
hosts主机名映射 - 提前从集群拷贝 3 份配置文件:
core-site.xml、hdfs-site.xml、hive-site.xml
二、模块一:配置 Hive 远程服务(服务端操作)
1. 核心原理
(1)HiveServer vs HiveServer2
- HiveServer:基于 Thrift 协议的初代远程服务,仅支持单客户端并发,无法满足多用户场景,已被淘汰。
- HiveServer2:Hive 0.11.0 重写后的版本,支持多客户端并发访问、身份认证,提供 JDBC、ODBC 等标准 API,是 Java/Python 等语言远程操作 Hive 的标准入口。
(2)两大核心服务
| 服务名称 | 核心作用 | 启动顺序 | 默认端口 | 依赖服务 |
|---|---|---|---|---|
| metastore | Hive 元数据服务,管理库、表、字段、存储位置等元数据,所有 Hive 操作都依赖元数据 | 必须先启动 | 9083 | MySQL、HDFS |
| hiveserver2 | Hive 远程服务,对外提供 JDBC 连接入口,接收客户端 HQL 请求并执行 | 依赖 metastore,后启动 | 10000 | metastore 服务 |
⚠️ 教学重点强调:必须严格按照「metastore → hiveserver2」的顺序启动,且启动后等待 30 秒再测试,服务初始化需要时间。
2. 服务端三大核心配置文件
以下配置文件部署在 Linux 集群节点上,是 Hadoop、Hive 服务启动的基础。
(1)core-site.xml(Hadoop 核心配置)
部署路径:/usr/local/hadoop-3.1.4/etc/hadoop/core-site.xml(所有节点都要有)
完整配置:
<?xml version="1.0" encoding="UTF-8"?>
<?xml-stylesheet type="text/xsl" href="configuration.xsl"?>
<configuration>
<!-- 默认文件系统,指定NameNode地址 -->
<property>
<name>fs.defaultFS</name>
<value>hdfs://master:9820</value>
</property>
<!-- Hadoop临时目录 -->
<property>
<name>hadoop.tmp.dir</name>
<value>/var/log/hadoop/tmp</value>
<description>A base for other temporary directories.</description>
</property>
<!-- 允许root用户从任意主机代理请求 -->
<property>
<name>hadoop.proxyuser.root.hosts</name>
<value>*</value>
</property>
<!-- 允许root用户代理任意用户组 -->
<property>
<name>hadoop.proxyuser.root.groups</name>
<value>*</value>
</property>
</configuration>关键说明:
- 两个
hadoop.proxyuser配置是 HiveServer2 能正常访问 HDFS 的核心,*为通配,生产环境建议配置具体 IP 和用户组。
(2)hdfs-site.xml(HDFS 配置)
部署路径:/usr/local/hadoop-3.1.4/etc/hadoop/hdfs-site.xml(所有节点都要有)
完整配置:
<?xml version="1.0" encoding="UTF-8"?>
<?xml-stylesheet type="text/xsl" href="configuration.xsl"?>
<configuration>
<!-- NameNode元数据存储目录 -->
<property>
<name>dfs.namenode.name.dir</name>
<value>file:///data/hadoop/hdfs/name</value>
</property>
<!-- DataNode数据存储目录 -->
<property>
<name>dfs.datanode.data.dir</name>
<value>file:///data/hadoop/hdfs/data</value>
</property>
<!-- SecondaryNameNode访问地址 -->
<property>
<name>dfs.namenode.secondary.http-address</name>
<value>master:9868</value>
</property>
<!-- 数据副本数量 -->
<property>
<name>dfs.replication</name>
<value>3</value>
</property>
<!-- 开启WebHDFS,支持HTTP方式访问HDFS -->
<property>
<name>dfs.webhdfs.enabled</name>
<value>true</value>
</property>
</configuration>(3)hive-site.xml(Hive 核心配置)
部署路径:/usr/local/hive-3.1.2/conf/hive-site.xml(master 节点)
完整配置:
<?xml version="1.0"?>
<?xml-stylesheet type="text/xsl" href="configuration.xsl"?>
<configuration>
<!-- 元数据MySQL连接地址 -->
<property>
<name>javax.jdo.option.ConnectionURL</name>
<value>jdbc:mysql://master:3306/hive?createDatabaseIfNotExist=true</value>
</property>
<!-- MySQL驱动类 -->
<property>
<name>javax.jdo.option.ConnectionDriverName</name>
<value>com.mysql.jdbc.Driver</value>
</property>
<!-- 元数据持久化工厂类 -->
<property>
<name>javax.jdo.PersistenceManagerFactoryClass</name>
<value>org.datanucleus.api.jdo.JDOPersistenceManagerFactory</value>
</property>
<!-- 提交后自动分离 -->
<property>
<name>javax.jdo.option.DetachAllOnCommit</name>
<value>true</value>
</property>
<!-- 非事务读取 -->
<property>
<name>javax.jdo.option.NonTransactionalRead</name>
<value>true</value>
</property>
<!-- MySQL用户名 -->
<property>
<name>javax.jdo.option.ConnectionUserName</name>
<value>root</value>
</property>
<!-- MySQL密码 -->
<property>
<name>javax.jdo.option.ConnectionPassword</name>
<value>123456</value>
</property>
<!-- 多线程支持 -->
<property>
<name>javax.jdo.option.Multithreaded</name>
<value>true</value>
</property>
<!-- 连接池类型 -->
<property>
<name>datanucleus.connectionPoolingType</name>
<value>BoneCP</value>
</property>
<!-- Hive数据仓库在HDFS上的存储路径 -->
<property>
<name>hive.metastore.warehouse.dir</name>
<value>/user/hive/warehouse</value>
</property>
<!-- HiveServer2服务端口 -->
<property>
<name>hive.server2.thrift.port</name>
<value>10000</value>
</property>
<!-- HiveServer2绑定主机 -->
<property>
<name>hive.server2.thrift.bind.host</name>
<value>master</value>
</property>
<!-- 远程元数据服务地址 -->
<property>
<name>hive.metastore.uris</name>
<value>thrift://master:9083</value>
</property>
<!-- Hive临时目录(需手动创建) -->
<property>
<name>system:java.io.tmpdir</name>
<value>/usr/local/hive-3.1.2/bin/hive-data/tmp</value>
</property>
<!-- 系统用户名 -->
<property>
<name>system:user.name</name>
<value>root</value>
</property>
</configuration>关键说明:
hive.metastore.uris:配置独立 metastore 服务地址,是启动远程元数据服务的核心。system:java.io.tmpdir:指定的目录必须提前手动创建,否则启动服务会报错。
3. 完整配置与启动流程
步骤 1:配置 Hadoop 代理用户权限
在 master 节点修改core-site.xml,添加两个代理用户配置项。
步骤 2:分发配置到所有子节点
scp /usr/local/hadoop-3.1.4/etc/hadoop/core-site.xml slave1:/usr/local/hadoop-3.1.4/etc/hadoop/
scp /usr/local/hadoop-3.1.4/etc/hadoop/core-site.xml slave2:/usr/local/hadoop-3.1.4/etc/hadoop/
scp /usr/local/hadoop-3.1.4/etc/hadoop/core-site.xml slave3:/usr/local/hadoop-3.1.4/etc/hadoop/步骤 3:重启 Hadoop 集群
cd /usr/local/hadoop-3.1.4/sbin
./stop-all.sh
./start-all.sh
./mr-jobhistory-daemon.sh start historyserver步骤 4:创建 Hive 临时目录
mkdir -p /usr/local/hive-3.1.2/bin/hive-data/tmp步骤 5:启动 Hive 远程服务(严格顺序)
cd /usr/local/hive-3.1.2/bin
# 1. 后台启动元数据服务
hive --service metastore &
# 执行后按回车回到命令行
# 2. 查看进程,确认出现RunJar进程
jps
# 3. 后台启动HiveServer2,日志输出到nohup.out
nohup hive --service hiveserver2 &
# 执行后按回车回到命令行
# 4. 再次查看进程,确认两个RunJar进程都存在
jps步骤 6:验证服务可用性
方法 1:端口监听验证
netstat -nltp | grep 10000方法 2:beeline 客户端验证
beeline
!connect jdbc:hive2://master:10000/default
# 输入用户名root,密码123456,能正常进入命令行则服务正常4. 高频踩坑与解决方案
| 问题现象 | 根本原因 | 解决方法 |
|---|---|---|
| 启动 metastore 报错,找不到目录 | system:java.io.tmpdir指定的临时目录未创建 | 手动创建对应目录 |
连接时报Failed to open new session | Hadoop 代理用户配置未生效 | 检查 core-site.xml,确认已分发所有节点并重启 Hadoop |
报Connection refused | 服务启动顺序错误,或服务未完全初始化 | 严格按 metastore→hiveserver2 顺序启动,等待 30 秒再测试 |
| Windows 无法连接 | 防火墙拦截,或本地 hosts 未配置 | 关闭所有节点防火墙;本地添加 master 主机名映射 |
| beeline 连接报错 | MySQL 服务未启动,或元数据库异常 | 检查 MySQL 服务状态,验证 hive-site.xml 的数据库配置 |
三、模块二:搭建 Hive Java 远程开发环境(本地操作)
1. 创建 Maven 项目
- 打开 IDEA,点击「New Project」
- 左侧选择「Maven」,Project SDK 选择 1.8,点击「Next」
- 项目名填写
HiveJavaAPI,选择存储路径,点击「Finish」
2. pom.xml 依赖配置
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>org.example</groupId>
<artifactId>HiveJavaAPI</artifactId>
<version>1.0-SNAPSHOT</version>
<properties>
<maven.compiler.source>8</maven.compiler.source>
<maven.compiler.target>8</maven.compiler.target>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
</properties>
<!-- 阿里云镜像,加速依赖下载 -->
<repositories>
<repository>
<id>aliyun-central</id>
<name>阿里云公共仓库</name>
<url>https://maven.aliyun.com/repository/public</url>
</repository>
</repositories>
<dependencies>
<!-- Hadoop公共基础依赖 -->
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-common</artifactId>
<version>3.1.4</version>
</dependency>
<!-- Hive执行核心依赖 -->
<dependency>
<groupId>org.apache.hive</groupId>
<artifactId>hive-exec</artifactId>
<version>3.1.2</version>
</dependency>
<!-- Hive JDBC驱动 -->
<dependency>
<groupId>org.apache.hive</groupId>
<artifactId>hive-jdbc</artifactId>
<version>3.1.2</version>
<!-- 排除冲突包,避免启动报错 -->
<exclusions>
<exclusion>
<groupId>org.glassfish</groupId>
<artifactId>javax.el</artifactId>
</exclusion>
<exclusion>
<groupId>org.eclipse.jetty</groupId>
<artifactId>jetty-runner</artifactId>
</exclusion>
</exclusions>
</dependency>
</dependencies>
</project>配置完成后,右键pom.xml → 「Maven」 → 「Reload project」,等待依赖下载完成。
3. 重点知识点:客户端配置文件部署
(1)为什么要放在src/main/resources?
这是大数据 Java 开发的标准规范,核心原因有 4 点:
Maven 机制:自动加入类路径
Maven 项目中,
src/main/resources是标准资源目录,编译、运行、打包时,会自动将目录下所有文件复制到 Java 运行类路径(classpath)。Hadoop/Hive 框架启动时,会默认从 classpath 扫描加载配置文件,无需额外写代码。客户端需要识别集群信息
Java 程序作为 Hive 客户端,需要通过配置文件知道:
HDFS 集群地址(
hdfs://master:9820)Hive 服务端口、元数据地址
代理用户权限规则
没有这些配置,客户端无法正常连接集群、访问 HDFS。
代码零改动,跨环境兼容
放在
resources后,框架自动加载,代码无需手动指定文件路径。换电脑、换集群,只需替换 3 个 XML 文件,不用改 Java 代码,可移植性极强。打包部署友好
项目打包成 Jar 包后,
resources下的配置文件会一起打入 Jar 中,上传到服务器直接运行,无需单独拷贝配置文件。
(2)服务端配置 vs 客户端配置 对比
| 类型 | 部署位置 | 作用 | 内容要求 |
|---|---|---|---|
| 服务端配置 | Linux 集群节点的 hadoop/conf、hive/conf 目录 | 启动 Hadoop、Hive、metastore、hiveserver2 服务 | 服务运行的完整配置 |
| 客户端配置 | 本地 IDEA 项目的src/main/resources目录 | 给 Java 程序读取,指导客户端连接集群 | 与服务端配置完全一致,直接拷贝即可 |
操作要求:从 Linux 集群原样拷贝 3 个 XML 文件,不要修改内容,保证客户端与集群配置一致。
(3)标准项目目录结构
HiveJavaAPI
└── src
└── main
├── java // 存放Java代码(ConnectionTest、HiveHelper等)
└── resources // 存放3个客户端配置文件
├── core-site.xml
├── hdfs-site.xml
└── hive-site.xml(4)不放配置文件的典型报错
- 代理用户权限异常:
User: root is not allowed to impersonate root - HDFS 地址无法识别:
Unknown host master、无效的 HDFS URI - 连接 HiveServer2 失败、LOAD DATA 异常
- 部分功能可用,但读写 HDFS 数据报错
4. JDBC 核心接口详解
JDBC(Java DataBase Connectivity)是 Java 访问数据库的统一标准 API,Hive 提供了专属的 JDBC 驱动。
表格
| 接口 / 类 | 核心作用 | 核心方法 | 使用场景 |
|---|---|---|---|
Driver | JDBC 驱动接口,Hive 实现类为HiveDriver | Class.forName("全类名") | 加载 Hive 驱动 |
DriverManager | 驱动管理工具类 | getConnection(url, user, pwd) | 创建数据库连接 |
Connection | 数据库连接会话 | createStatement()、close() | 管理连接、创建执行对象 |
Statement | 静态 SQL 执行对象 | execute()、executeQuery() | 执行 DDL、DML、查询语句 |
ResultSet | 查询结果集 | next()、getString() | 遍历查询返回的数据 |
易错点区分:
execute():返回 boolean,执行 DDL(建库、建表)、DML(加载数据)语句executeQuery():返回 ResultSet,仅执行 SELECT 查询语句,不可混用
5. 连接测试程序编写
(1)高频易错点:类名不能为Connection
教材示例类名为Connection,会与java.sql.Connection接口重名,导致编译冲突,必须修改为ConnectionTest,这是上机第一大高频报错点。
(2)完整测试代码
在src/main/java下创建ConnectionTest.java:
java
运行
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.Statement;
/**
* Hive连接测试类
* 注意:类名禁止命名为Connection,否则与JDK接口冲突
*/
public class ConnectionTest {
public static void main(String[] args) {
// 连接核心参数
String driver = "org.apache.hive.jdbc.HiveDriver";
String url = "jdbc:hive2://master:10000/default";
String username = "root";
String password = "123456";
Connection conn = null;
Statement stmt = null;
try {
// 1. 加载JDBC驱动
Class.forName(driver);
System.out.println("✅ Hive驱动加载成功");
// 2. 获取数据库连接
conn = DriverManager.getConnection(url, username, password);
System.out.println("✅ Hive连接成功:" + conn);
// 3. 创建Statement执行对象
stmt = conn.createStatement();
// 4. 执行HQL,创建测试数据库
String sql = "CREATE DATABASE IF NOT EXISTS test";
stmt.execute(sql);
System.out.println("✅ 测试数据库test创建成功");
} catch (Exception e) {
System.err.println("❌ 程序执行失败:");
e.printStackTrace();
} finally {
// 5. 逆序关闭资源:先Statement,后Connection
try {
if (stmt != null) stmt.close();
if (conn != null) conn.close();
System.out.println("✅ 连接已关闭");
} catch (Exception e) {
e.printStackTrace();
}
}
}
}(3)运行验证标准
- 控制台输出
Process finished with exit code 0 - 依次打印「驱动加载成功、连接成功、数据库创建成功、连接已关闭」
- Hive CLI 中执行
SHOW DATABASES;能看到test数据库
6. 环境搭建常见问题
表格
| 报错信息 | 原因 | 解决方法 |
|---|---|---|
ClassNotFoundException: org.apache.hive.jdbc.HiveDriver | 依赖未下载,或未刷新 Maven | 执行 Reload project,检查网络和镜像配置 |
| 类名 Connection 报红、编译错误 | 类名与 JDK 接口重名 | 修改类名为 ConnectionTest |
Connection refused: connect | 服务未启动、端口不通、主机名错误 | 检查服务状态、网络连通性、hosts 配置 |
| 代理用户权限报错 | 客户端配置文件未放 resources | 将 3 个 XML 文件放入 src/main/resources 目录 |
四、模块三:编写程序实现广电数据存储
1. 工具类封装设计思想
- 封装复用:将连接、关闭、建库、建表等通用操作封装为方法,避免重复代码
- 统一配置:连接地址、账号等常量统一维护,修改方便
- 异常处理:每个方法内部捕获异常,便于定位问题
- 单例连接:同一个工具类对象复用连接,减少连接创建开销
2. 广电 5 张业务表说明
表格
| 表名 | 业务含义 | 核心字段 |
|---|---|---|
mediamatch_usermsg | 用户基本信息表 | 用户编号、业务类型、用户状态、用户等级、地址、开户时间 |
mediamatch_userevent | 用户状态变更表 | 用户编号、状态名称、变更时间、用户等级、开户时间 |
mmconsume_billevents | 用户账单表 | 用户编号、账期、应付金额、优惠金额、业务类型 |
order_index | 用户订单表 | 用户编号、订单号、产品名称、订单时间、费用 |
media_index | 用户收视行为表 | 用户编号、节目名称、观看时长、开始 / 结束时间、资源类型 |
3. HiveHelper 工具类(教材原版)
java
运行
import java.sql.*;
/**
* Hive操作工具类(教材原版)
*/
public class HiveHelper {
// 连接配置常量
private static String driverName = "org.apache.hive.jdbc.HiveDriver";
private static String url = "jdbc:hive2://master:10000/DEFAULT";
private static String username = "root";
private static String password = "123456";
// 成员变量:连接、执行对象、结果集
private Connection conn = null;
private Statement stmt = null;
private ResultSet rs = null;
/**
* 获取数据库连接
*/
public Connection getConn() throws ClassNotFoundException, SQLException {
if (null == conn) {
Class.forName(driverName);
conn = DriverManager.getConnection(url, username, password);
}
return conn;
}
/**
* 关闭连接
*/
public void close() {
try {
if (null != conn && !conn.isClosed())
conn.close();
} catch(SQLException e){
e.printStackTrace();
}finally {
conn = null;
}
}
/**
* 创建数据库
*/
public void createDatabase(String dbName) {
try {
stmt = getConn().createStatement();
stmt.execute("CREATE DATABASE IF NOT EXISTS " + dbName);
System.out.println("数据库 " + dbName + " 创建成功");
} catch (Exception e) {
e.printStackTrace();
}
}
/**
* 创建用户基本信息表
*/
public void createTable1(String dbName) {
try {
stmt = getConn().createStatement();
stmt.execute("USE " + dbName);
String sql = "CREATE TABLE IF NOT EXISTS mediamatch_usermsg(" +
"terminal_no STRING," +
"phone_no STRING," +
"sm_name STRING," +
"run_name STRING," +
"sm_code STRING," +
"owner_name STRING," +
"owner_code STRING," +
"run_time STRING," +
"addressoj STRING," +
"open_time STRING," +
"force STRING)" +
"ROW FORMAT DELIMITED FIELDS TERMINATED BY ';'";
stmt.execute(sql);
System.out.println("用户基本信息表创建成功");
} catch (Exception e) {
e.printStackTrace();
}
}
/**
* 创建用户状态变更表
*/
public void createTable2(String dbName) {
try {
stmt = getConn().createStatement();
stmt.execute("USE " + dbName);
String sql = "CREATE TABLE IF NOT EXISTS mediamatch_userevent(" +
"phone_no STRING," +
"run_name STRING," +
"run_time STRING," +
"owner_name STRING," +
"owner_code STRING," +
"open_time STRING," +
"sm_name STRING)" +
"ROW FORMAT DELIMITED FIELDS TERMINATED BY ';'";
stmt.execute(sql);
System.out.println("用户状态变更表创建成功");
} catch (Exception e) {
e.printStackTrace();
}
}
/**
* 创建账单数据表
*/
public void createTable3(String dbName) {
try {
stmt = getConn().createStatement();
stmt.execute("USE " + dbName);
String sql = "CREATE TABLE IF NOT EXISTS mmconsume_billevents(" +
"terminal_no STRING," +
"phone_no STRING," +
"fee_code STRING," +
"year_month STRING," +
"owner_name STRING," +
"owner_code STRING," +
"sm_name STRING," +
"should_pay STRING," +
"favour_fee STRING)" +
"ROW FORMAT DELIMITED FIELDS TERMINATED BY ';'";
stmt.execute(sql);
System.out.println("账单数据表创建成功");
} catch (Exception e) {
e.printStackTrace();
}
}
/**
* 创建订单数据表
*/
public void createTable4(String dbName) {
try {
stmt = getConn().createStatement();
stmt.execute("USE " + dbName);
String sql = "CREATE TABLE IF NOT EXISTS order_index(" +
"phone_no STRING," +
"owner_name STRING," +
"optdate STRING," +
"prodname STRING," +
"sm_name STRING," +
"offerid STRING," +
"offername STRING," +
"business_name STRING," +
"owner_code STRING," +
"prodprcid STRING," +
"prodprcname STRING," +
"effdate STRING," +
"expdate STRING," +
"orderdate STRING," +
"cost STRING," +
"mode_time STRING," +
"prodstatus STRING," +
"run_name STRING," +
"orderno STRING," +
"offertype STRING)" +
"ROW FORMAT DELIMITED FIELDS TERMINATED BY ';'";
stmt.execute(sql);
System.out.println("订单数据表创建成功");
} catch (Exception e) {
e.printStackTrace();
}
}
/**
* 创建收视行为表
*/
public void createTable5(String dbName) {
try {
stmt = getConn().createStatement();
stmt.execute("USE " + dbName);
String sql = "CREATE TABLE IF NOT EXISTS media_index(" +
"terminal_no STRING," +
"phone_no STRING," +
"duration STRING," +
"station_name STRING," +
"origin_time STRING," +
"end_time STRING," +
"owner_code STRING," +
"owner_name STRING," +
"vod_cat_tags ARRAY<STRUCT<level1_name: STRING,level2_name: STRING,level3_name: STRING,level4_name: STRING, level5_name: STRING>>," +
"resolution STRING," +
"audio_lang STRING," +
"region STRING," +
"res_name STRING," +
"res_type STRING," +
"vod_title STRING," +
"category_name STRING," +
"program_title STRING," +
"sm_name STRING)" +
"ROW FORMAT DELIMITED FIELDS TERMINATED BY ';'";
stmt.execute(sql);
System.out.println("收视行为表创建成功");
} catch (Exception e) {
e.printStackTrace();
}
}
/**
* 加载本地CSV数据到Hive表
*/
public void loadData(String localFile, String tbName) {
String sql = "LOAD DATA LOCAL INPATH '" + localFile + "' OVERWRITE INTO TABLE " + tbName;
System.out.println("执行SQL:" + sql);
try {
stmt = getConn().createStatement();
stmt.execute(sql);
System.out.println("数据加载完成:" + localFile + " → " + tbName);
} catch (Exception e) {
e.printStackTrace();
}
}
}4. 数据存储测试类
java
运行
public class HiveTest {
public static void main(String[] args) {
// 业务参数
String dbName = "ZJSM2";
String localFile1 = "/opt/data/mediamatch_usermsg.csv";
String tbName1 = "mediamatch_usermsg";
String localFile2 = "/opt/data/mediamatch_userevent.csv";
String tbName2 = "mediamatch_userevent";
String localFile3 = "/opt/data/mmconsume_billevents.csv";
String tbName3 = "mmconsume_billevents";
String localFile4 = "/opt/data/order_index.csv";
String tbName4 = "order_index";
String localFile5 = "/opt/data/media_index.csv";
String tbName5 = "media_index";
HiveHelper helper = new HiveHelper();
try {
// 1. 创建数据库
helper.createDatabase(dbName);
// 2. 创建5张业务表
helper.createTable1(dbName);
helper.createTable2(dbName);
helper.createTable3(dbName);
helper.createTable4(dbName);
helper.createTable5(dbName);
// 3. 加载数据
helper.loadData(localFile1, tbName1);
helper.loadData(localFile2, tbName2);
helper.loadData(localFile3, tbName3);
helper.loadData(localFile4, tbName4);
helper.loadData(localFile5, tbName5);
System.out.println("=== 广电数据存储全流程执行完成 ===");
} finally {
helper.close();
}
}
}5. 教材原版代码缺陷分析(提升知识点)
表格
| 缺陷类型 | 具体问题 | 风险 |
|---|---|---|
| 资源泄漏 | stmt、rs作为成员变量,每次调用方法都会被覆盖,旧资源未关闭 | 长时间运行会耗尽数据库连接,导致服务不可用 |
| 关闭方法不完整 | close()只关闭了 Connection,未关闭 Statement 和 ResultSet | 资源泄漏,占用数据库服务器资源 |
| 代码冗余 | 每个建表方法都重复写stmt.execute("USE " + dbName) | 代码冗余,维护成本高 |
| 扩展性差 | selectAll仅支持单张表,新增表需修改方法 | 复用性低 |
6. 核心注意事项
- 路径问题:
LOAD DATA LOCAL INPATH的路径是Linux 服务器本地路径,不是 Windows 本地路径,必须提前上传 CSV 文件。 - 分隔符一致:建表语句的分隔符
;必须与 CSV 文件实际分隔符一致,否则数据全为 NULL。 - 权限问题:若报权限错误,执行
hdfs dfs -chmod -R 777 /user/hive/warehouse。
五、模块四:编写程序实现广电数据查询与处理
1. 数据查询方法
在HiveHelper中新增查询方法:
java
运行
/**
* 查询用户状态变更表所有数据
*/
public void selectAll(String tbName) {
String sql = "SELECT * FROM " + tbName;
if (tbName.equals("mediamatch_userevent")) {
try {
stmt = getConn().createStatement();
rs = stmt.executeQuery(sql);
System.out.println("=== 用户状态变更表数据 ===");
while (rs.next()) {
System.out.println(
rs.getString("phone_no") + "\t" +
rs.getString("run_name") + "\t" +
rs.getString("run_time") + "\t" +
rs.getString("owner_name") + "\t" +
rs.getString("owner_code") + "\t" +
rs.getString("open_time")
);
}
} catch (Exception e) {
e.printStackTrace();
}
}
}2. 三大数据清洗方法
通过CTAS(CREATE TABLE AS SELECT)语法,将清洗逻辑封装为方法,生成清洗后的新表。
(1)用户基本数据清洗
清洗规则:排除无效等级、无效业务类型、无效用户状态
java
运行
public void cleanTable1(String dbName) {
try {
stmt = getConn().createStatement();
stmt.execute("USE " + dbName);
String sql = "CREATE TABLE IF NOT EXISTS mediamatch_usermsg_clean " +
"AS " +
"SELECT * FROM mediamatch_usermsg " +
"WHERE " +
"owner_code NOT IN ('02','09','10') " +
"AND owner_name NOT IN ('EA级','EB级','EC级','ED级','EE级') " +
"AND sm_name IN ('数字电视','互动电视','珠江宽频','甜果电视') " +
"AND run_name IN ('正常','欠费暂停','主动暂停','主动销户')";
stmt.execute(sql);
System.out.println("用户基本数据清洗完成");
} catch (Exception e) {
e.printStackTrace();
}
}(2)收视行为数据清洗
清洗规则:时长 20 秒~5 小时;直播排除整点异常数据;点播仅过滤时长
java
运行
public void cleanTable2(String dbName) {
try {
stmt = getConn().createStatement();
stmt.execute("USE " + dbName);
String sql = "CREATE TABLE IF NOT EXISTS media_index_clean " +
"AS " +
"SELECT * FROM media_index " +
"WHERE " +
"(CAST(duration AS double)/1000 >= 20 " +
"AND CAST(duration AS double)/(1000*60*60) < 5 " +
"AND res_type='0' " +
"AND origin_time NOT LIKE '%00' " +
"AND end_time NOT LIKE '%00') " +
"OR " +
"(CAST(duration AS double)/1000 >= 20 " +
"AND CAST(duration AS double)/(1000*60*60) < 5 " +
"AND res_type='1')";
stmt.execute(sql);
System.out.println("收视行为数据清洗完成");
} catch (Exception e) {
e.printStackTrace();
}
}(3)账单数据清洗
清洗规则:排除应付金额为负数的异常数据
java
运行
public void cleanTable3(String dbName) {
try {
stmt = getConn().createStatement();
stmt.execute("USE " + dbName);
String sql = "CREATE TABLE IF NOT EXISTS mmconsume_billevents_clean " +
"AS " +
"SELECT * FROM mmconsume_billevents " +
"WHERE should_pay >= 0";
stmt.execute(sql);
System.out.println("账单数据清洗完成");
} catch (Exception e) {
e.printStackTrace();
}
}3. 全流程测试执行
在HiveTest的 main 方法中追加:
java
运行
// 4. 查询用户状态变更表数据
helper.selectAll(tbName2);
// 5. 执行数据清洗
helper.cleanTable1(dbName);
helper.cleanTable2(dbName);
helper.cleanTable3(dbName);
System.out.println("=== 广电数据查询与清洗全流程执行完成 ===");4. 结果验证
在 Hive CLI 中执行以下命令验证清洗结果:
sql
-- 查看清洗后的数据行数
SELECT COUNT(*) FROM mediamatch_usermsg_clean;
SELECT COUNT(*) FROM media_index_clean;
SELECT COUNT(*) FROM mmconsume_billevents_clean;
-- 查看前5条数据
SELECT * FROM mediamatch_usermsg_clean LIMIT 5;六、IDEA 程序调试方法
1. 调试基本步骤
- 设置断点:点击代码行号右侧空白处,出现红色圆点即为断点,程序执行到该行会自动暂停。
- 启动调试:右键代码 → 「Debug」,或快捷键
Shift+F9。 - 控制执行:使用快捷键控制程序运行,观察变量变化。
- 查看变量:底部 Debugger 窗口的 Variables 面板,可查看所有变量的实时值。
2. 核心调试快捷键
表格
| 快捷键 | 功能 | 使用场景 |
|---|---|---|
F8 | 步过:执行当前行,跳到下一行,不进入方法内部 | 逐行执行,查看每一步执行结果 |
F7 | 步入:进入当前调用的方法内部 | 查看工具类方法内部逻辑,定位方法内报错 |
F9 | 恢复执行:继续运行,直到下一个断点 | 跳过无关代码,快速定位问题 |
Alt+F8 | 计算表达式:查看变量的实时值 | 查看 SQL 完整内容,验证 SQL 正确性 |
3. 典型调试场景
- SQL 执行报错:在
stmt.execute(sql)行设断点,用Alt+F8查看 sql 变量,复制到 Hive CLI 执行验证。 - 查询无结果:断点查看 SQL 语句,检查表名、数据库名、查询条件是否正确。
- 空指针异常:断点查看哪个变量为 null,追溯变量的赋值逻辑。
七、进阶优化与拓展
1. 代码质量优化
(1)资源管理优化
将Statement、ResultSet改为方法内局部变量,每个方法执行完立即关闭,彻底解决资源泄漏。
java
运行
/**
* 通用资源关闭方法
*/
public void close(ResultSet rs, Statement stmt) {
try {
if (rs != null && !rs.isClosed()) rs.close();
if (stmt != null && !stmt.isClosed()) stmt.close();
} catch (SQLException e) {
e.printStackTrace();
}
}(2)日志配置优化
在src/main/resources下添加log4j.properties,解决日志警告、减少冗余输出:
properties
log4j.rootLogger=WARN, console
log4j.appender.console=org.apache.log4j.ConsoleAppender
log4j.appender.console.layout=org.apache.log4j.PatternLayout
log4j.appender.console.layout.ConversionPattern=%d{yyyy-MM-dd HH:mm:ss} [%t] %-5p %c{1} - %m%n
log4j.logger.org.apache.hive=WARN
log4j.logger.org.apache.hadoop=WARN
log4j.logger.org.eclipse.jetty=WARN(3)中文乱码优化
在连接 URL 中添加编码参数,避免中文查询乱码:
java
运行
String url = "jdbc:hive2://master:10000/default?useUnicode=true&characterEncoding=UTF-8";2. 企业级拓展方向
- 连接池优化:引入 HikariCP、Druid 等连接池管理连接,避免频繁创建销毁连接,提升并发性能。
- 配置外置化:将连接地址、账号密码放到 properties 文件中,不硬编码在 Java 代码里,方便多环境切换。
- SQL 注入防护:使用
PreparedStatement预编译 SQL,替代字符串拼接,防止 SQL 注入风险。 - 定时任务:结合 Quartz、Spring Task 实现定时自动执行数据清洗任务。
八、全章常见问题汇总
1. 服务端启动类
- 启动 metastore 报错找不到目录 → 手动创建
system:java.io.tmpdir指定的目录 - 两个服务启动后只有一个 RunJar → 启动顺序错误,先启动的服务被异常终止,按顺序重新启动
- 10000 端口一直不被监听 → 服务初始化慢,等待 1~2 分钟再查看
2. 客户端连接类
- 代理用户权限报错 → 检查 3 个 XML 是否放入 resources 目录,检查配置内容是否正确
- 无法解析 master 主机名 → 本地 Windows 未配置 hosts 映射
- 依赖下载慢 → 配置阿里云 Maven 镜像
3. 代码运行类
- 建表成功但找不到表 → 没有执行
USE 数据库名,建表到了 default 库 - LOAD DATA 成功但表无数据 → 文件路径错误、分隔符不一致、文件权限不足
- 编译报错类名冲突 → 修改 Connection 为 ConnectionTest
九、课后自测与作业
1. 基础简答题
- 简述 HiveServer2 相比 HiveServer 的优势。
- 简述启动 Hive 远程服务的顺序,并说明原因。
- 为什么要将
core-site.xml等配置文件放到src/main/resources目录? - 简述 JDBC 中 Statement 和 PreparedStatement 的区别。
2. 实操题
- 独立完成 Hive 远程服务配置与启动,并验证服务可用性。
- 在 HiveHelper 中新增方法
countUserByLevel,统计用户基本表中每个用户等级的用户数量,并打印结果。 - 封装方法实现直播频道观看用户数 Top10 统计,将第 6 章的 Top10 逻辑改为 Java 程序实现。
3. 拓展思考题
- 思考如何实现程序定时自动执行数据清洗任务。
- 调研 Hive JDBC 连接池的配置方法,尝试引入 HikariCP 优化连接管理。