第 8 章 广电用户数据存储与处理的程序开发
开篇自学说明
本章是 Hive 从「命令行手工操作」到「企业级程序自动化处理」的核心跨越。前面 7 章我们都是在 Hive CLI 中手动输入 HQL 语句完成建表、查询、清洗等操作,效率低、无法复用、难以集成到业务系统中。而本章的核心目标,是通过 Java 语言基于 JDBC 接口,实现 Hive 的远程程序调用,把之前所有的手工操作封装成可复用、可自动化执行的程序,完成广电数据的全流程自动化存储、查询与处理。
学完本章,你将掌握 4 大核心能力:
- Hive 远程服务 HiveServer2 的配置、启动与运维
- 基于 IDEA 搭建 Hive Java 开发环境,完成项目依赖配置
- 通过 JDBC 接口实现 Hive 的程序调用,完成数据库创建、表结构创建、数据批量加载
- 用 Java 程序封装第 4~7 章的所有数据查询、数据清洗逻辑,实现广电数据的自动化处理
前置必备准备
- 已完成 Hadoop 3.3.6、Hive 3.1.2 集群的搭建与正常运行
- 已完成广电 5 张核心业务表的创建与原始数据导入
- 本地 Windows 环境已安装 IntelliJ IDEA 2021.3+、JDK 1.8
- 本地 Windows 能正常访问 Hadoop 集群的 master 节点,已配置 hosts 主机名映射
模块一:Hive 远程服务配置与启动(任务 8.1 核心)
要实现 Java 程序远程操作 Hive,必须先开启 Hive 的远程服务。Hive 提供了 HiveServer2 服务,基于 Thrift 协议,支持多客户端并发访问,是 Java 程序通过 JDBC 连接 Hive 的核心基础。
一、核心服务说明
| 服务名称 | 核心作用 | 启动优先级 |
|---|---|---|
metastore | Hive 元数据服务,管理 Hive 表的元数据信息(库名、表名、字段、存储位置等),所有对 Hive 的操作都必须先访问元数据 | 必须先启动,再启动 HiveServer2 |
hiveserver2 | Hive 远程服务,开放 10000 端口,对外提供 JDBC 连接入口,支持 Java/Python 等客户端远程提交 HQL 执行 | 依赖 metastore 服务,后启动 |
二、完整配置步骤
步骤 1:配置 Hadoop 代理用户权限
HiveServer2 需要通过 Hadoop 的代理用户权限访问 HDFS,必须先修改 Hadoop 的核心配置文件,允许 root 用户代理所有主机、所有用户组的请求。
<!-- 代码8-1 进入master节点,修改core-site.xml配置文件 -->
<!-- 执行命令:vim /usr/local/hadoop-3.3.6/etc/hadoop/core-site.xml -->
<!-- 在configuration标签内添加以下配置 -->
<property>
<name>hadoop.proxyuser.root.hosts</name>
<value>*</value>
</property>
<property>
<name>hadoop.proxyuser.root.groups</name>
<value>*</value>
</property>hadoop.proxyuser.root.hosts:*表示允许 root 用户从任意主机访问hadoop.proxyuser.root.groups:*表示允许 root 用户代理任意用户组
步骤 2:分发配置文件到所有子节点
修改后的配置必须同步到集群所有节点,否则配置不生效。
# 代码8-2 分发core-site.xml到slave1、slave2、slave3子节点
scp -r /usr/local/hadoop-3.3.6/etc/hadoop/core-site.xml slave1:/usr/local/hadoop-3.3.6/etc/hadoop/
scp -r /usr/local/hadoop-3.3.6/etc/hadoop/core-site.xml slave2:/usr/local/hadoop-3.3.6/etc/hadoop/
scp -r /usr/local/hadoop-3.3.6/etc/hadoop/core-site.xml slave3:/usr/local/hadoop-3.3.6/etc/hadoop/步骤 3:重启 Hadoop 集群,让配置生效
# 进入Hadoop的sbin目录,停止集群
cd /usr/local/hadoop-3.3.6/sbin
./stop-all.sh
# 重新启动Hadoop集群
./start-all.sh
# 启动历史服务
./mr-jobhistory-daemon.sh start historyserver步骤 4:启动 Hive 远程服务
必须严格按照「先启 metastore,再启 hiveserver2」的顺序启动,否则会出现连接失败。
# 代码8-3 进入Hive的bin目录
cd /usr/local/hive-3.1.2/bin
# 1. 后台启动Hive元数据服务metastore
hive --service metastore &
# 执行后按回车,回到命令行
# 2. 查看进程,确认metastore启动成功(会出现RunJar进程)
jps
# 3. 后台启动HiveServer2远程服务,输出日志到nohup.out
nohup hive --service hiveserver2 &
# 执行后按回车,回到命令行
# 4. 再次查看进程,确认两个服务都启动成功
jps✅ 启动成功验证:执行jps命令后,能看到 2 个 RunJar 进程,分别对应 metastore 和 hiveserver2 服务。
步骤 5:验证 HiveServer2 服务是否正常
# 1. 查看10000端口是否被监听(HiveServer2默认端口)
netstat -nltp | grep 10000
# 2. 用beeline客户端测试本地连接,验证服务可用
beeline
# 进入beeline后,执行连接命令
!connect jdbc:hive2://master:10000/default
# 输入用户名root,密码123456,能正常进入Hive命令行,说明服务正常三、新手高频踩坑避坑指南
❌ 坑 1:先启动 hiveserver2,再启动 metastore,导致连接失败
✅ 解决:必须严格按照「metastore → hiveserver2」的顺序启动,两个服务都启动后,等待 30 秒再测试连接,服务启动需要时间
❌ 坑 2:修改 core-site.xml 后没有分发到子节点,也没有重启 Hadoop,配置不生效
✅ 解决:修改配置后必须分发到所有节点,重启 Hadoop 集群,再启动 Hive 服务
❌ 坑 3:防火墙没有关闭,10000 端口被拦截,Windows 无法远程连接
✅ 解决:关闭集群所有节点的防火墙,或开放 10000 端口
bash# 临时关闭防火墙 systemctl stop firewalld # 永久关闭防火墙 systemctl disable firewalld❌ 坑 4:本地 Windows 没有配置 master 的 hosts 映射,无法识别主机名
✅ 解决:修改 Windows 的
C:\Windows\System32\drivers\etc\hosts文件,添加 master 节点的 IP 和主机名映射,例如:
192.168.128.130 master
模块二:IDEA 开发环境搭建(任务 8.2 核心)
本模块将在本地 Windows 的 IDEA 中,搭建 Hive Java 开发环境,完成 Maven 项目创建、依赖配置、JDBC 连接测试,实现 Windows 远程连接 Linux 集群的 Hive 服务。
一、创建 Maven 项目
- 打开 IDEA,点击欢迎界面的New Project
- 左侧选择Maven,Project SDK 选择 1.8 版本的 JDK,点击 Next
- 项目名称填写
HiveJavaAPI,选择项目存储路径(如D:\HiveJavaAPI),点击 Finish - 项目创建完成后,自动生成 Maven 项目结构,核心文件为
pom.xml(Maven 依赖配置文件)
二、配置 Maven 依赖
打开项目的pom.xml文件,添加 Hadoop 和 Hive 相关依赖,解决 jar 包冲突问题。
<?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>
<name>HiveJavaAPI</name>
<description>Java 调用 HiveServer2 实训项目</description>
<!-- 全局统一配置:JDK、编码 -->
<properties>
<maven.compiler.source>8</maven.compiler.source>
<maven.compiler.target>8</maven.compiler.target>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding>
</properties>
<!-- 依赖仓库:国内阿里云镜像,加速Jar包下载 -->
<repositories>
<repository>
<id>aliyun-central</id>
<name>阿里云公共仓库</name>
<url>https://maven.aliyun.com/repository/public</url>
<releases>
<enabled>true</enabled>
</releases>
<snapshots>
<enabled>true</enabled>
</snapshots>
</repository>
</repositories>
<!-- 核心依赖 -->
<dependencies>
<!-- 1. Hadoop 公共基础依赖 Hadoop3.3.6 -->
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-common</artifactId>
<version>3.3.6</version>
</dependency>
<!-- 2. Hive 执行核心依赖 Hive3.1.2 -->
<dependency>
<groupId>org.apache.hive</groupId>
<artifactId>hive-exec</artifactId>
<version>3.1.2</version>
</dependency>
<!-- 3. Hive JDBC 驱动(连接HiveServer2核心,排除冲突包) -->
<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>
<!-- 4. 日志依赖(排查报错必备,运行时生效) -->
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-simple</artifactId>
<version>1.7.25</version>
<scope>runtime</scope>
</dependency>
</dependencies>
</project>配置完成后,右键 pom.xml → 选择 Maven → Reload project,等待依赖自动下载完成。下载完成后,可在左侧 External Libraries 中看到所有依赖包。
三、JDBC 核心接口详解
JDBC(Java DataBase Connectivity)是 Java 访问数据库的统一 API,Hive 提供了专属的 Hive-JDBC 驱动,我们通过这套 API 实现对 Hive 的远程操作。核心接口如下,新手必须掌握:
| 接口名称 | 核心作用 | 常用方法 |
|---|---|---|
Driver | JDBC 驱动接口,Hive 的实现类为org.apache.hive.jdbc.HiveDriver | Class.forName(driverName)加载驱动 |
DriverManager | 驱动管理类,用于创建数据库连接 | getConnection(url, username, password)获取连接对象 |
Connection | 数据库连接对象,代表 Java 程序和 Hive 的连接会话 | createStatement()创建 Statement 对象、close()关闭连接 |
Statement | SQL 语句执行对象,用于向 Hive 提交静态 HQL 语句 | execute()执行 DDL 语句(建库、建表)、executeQuery()执行查询语句、executeUpdate()执行数据加载语句 |
PreparedStatement | 预编译 SQL 语句对象,Statement 的子接口,支持占位符?,防止 SQL 注入 | setString()给占位符赋值、executeQuery()执行查询 |
ResultSet | 查询结果集对象,封装了 SELECT 查询的返回结果 | next()移动游标到下一行、getString()获取字符串类型字段值、getInt()获取整数类型字段值 |
四、Hive 连接测试程序编写
我们编写第一个 Java 程序,实现远程连接 Hive,并创建一个测试数据库,验证开发环境是否搭建成功。
步骤 1:创建 Java 类
在 IDEA 中,找到src/main/java目录,右键 → New → Java Class,类名填写ConnectionTest,按回车创建。
步骤 2:编写连接测试代码
// 代码8-5 完整的连接测试代码
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.Statement;
public class ConnectionTest {
public static void main(String[] args) {
// 1. 定义连接核心参数
// Hive JDBC驱动类名,固定写法
String driverName = "org.apache.hive.jdbc.HiveDriver";
// Hive连接地址:jdbc:hive2://master主机名:10000/要连接的数据库名
String url = "jdbc:hive2://master:10000/default";
// Hive连接用户名,默认root
String username = "root";
// Hive连接密码,默认123456,无密码填空字符串
String password = "123456";
Connection connection = null;
Statement statement = null;
try {
// 2. 加载JDBC驱动
Class.forName(driverName);
System.out.println("Hive驱动加载成功");
// 3. 通过DriverManager获取数据库连接
connection = DriverManager.getConnection(url, username, password);
System.out.println("Hive连接成功:" + connection);
// 4. 创建Statement对象,用于执行HQL语句
statement = connection.createStatement();
// 5. 执行HQL语句,创建测试数据库test
String createDbSql = "CREATE DATABASE IF NOT EXISTS test";
statement.execute(createDbSql);
System.out.println("测试数据库test创建成功");
} catch (Exception e) {
// 捕获异常,打印错误信息
e.printStackTrace();
} finally {
// 6. 关闭资源,释放连接,必须按「statement → connection」的顺序关闭
try {
if (statement != null) statement.close();
if (connection != null) connection.close();
System.out.println("连接已关闭");
} catch (Exception e) {
e.printStackTrace();
}
}
}
}步骤 3:运行程序并验证
- 右键代码空白处 → 选择 Run 'ConnectionTest.main ()',运行程序
- 控制台输出
Process finished with exit code 0,且打印「驱动加载成功、连接成功、数据库创建成功」,说明程序运行正常 - 回到 Linux 的 Hive CLI,执行
SHOW DATABASES;,能看到新建的 test 数据库,说明远程连接完全正常
五、开发环境搭建避坑指南
❌ 坑 1:pom.xml 依赖下载失败,出现红色报错
✅ 解决:配置 Maven 国内镜像源(阿里云镜像),检查网络是否正常,重新 Reload project
❌ 坑 2:运行报错
ClassNotFoundException: org.apache.hive.jdbc.HiveDriver✅ 解决:检查 pom.xml 的 hive-jdbc 依赖是否正确添加,是否完成了 Maven Reload,依赖是否下载成功
❌ 坑 3:运行报错
Connection refused: connect✅ 解决:
- 检查 HiveServer2 服务是否正常启动,10000 端口是否被监听
- 检查 Windows 能否 ping 通 master 节点的 IP
- 检查防火墙是否关闭,10000 端口是否开放
- 检查 url 中的主机名和端口是否正确
❌ 坑 4:运行报错
HiveSQLException: Failed to open new session✅ 解决:检查 core-site.xml 的代理用户配置是否正确,是否分发到所有节点,是否重启了 Hadoop 集群
模块三:编写程序实现广电数据的存储(任务 8.3 核心)
本模块将封装一个 Hive 操作工具类HiveHelper,把 Hive 的连接、关闭、建库、建表、数据加载等常用操作封装成可复用的方法,实现广电 5 张核心业务表的自动化创建与数据装载。
一、封装 HiveHelper 工具类
在src/main/java目录下创建 Java 类HiveHelper,封装所有 Hive 操作的核心方法。
// 代码8-6 HiveHelper工具类完整代码
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.ResultSet;
import java.sql.Statement;
/**
* Hive操作工具类,封装Hive的连接、建库、建表、数据加载、查询等操作
*/
public class HiveHelper {
// 固定配置参数,统一管理,避免重复代码
private static final String DRIVER_NAME = "org.apache.hive.jdbc.HiveDriver";
private static final String URL = "jdbc:hive2://master:10000/default";
private static final String USERNAME = "root";
private static final String PASSWORD = "123456";
// 连接对象、SQL执行对象、查询结果集
private Connection conn = null;
private Statement stmt = null;
private ResultSet rs = null;
/**
* 获取Hive数据库连接
* @return Connection连接对象
* @throws Exception 驱动加载、连接失败异常
*/
public Connection getConn() throws Exception {
// 单例模式,避免重复创建连接
if (null == conn || conn.isClosed()) {
// 加载驱动
Class.forName(DRIVER_NAME);
// 获取连接
conn = DriverManager.getConnection(URL, USERNAME, PASSWORD);
}
return conn;
}
/**
* 关闭所有资源,释放连接
*/
public void close() {
try {
if (rs != null && !rs.isClosed()) rs.close();
if (stmt != null && !stmt.isClosed()) stmt.close();
if (conn != null && !conn.isClosed()) conn.close();
} catch (Exception e) {
e.printStackTrace();
} finally {
rs = null;
stmt = null;
conn = null;
}
}
/**
* 创建数据库
* @param dbName 数据库名称
*/
public void createDatabase(String dbName) {
try {
stmt = getConn().createStatement();
String sql = "CREATE DATABASE IF NOT EXISTS " + dbName;
stmt.execute(sql);
System.out.println("数据库" + dbName + "创建成功");
} catch (Exception e) {
e.printStackTrace();
}
}
/**
* 切换到指定数据库
* @param dbName 数据库名称
*/
public void useDatabase(String dbName) {
try {
stmt = getConn().createStatement();
stmt.execute("USE " + dbName);
System.out.println("切换到数据库:" + dbName);
} catch (Exception e) {
e.printStackTrace();
}
}
/**
* 创建用户状态变更数据表mediamatch_userevent
* @param dbName 数据库名称
*/
public void createUserEventTable(String dbName) {
try {
// 先切换到目标数据库
useDatabase(dbName);
stmt = getConn().createStatement();
// 建表HQL,和第3章CLI中的建表语句完全一致
String createTableSql = "CREATE TABLE IF NOT EXISTS mediamatch_userevent(" +
"phone_no STRING COMMENT '用户编号'," +
"run_name STRING COMMENT '用户状态'," +
"run_time STRING COMMENT '状态变更时间'," +
"owner_name STRING COMMENT '用户等级名称'," +
"owner_code STRING COMMENT '用户等级编号'," +
"open_time STRING COMMENT '开户时间')" +
"ROW FORMAT DELIMITED FIELDS TERMINATED BY ';'";
// 执行建表语句
stmt.execute(createTableSql);
System.out.println("用户状态变更表mediamatch_userevent创建成功");
} catch (Exception e) {
e.printStackTrace();
}
}
/**
* 加载本地CSV数据到Hive表
* @param localFile Linux服务器上的CSV文件绝对路径
* @param tbName 目标表名
*/
public void loadData(String localFile, String tbName) {
try {
stmt = getConn().createStatement();
// LOAD DATA语句,OVERWRITE表示覆盖表中原有数据
String loadSql = "LOAD DATA LOCAL INPATH '" + localFile + "' OVERWRITE INTO TABLE " + tbName;
System.out.println("执行数据加载SQL:" + loadSql);
stmt.execute(loadSql);
System.out.println("数据加载完成,文件:" + localFile + " → 表:" + tbName);
} catch (Exception e) {
e.printStackTrace();
}
}
}✅ 工具类说明:
- 把连接参数、重复代码统一封装,后续新增表、新增操作只需添加对应的方法即可
- 所有方法都做了异常处理,方便定位问题
- 建表语句和第 3 章 CLI 中的语句完全一致,保证表结构统一
- 其他 4 张表(用户基本表、账单表、订单表、收视行为表)的创建方法,只需参照
createUserEventTable方法,修改建表 HQL 即可
二、编写测试类,实现广电数据自动化存储
在src/main/java目录下创建测试类HiveDataStorageTest,调用 HiveHelper 工具类,完成广电数据库创建、表创建、数据加载的全流程自动化执行。
// 代码8-7 数据存储测试类完整代码
public class HiveDataStorageTest {
public static void main(String[] args) {
// 1. 定义核心参数
String dbName = "ZJSM_GUANGDIAN"; // 广电业务数据库名
// Linux服务器上的CSV数据文件路径,必须提前把数据文件上传到该路径
String userEventDataPath = "/opt/data/mediamatch_userevent.csv";
String userMsgDataPath = "/opt/data/mediamatch_usermsg.csv";
String billDataPath = "/opt/data/mmconsume_billevents.csv";
String orderDataPath = "/opt/data/order_index.csv";
String mediaDataPath = "/opt/data/media_index.csv";
// 2. 创建HiveHelper工具类对象
HiveHelper helper = new HiveHelper();
try {
// 3. 执行全流程操作
// 3.1 创建广电业务数据库
helper.createDatabase(dbName);
// 3.2 创建5张核心业务表
helper.createUserEventTable(dbName);
// 此处可补充其他4张表的创建方法调用
// helper.createUserMsgTable(dbName);
// helper.createBillTable(dbName);
// helper.createOrderTable(dbName);
// helper.createMediaTable(dbName);
// 3.3 加载CSV数据到对应表中
helper.loadData(userEventDataPath, "mediamatch_userevent");
// 此处可补充其他表的数据加载
// helper.loadData(userMsgDataPath, "mediamatch_usermsg");
} catch (Exception e) {
e.printStackTrace();
} finally {
// 4. 关闭连接,释放资源
helper.close();
System.out.println("=== 广电数据存储全流程执行完成 ===");
}
}
}三、核心避坑指南
❌ 坑 1:LOAD DATA 语句执行成功,但表中无数据
✅ 解决:
LOCAL INPATH后面的路径是Linux 服务器上的文件路径,不是 Windows 本地的路径,必须提前把 CSV 文件上传到 Linux 服务器的对应目录- 检查 Linux 上的文件路径是否正确,文件是否存在,权限是否正确
- 检查建表语句的字段分隔符是否和 CSV 文件的分隔符一致
❌ 坑 2:建表语句执行报错,提示语法错误
✅ 解决:检查 Java 字符串中的建表 SQL 是否有拼写错误,字符串拼接是否正确,关键字是否正确,字段类型是否符合 Hive 规范
❌ 坑 3:程序执行报错
Permission denied: user=anonymous, access=WRITE✅ 解决:检查 HDFS 的目录权限,执行
hdfs dfs -chmod -R 777 /user/hive/warehouse,给 Hive 仓库目录开放读写权限
模块四:编写程序实现广电数据的查询与处理(任务 8.4 核心)
本模块将在 HiveHelper 工具类中扩展数据查询、数据清洗的方法,把第 4~7 章的 HQL 查询、数据清洗逻辑封装成 Java 程序,实现广电数据的自动化查询与处理。
一、扩展 HiveHelper 工具类,新增数据查询方法
在 HiveHelper 类中新增selectAllUserEvent方法,实现用户状态变更表的数据查询,并打印到控制台。
// 代码8-11 数据查询方法,添加到HiveHelper类中
/**
* 查询用户状态变更表的所有数据,并打印到控制台
* @param tbName 表名
*/
public void selectAllUserEvent(String tbName) {
String sql = "SELECT * FROM " + tbName;
try {
stmt = getConn().createStatement();
// 执行查询语句,返回结果集ResultSet
rs = stmt.executeQuery(sql);
System.out.println("=== 用户状态变更表数据 ===");
System.out.println("用户编号\t用户状态\t状态变更时间\t用户等级\t等级编号\t开户时间");
// 遍历结果集,next()方法移动游标,有下一行数据返回true,无数据返回false
while (rs.next()) {
// 通过字段名获取对应的值,拼接成字符串打印
String line = 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");
System.out.println(line);
}
} catch (Exception e) {
e.printStackTrace();
}
}二、扩展数据清洗方法,实现自动化数据清洗
把第 7 章的 3 大核心数据清洗逻辑,封装成 Java 方法,实现无效数据的自动化清洗。
1. 用户基本数据清洗方法
// 代码8-12 用户基本数据清洗方法,添加到HiveHelper类中
/**
* 清洗用户基本表的无效数据,创建清洗后的新表mediamatch_usermsg_clean
* @param dbName 数据库名称
*/
public void cleanUserMsgData(String dbName) {
try {
useDatabase(dbName);
stmt = getConn().createStatement();
// 清洗逻辑和第7章CLI中的HQL完全一致
String cleanSql = "CREATE TABLE IF NOT EXISTS mediamatch_usermsg_clean " +
"AS " +
"SELECT * FROM mediamatch_usermsg " +
"WHERE " +
"owner_code NOT IN ('2','9','10') " +
"AND " +
"owner_name NOT IN ('EA级','EB级','EC级','ED级','EE级') " +
"AND " +
"sm_name IN ('数字电视','互动电视','珠江宽频','甜果电视') " +
"AND " +
"run_name IN ('正常','欠费暂停','主动暂停','主动销户')";
stmt.execute(cleanSql);
System.out.println("用户基本数据清洗完成,清洗表mediamatch_usermsg_clean创建成功");
} catch (Exception e) {
e.printStackTrace();
}
}2. 收视行为数据清洗方法
// 代码8-13 收视行为数据清洗方法,添加到HiveHelper类中
/**
* 清洗收视行为表的无效数据,创建清洗后的新表media_index_clean
* @param dbName 数据库名称
*/
public void cleanMediaData(String dbName) {
try {
useDatabase(dbName);
stmt = getConn().createStatement();
String cleanSql = "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(cleanSql);
System.out.println("收视行为数据清洗完成,清洗表media_index_clean创建成功");
} catch (Exception e) {
e.printStackTrace();
}
}3. 账单数据清洗方法
// 代码8-14 账单数据清洗方法,添加到HiveHelper类中
/**
* 清洗账单表的无效数据,创建清洗后的新表mmconsume_billevents_clean
* @param dbName 数据库名称
*/
public void cleanBillData(String dbName) {
try {
useDatabase(dbName);
stmt = getConn().createStatement();
String cleanSql = "CREATE TABLE IF NOT EXISTS mmconsume_billevents_clean " +
"AS " +
"SELECT * FROM mmconsume_billevents " +
"WHERE should_pay >= 0";
stmt.execute(cleanSql);
System.out.println("账单数据清洗完成,清洗表mmconsume_billevents_clean创建成功");
} catch (Exception e) {
e.printStackTrace();
}
}三、编写测试类,实现数据查询与清洗的自动化执行
创建测试类HiveDataProcessTest,调用 HiveHelper 的方法,实现数据查询与清洗的全流程执行。
public class HiveDataProcessTest {
public static void main(String[] args) {
String dbName = "ZJSM_GUANGDIAN";
HiveHelper helper = new HiveHelper();
try {
// 1. 查询用户状态变更表数据
helper.selectAllUserEvent("mediamatch_userevent");
// 2. 执行全量数据清洗
helper.cleanUserMsgData(dbName);
helper.cleanMediaData(dbName);
helper.cleanBillData(dbName);
System.out.println("=== 广电数据查询与清洗全流程执行完成 ===");
} catch (Exception e) {
e.printStackTrace();
} finally {
// 关闭连接
helper.close();
}
}
}运行程序后,控制台会打印查询结果,同时在 Hive 中自动创建 3 张清洗后的表,回到 Hive CLI 执行SELECT * FROM mediamatch_usermsg_clean LIMIT 5;即可验证清洗结果。
模块五:IDEA 程序调试方法与常见问题排查
一、IDEA 程序调试核心步骤
当程序运行结果不符合预期、出现报错时,通过断点调试可以快速定位问题。
设置断点:在代码行号的右侧空白处单击,出现红色圆点,即为断点,程序执行到这一行会自动暂停
进入调试模式:右键代码空白处 → 选择 Debug ' 类名.main ()',启动调试模式
调试核心快捷键
快捷键 功能 新手常用场景 F8 步过,执行当前行,跳到下一行,不进入方法内部 逐行执行代码,查看每一步的执行结果 F7 步入,进入当前调用的方法内部 查看工具类方法内部的执行情况,定位方法内的报错 F9 恢复程序执行,跳到下一个断点 快速跳转到下一个断点,跳过不需要逐行查看的代码 Alt+F8 计算表达式,查看变量的实时值 查看 SQL 语句、变量的实时内容,定位 SQL 拼写错误 查看变量:调试模式下,底部 Debugger 窗口的 Variables 面板,可以看到所有变量的实时值,比如 SQL 语句的内容、连接对象是否正常、参数是否正确。
二、常见问题排查指南
| 问题现象 | 常见原因 | 解决方法 |
|---|---|---|
| 程序执行建表语句成功,但 Hive 中找不到表 | 程序连接的数据库和 Hive CLI 的数据库不一致 | 检查程序中是否执行了USE 数据库名,确认建表的数据库是否正确 |
| 数据清洗方法执行成功,但清洗表中无数据 | 清洗条件错误,或原始表中无符合条件的数据 | 调试模式下查看 cleanSql 的完整内容,复制到 Hive CLI 中执行,查看是否有结果返回 |
程序运行报错SQLFeatureNotSupportedException | 使用了 Statement 的不支持方法,比如executeUpdate执行查询语句 | DDL 语句、LOAD DATA 语句用execute(),SELECT 查询语句用executeQuery() |
程序运行一段时间后报错Connection is closed | 连接被提前关闭,或长时间无操作被释放 | 检查 close () 方法的调用位置,确保所有 SQL 执行完成后再关闭连接 |
| 查询结果打印乱码 | 程序编码和 Hive 的编码不一致 | IDEA 中设置项目编码为 UTF-8,修改 pom.xml 添加<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding> |
模块六:入门自测题(学完检验成果)
题目 1
请在 HiveHelper 工具类中,新增一个方法countUserByLevel,实现统计用户基本表中每个用户等级的用户数量,并打印到控制台。
参考答案
/**
* 统计每个用户等级的用户数量
* @param dbName 数据库名称
*/
public void countUserByLevel(String dbName) {
try {
useDatabase(dbName);
stmt = getConn().createStatement();
String sql = "SELECT owner_name, COUNT(DISTINCT phone_no) AS user_num FROM mediamatch_usermsg GROUP BY owner_name";
rs = stmt.executeQuery(sql);
System.out.println("=== 用户等级分布统计 ===");
System.out.println("用户等级\t用户数量");
while (rs.next()) {
System.out.println(rs.getString("owner_name") + "\t" + rs.getInt("user_num"));
}
} catch (Exception e) {
e.printStackTrace();
}
}题目 2
请编写 Java 程序,实现统计直播频道 Top10 的功能,把第 6 章的 Top10 统计逻辑封装成程序,打印出观看用户数最高的 10 个直播频道。
参考答案
/**
* 统计直播频道观看用户数Top10
* @param dbName 数据库名称
*/
public void getLiveChannelTop10(String dbName) {
try {
useDatabase(dbName);
stmt = getConn().createStatement();
String sql = "SELECT station_name, COUNT(DISTINCT phone_no) AS user_num " +
"FROM media_index_clean " +
"WHERE res_type = '0' " +
"GROUP BY station_name " +
"ORDER BY user_num DESC " +
"LIMIT 10";
rs = stmt.executeQuery(sql);
System.out.println("=== 直播频道观看用户数Top10 ===");
System.out.println("频道名称\t观看用户数");
while (rs.next()) {
System.out.println(rs.getString("station_name") + "\t" + rs.getInt("user_num"));
}
} catch (Exception e) {
e.printStackTrace();
}
}
// 测试类中调用
public static void main(String[] args) {
String dbName = "ZJSM_GUANGDIAN";
HiveHelper helper = new HiveHelper();
try {
helper.getLiveChannelTop10(dbName);
} catch (Exception e) {
e.printStackTrace();
} finally {
helper.close();
}
}题目 3
请编写程序,实现将清洗后的用户基本表数据,导出到 HDFS 的/opt/guangdian_clean/user_data目录,字段用逗号分隔。
参考答案
/**
* 导出清洗后的用户数据到HDFS
* @param dbName 数据库名称
* @param hdfsPath 导出的HDFS路径
*/
public void exportCleanUserData(String dbName, String hdfsPath) {
try {
useDatabase(dbName);
stmt = getConn().createStatement();
String exportSql = "INSERT OVERWRITE DIRECTORY '" + hdfsPath + "' " +
"ROW FORMAT DELIMITED FIELDS TERMINATED BY ',' " +
"SELECT * FROM mediamatch_usermsg_clean";
stmt.execute(exportSql);
System.out.println("用户数据导出完成,导出路径:" + hdfsPath);
} catch (Exception e) {
e.printStackTrace();
}
}
// 测试调用
helper.exportCleanUserData(dbName, "/opt/guangdian_clean/user_data");