Skip to content

第 8 章 广电用户数据存储与处理的程序开发 ​

开篇自学说明 ​

本章是 Hive 从「命令行手工操作」到「企业级程序自动化处理」的核心跨越。前面 7 章我们都是在 Hive CLI 中手动输入 HQL 语句完成建表、查询、清洗等操作,效率低、无法复用、难以集成到业务系统中。而本章的核心目标,是通过 Java 语言基于 JDBC 接口,实现 Hive 的远程程序调用,把之前所有的手工操作封装成可复用、可自动化执行的程序,完成广电数据的全流程自动化存储、查询与处理。

学完本章,你将掌握 4 大核心能力:

  1. Hive 远程服务 HiveServer2 的配置、启动与运维
  2. 基于 IDEA 搭建 Hive Java 开发环境,完成项目依赖配置
  3. 通过 JDBC 接口实现 Hive 的程序调用,完成数据库创建、表结构创建、数据批量加载
  4. 用 Java 程序封装第 4~7 章的所有数据查询、数据清洗逻辑,实现广电数据的自动化处理

前置必备准备

  1. 已完成 Hadoop 3.3.6、Hive 3.1.2 集群的搭建与正常运行
  2. 已完成广电 5 张核心业务表的创建与原始数据导入
  3. 本地 Windows 环境已安装 IntelliJ IDEA 2021.3+、JDK 1.8
  4. 本地 Windows 能正常访问 Hadoop 集群的 master 节点,已配置 hosts 主机名映射

模块一:Hive 远程服务配置与启动(任务 8.1 核心) ​

要实现 Java 程序远程操作 Hive,必须先开启 Hive 的远程服务。Hive 提供了 HiveServer2 服务,基于 Thrift 协议,支持多客户端并发访问,是 Java 程序通过 JDBC 连接 Hive 的核心基础。

一、核心服务说明 ​

服务名称核心作用启动优先级
metastoreHive 元数据服务,管理 Hive 表的元数据信息(库名、表名、字段、存储位置等),所有对 Hive 的操作都必须先访问元数据必须先启动,再启动 HiveServer2
hiveserver2Hive 远程服务,开放 10000 端口,对外提供 JDBC 连接入口,支持 Java/Python 等客户端远程提交 HQL 执行依赖 metastore 服务,后启动

二、完整配置步骤 ​

步骤 1:配置 Hadoop 代理用户权限 ​

HiveServer2 需要通过 Hadoop 的代理用户权限访问 HDFS,必须先修改 Hadoop 的核心配置文件,允许 root 用户代理所有主机、所有用户组的请求。

xml
<!-- 代码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:分发配置文件到所有子节点 ​

修改后的配置必须同步到集群所有节点,否则配置不生效。

bash
# 代码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 集群,让配置生效 ​

bash
# 进入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」的顺序启动,否则会出现连接失败。

bash
# 代码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 服务是否正常 ​

bash
# 1. 查看10000端口是否被监听(HiveServer2默认端口)
netstat -nltp | grep 10000

# 2. 用beeline客户端测试本地连接,验证服务可用
beeline
# 进入beeline后,执行连接命令
!connect jdbc:hive2://master:10000/default
# 输入用户名root,密码123456,能正常进入Hive命令行,说明服务正常

三、新手高频踩坑避坑指南 ​

  1. ❌ 坑 1:先启动 hiveserver2,再启动 metastore,导致连接失败

    ✅ 解决:必须严格按照「metastore → hiveserver2」的顺序启动,两个服务都启动后,等待 30 秒再测试连接,服务启动需要时间

  2. ❌ 坑 2:修改 core-site.xml 后没有分发到子节点,也没有重启 Hadoop,配置不生效

    ✅ 解决:修改配置后必须分发到所有节点,重启 Hadoop 集群,再启动 Hive 服务

  3. ❌ 坑 3:防火墙没有关闭,10000 端口被拦截,Windows 无法远程连接

    ✅ 解决:关闭集群所有节点的防火墙,或开放 10000 端口

    bash
    # 临时关闭防火墙
    systemctl stop firewalld
    # 永久关闭防火墙
    systemctl disable firewalld
  4. ❌ 坑 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 项目 ​

  1. 打开 IDEA,点击欢迎界面的New Project
  2. 左侧选择Maven,Project SDK 选择 1.8 版本的 JDK,点击 Next
  3. 项目名称填写HiveJavaAPI,选择项目存储路径(如D:\HiveJavaAPI),点击 Finish
  4. 项目创建完成后,自动生成 Maven 项目结构,核心文件为pom.xml(Maven 依赖配置文件)

二、配置 Maven 依赖 ​

打开项目的pom.xml文件,添加 Hadoop 和 Hive 相关依赖,解决 jar 包冲突问题。

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>
    <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 的远程操作。核心接口如下,新手必须掌握:

接口名称核心作用常用方法
DriverJDBC 驱动接口,Hive 的实现类为org.apache.hive.jdbc.HiveDriverClass.forName(driverName)加载驱动
DriverManager驱动管理类,用于创建数据库连接getConnection(url, username, password)获取连接对象
Connection数据库连接对象,代表 Java 程序和 Hive 的连接会话createStatement()创建 Statement 对象、close()关闭连接
StatementSQL 语句执行对象,用于向 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:编写连接测试代码 ​

java
// 代码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:运行程序并验证 ​

  1. 右键代码空白处 → 选择 Run 'ConnectionTest.main ()',运行程序
  2. 控制台输出Process finished with exit code 0,且打印「驱动加载成功、连接成功、数据库创建成功」,说明程序运行正常
  3. 回到 Linux 的 Hive CLI,执行SHOW DATABASES;,能看到新建的 test 数据库,说明远程连接完全正常

五、开发环境搭建避坑指南 ​

  1. ❌ 坑 1:pom.xml 依赖下载失败,出现红色报错

    ✅ 解决:配置 Maven 国内镜像源(阿里云镜像),检查网络是否正常,重新 Reload project

  2. ❌ 坑 2:运行报错ClassNotFoundException: org.apache.hive.jdbc.HiveDriver

    ✅ 解决:检查 pom.xml 的 hive-jdbc 依赖是否正确添加,是否完成了 Maven Reload,依赖是否下载成功

  3. ❌ 坑 3:运行报错Connection refused: connect

    ✅ 解决:

    • 检查 HiveServer2 服务是否正常启动,10000 端口是否被监听
    • 检查 Windows 能否 ping 通 master 节点的 IP
    • 检查防火墙是否关闭,10000 端口是否开放
    • 检查 url 中的主机名和端口是否正确
  4. ❌ 坑 4:运行报错HiveSQLException: Failed to open new session

    ✅ 解决:检查 core-site.xml 的代理用户配置是否正确,是否分发到所有节点,是否重启了 Hadoop 集群


模块三:编写程序实现广电数据的存储(任务 8.3 核心) ​

本模块将封装一个 Hive 操作工具类HiveHelper,把 Hive 的连接、关闭、建库、建表、数据加载等常用操作封装成可复用的方法,实现广电 5 张核心业务表的自动化创建与数据装载。

一、封装 HiveHelper 工具类 ​

在src/main/java目录下创建 Java 类HiveHelper,封装所有 Hive 操作的核心方法。

java
// 代码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 工具类,完成广电数据库创建、表创建、数据加载的全流程自动化执行。

java
// 代码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. ❌ 坑 1:LOAD DATA 语句执行成功,但表中无数据

    ✅ 解决:

    • LOCAL INPATH后面的路径是Linux 服务器上的文件路径,不是 Windows 本地的路径,必须提前把 CSV 文件上传到 Linux 服务器的对应目录
    • 检查 Linux 上的文件路径是否正确,文件是否存在,权限是否正确
    • 检查建表语句的字段分隔符是否和 CSV 文件的分隔符一致
  2. ❌ 坑 2:建表语句执行报错,提示语法错误

    ✅ 解决:检查 Java 字符串中的建表 SQL 是否有拼写错误,字符串拼接是否正确,关键字是否正确,字段类型是否符合 Hive 规范

  3. ❌ 坑 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方法,实现用户状态变更表的数据查询,并打印到控制台。

java
// 代码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. 用户基本数据清洗方法 ​

java
// 代码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. 收视行为数据清洗方法 ​

java
// 代码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. 账单数据清洗方法 ​

java
// 代码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 的方法,实现数据查询与清洗的全流程执行。

java
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 程序调试核心步骤 ​

当程序运行结果不符合预期、出现报错时,通过断点调试可以快速定位问题。

  1. 设置断点:在代码行号的右侧空白处单击,出现红色圆点,即为断点,程序执行到这一行会自动暂停

  2. 进入调试模式:右键代码空白处 → 选择 Debug ' 类名.main ()',启动调试模式

  3. 调试核心快捷键

    快捷键功能新手常用场景
    F8步过,执行当前行,跳到下一行,不进入方法内部逐行执行代码,查看每一步的执行结果
    F7步入,进入当前调用的方法内部查看工具类方法内部的执行情况,定位方法内的报错
    F9恢复程序执行,跳到下一个断点快速跳转到下一个断点,跳过不需要逐行查看的代码
    Alt+F8计算表达式,查看变量的实时值查看 SQL 语句、变量的实时内容,定位 SQL 拼写错误
  4. 查看变量:调试模式下,底部 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,实现统计用户基本表中每个用户等级的用户数量,并打印到控制台。

参考答案
java
/**
 * 统计每个用户等级的用户数量
 * @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 个直播频道。

参考答案
java
/**
 * 统计直播频道观看用户数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目录,字段用逗号分隔。

参考答案
java
/**
 * 导出清洗后的用户数据到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");

基于 Vite 强力驱动 | 纯静态轻量托管