Skip to content

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

一、章节整体概述 ​

1. 章节定位与作用 ​

本章是 Hive 学习从手工命令行操作向企业级程序自动化处理的核心转折点。前 7 章均通过 Hive CLI 手动输入 HQL 完成建表、查询、清洗,存在效率低、无法复用、难以集成到业务系统的局限。

本章通过 Java 语言基于 JDBC 接口实现 Hive 远程调用,将手工操作封装为可复用、可自动化执行的程序,完成广电数据全流程自动化存储、查询与处理,是衔接大数据平台操作与 Java 后端开发的关键章节,也是后续大数据项目开发的核心基础。

2. 课时安排(共 10 学时) ​

课时类型时长核心内容
理论课2 学时HiveServer2 原理、JDBC 接口体系、工具类封装设计思想、客户端配置机制
实验课 13 学时Hive 远程服务配置、IDEA 开发环境搭建、配置文件部署、连接测试
实验课 23 学时广电数据存储程序开发、数据查询与清洗程序开发
调试与答疑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)两大核心服务 ​

服务名称核心作用启动顺序默认端口依赖服务
metastoreHive 元数据服务,管理库、表、字段、存储位置等元数据,所有 Hive 操作都依赖元数据必须先启动9083MySQL、HDFS
hiveserver2Hive 远程服务,对外提供 JDBC 连接入口,接收客户端 HQL 请求并执行依赖 metastore,后启动10000metastore 服务

⚠️ 教学重点强调:必须严格按照「metastore → hiveserver2」的顺序启动,且启动后等待 30 秒再测试,服务初始化需要时间。

2. 服务端三大核心配置文件 ​

以下配置文件部署在 Linux 集群节点上,是 Hadoop、Hive 服务启动的基础。

(1)core-site.xml(Hadoop 核心配置) ​

部署路径:/usr/local/hadoop-3.1.4/etc/hadoop/core-site.xml(所有节点都要有)

完整配置:

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
<?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
<?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:分发配置到所有子节点 ​

bash
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 集群 ​

bash
cd /usr/local/hadoop-3.1.4/sbin
./stop-all.sh
./start-all.sh
./mr-jobhistory-daemon.sh start historyserver

步骤 4:创建 Hive 临时目录 ​

bash
mkdir -p /usr/local/hive-3.1.2/bin/hive-data/tmp

步骤 5:启动 Hive 远程服务(严格顺序) ​

bash
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:端口监听验证

bash
netstat -nltp | grep 10000

方法 2:beeline 客户端验证

bash
beeline
!connect jdbc:hive2://master:10000/default
# 输入用户名root,密码123456,能正常进入命令行则服务正常

4. 高频踩坑与解决方案 ​

问题现象根本原因解决方法
启动 metastore 报错,找不到目录system:java.io.tmpdir指定的临时目录未创建手动创建对应目录
连接时报Failed to open new sessionHadoop 代理用户配置未生效检查 core-site.xml,确认已分发所有节点并重启 Hadoop
报Connection refused服务启动顺序错误,或服务未完全初始化严格按 metastore→hiveserver2 顺序启动,等待 30 秒再测试
Windows 无法连接防火墙拦截,或本地 hosts 未配置关闭所有节点防火墙;本地添加 master 主机名映射
beeline 连接报错MySQL 服务未启动,或元数据库异常检查 MySQL 服务状态,验证 hive-site.xml 的数据库配置

三、模块二:搭建 Hive Java 远程开发环境(本地操作) ​

1. 创建 Maven 项目 ​

  1. 打开 IDEA,点击「New Project」
  2. 左侧选择「Maven」,Project SDK 选择 1.8,点击「Next」
  3. 项目名填写HiveJavaAPI,选择存储路径,点击「Finish」

2. pom.xml 依赖配置 ​

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 点:

  1. Maven 机制:自动加入类路径

    Maven 项目中,src/main/resources是标准资源目录,编译、运行、打包时,会自动将目录下所有文件复制到 Java 运行类路径(classpath)。Hadoop/Hive 框架启动时,会默认从 classpath 扫描加载配置文件,无需额外写代码。

  2. 客户端需要识别集群信息

    Java 程序作为 Hive 客户端,需要通过配置文件知道:

    • HDFS 集群地址(hdfs://master:9820)

    • Hive 服务端口、元数据地址

    • 代理用户权限规则

      没有这些配置,客户端无法正常连接集群、访问 HDFS。

  3. 代码零改动,跨环境兼容

    放在resources后,框架自动加载,代码无需手动指定文件路径。换电脑、换集群,只需替换 3 个 XML 文件,不用改 Java 代码,可移植性极强。

  4. 打包部署友好

    项目打包成 Jar 包后,resources下的配置文件会一起打入 Jar 中,上传到服务器直接运行,无需单独拷贝配置文件。

(2)服务端配置 vs 客户端配置 对比 ​

类型部署位置作用内容要求
服务端配置Linux 集群节点的 hadoop/conf、hive/conf 目录启动 Hadoop、Hive、metastore、hiveserver2 服务服务运行的完整配置
客户端配置本地 IDEA 项目的src/main/resources目录给 Java 程序读取,指导客户端连接集群与服务端配置完全一致,直接拷贝即可

操作要求:从 Linux 集群原样拷贝 3 个 XML 文件,不要修改内容,保证客户端与集群配置一致。

(3)标准项目目录结构 ​

plaintext
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 驱动。

表格

接口 / 类核心作用核心方法使用场景
DriverJDBC 驱动接口,Hive 实现类为HiveDriverClass.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)运行验证标准 ​

  1. 控制台输出Process finished with exit code 0
  2. 依次打印「驱动加载成功、连接成功、数据库创建成功、连接已关闭」
  3. 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. 核心注意事项 ​

  1. 路径问题:LOAD DATA LOCAL INPATH的路径是Linux 服务器本地路径,不是 Windows 本地路径,必须提前上传 CSV 文件。
  2. 分隔符一致:建表语句的分隔符;必须与 CSV 文件实际分隔符一致,否则数据全为 NULL。
  3. 权限问题:若报权限错误,执行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. 调试基本步骤 ​

  1. 设置断点:点击代码行号右侧空白处,出现红色圆点即为断点,程序执行到该行会自动暂停。
  2. 启动调试:右键代码 → 「Debug」,或快捷键Shift+F9。
  3. 控制执行:使用快捷键控制程序运行,观察变量变化。
  4. 查看变量:底部 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. 企业级拓展方向 ​

  1. 连接池优化:引入 HikariCP、Druid 等连接池管理连接,避免频繁创建销毁连接,提升并发性能。
  2. 配置外置化:将连接地址、账号密码放到 properties 文件中,不硬编码在 Java 代码里,方便多环境切换。
  3. SQL 注入防护:使用PreparedStatement预编译 SQL,替代字符串拼接,防止 SQL 注入风险。
  4. 定时任务:结合 Quartz、Spring Task 实现定时自动执行数据清洗任务。

八、全章常见问题汇总 ​

1. 服务端启动类 ​

  1. 启动 metastore 报错找不到目录 → 手动创建system:java.io.tmpdir指定的目录
  2. 两个服务启动后只有一个 RunJar → 启动顺序错误,先启动的服务被异常终止,按顺序重新启动
  3. 10000 端口一直不被监听 → 服务初始化慢,等待 1~2 分钟再查看

2. 客户端连接类 ​

  1. 代理用户权限报错 → 检查 3 个 XML 是否放入 resources 目录,检查配置内容是否正确
  2. 无法解析 master 主机名 → 本地 Windows 未配置 hosts 映射
  3. 依赖下载慢 → 配置阿里云 Maven 镜像

3. 代码运行类 ​

  1. 建表成功但找不到表 → 没有执行USE 数据库名,建表到了 default 库
  2. LOAD DATA 成功但表无数据 → 文件路径错误、分隔符不一致、文件权限不足
  3. 编译报错类名冲突 → 修改 Connection 为 ConnectionTest

九、课后自测与作业 ​

1. 基础简答题 ​

  1. 简述 HiveServer2 相比 HiveServer 的优势。
  2. 简述启动 Hive 远程服务的顺序,并说明原因。
  3. 为什么要将core-site.xml等配置文件放到src/main/resources目录?
  4. 简述 JDBC 中 Statement 和 PreparedStatement 的区别。

2. 实操题 ​

  1. 独立完成 Hive 远程服务配置与启动,并验证服务可用性。
  2. 在 HiveHelper 中新增方法countUserByLevel,统计用户基本表中每个用户等级的用户数量,并打印结果。
  3. 封装方法实现直播频道观看用户数 Top10 统计,将第 6 章的 Top10 逻辑改为 Java 程序实现。

3. 拓展思考题 ​

  1. 思考如何实现程序定时自动执行数据清洗任务。
  2. 调研 Hive JDBC 连接池的配置方法,尝试引入 HikariCP 优化连接管理。

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